独立出调用外部接口
This commit is contained in:
243
app/external_api_processor.py
Normal file
243
app/external_api_processor.py
Normal file
@@ -0,0 +1,243 @@
|
||||
"""
|
||||
异步调用外部接口的处理器
|
||||
"""
|
||||
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from sqlalchemy import func, select
|
||||
from fastapi import HTTPException
|
||||
import httpx
|
||||
import asyncio
|
||||
from typing import Dict, Any, Optional
|
||||
|
||||
|
||||
from app.database import CallbackData, ExternalApiLog
|
||||
from app.config import settings
|
||||
from app.logger import get_logger
|
||||
|
||||
logger = get_logger("external_api_processor")
|
||||
|
||||
|
||||
async def check_phone_number_threshold(
|
||||
db: AsyncSession,
|
||||
phone_number: str
|
||||
) -> tuple[bool, bool]:
|
||||
"""
|
||||
检查手机号在数据库中的出现次数是否超过阈值
|
||||
|
||||
Args:
|
||||
db: 数据库会话
|
||||
phone_number: 要检查的手机号
|
||||
|
||||
Returns:
|
||||
tuple[是否超过阈值, 是否查询成功]
|
||||
"""
|
||||
try:
|
||||
# 查询手机号在数据库中出现的次数
|
||||
count_query = select(func.count(CallbackData.id)).where(
|
||||
CallbackData.phone_number == phone_number
|
||||
)
|
||||
result = await db.execute(count_query)
|
||||
phone_count = result.scalar() or 0
|
||||
|
||||
logger.info(f"📊 手机号 {phone_number} 在数据库中出现次数: {phone_count}")
|
||||
|
||||
# 判断是否超过阈值
|
||||
exceeds_threshold = phone_count >= settings.count_threshold
|
||||
|
||||
if exceeds_threshold:
|
||||
logger.info(f"✅ 手机号 {phone_number} 出现次数 {phone_count} >= {settings.count_threshold},超过阈值")
|
||||
|
||||
return exceeds_threshold, True
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"❌ 查询手机号 {phone_number} 失败: {e}", exc_info=True)
|
||||
# 查询失败时返回False,表示查询未成功
|
||||
return False, False
|
||||
|
||||
|
||||
async def log_external_api_request(
|
||||
db: AsyncSession,
|
||||
callback_logs_id: int,
|
||||
request_url: str,
|
||||
request_headers: Dict[str, Any],
|
||||
request_body: Dict[str, Any],
|
||||
response_status: Optional[int] = None,
|
||||
response_headers: Optional[Dict[str, Any]] = None,
|
||||
response_body: Optional[str] = None,
|
||||
retry_count: int = 0
|
||||
):
|
||||
"""记录外部API请求日志"""
|
||||
try:
|
||||
api_log = ExternalApiLog(
|
||||
callback_logs_id=callback_logs_id,
|
||||
request_url=request_url,
|
||||
request_headers=request_headers,
|
||||
request_body=request_body,
|
||||
response_status=response_status,
|
||||
response_headers=response_headers,
|
||||
response_body=response_body,
|
||||
retry_count=retry_count
|
||||
)
|
||||
|
||||
db.add(api_log)
|
||||
await db.commit()
|
||||
|
||||
logger.info(f"✅ 外部API请求日志记录成功,ID: {api_log.id}")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"❌ 记录外部API请求日志失败: {e}", exc_info=True)
|
||||
|
||||
|
||||
async def call_external_api_with_retry(
|
||||
db: AsyncSession,
|
||||
request_body: Dict[str, Any],
|
||||
max_retries: int = None,
|
||||
callback_logs_id: int = None
|
||||
) -> tuple[bool, int]:
|
||||
"""
|
||||
调用外部API并支持重试机制
|
||||
|
||||
Args:
|
||||
db: 数据库会话
|
||||
request_body: 请求体
|
||||
max_retries: 最大重试次数
|
||||
callback_logs_id: 回调日志ID
|
||||
|
||||
Returns:
|
||||
tuple[是否成功, 实际重试次数]
|
||||
"""
|
||||
if max_retries is None:
|
||||
max_retries = settings.external_api_retry_max
|
||||
|
||||
logger.info(f"🌐 开始调用外部API: {settings.external_api_url}, 最大重试次数: {max_retries}")
|
||||
|
||||
headers = {
|
||||
"Content-Type": "application/json",
|
||||
"User-Agent": "AITalkCallbackService/1.0"
|
||||
}
|
||||
|
||||
for attempt in range(1, max_retries+1):
|
||||
try:
|
||||
logger.debug(f"📤 第{attempt + 1}次尝试调用外部API")
|
||||
|
||||
async with httpx.AsyncClient(timeout=30.0) as client:
|
||||
response = await client.post(
|
||||
settings.external_api_url,
|
||||
headers=headers,
|
||||
json=request_body
|
||||
)
|
||||
|
||||
logger.debug(f"📥 外部API响应: status={response.status_code}")
|
||||
|
||||
# 记录每次尝试的结果
|
||||
await log_external_api_request(
|
||||
db=db,
|
||||
callback_logs_id=callback_logs_id,
|
||||
request_url=settings.external_api_url,
|
||||
request_headers=headers,
|
||||
request_body=request_body,
|
||||
response_status=response.status_code,
|
||||
response_headers=dict(response.headers),
|
||||
response_body=response.text,
|
||||
retry_count=attempt - 1
|
||||
)
|
||||
|
||||
# 检查响应状态
|
||||
if response.status_code == 200:
|
||||
logger.info(f"✅ 外部API调用成功,状态码: {response.status_code}")
|
||||
return True, attempt - 1
|
||||
else:
|
||||
logger.warning(f"⚠️ 外部API返回非成功状态码: {response.status_code}")
|
||||
|
||||
# 如果是客户端错误(4xx),不重试
|
||||
if 400 <= response.status_code < 500:
|
||||
logger.error(f"❌ 客户端错误,不重试: {response.status_code}")
|
||||
return False, attempt - 1
|
||||
|
||||
except httpx.TimeoutException:
|
||||
logger.warning(f"⏰ 外部API调用超时,第{attempt}次尝试")
|
||||
except httpx.RequestError as e:
|
||||
logger.warning(f"🌐 外部API请求错误,第{attempt}次尝试: {e}")
|
||||
except Exception as e:
|
||||
logger.error(f"❌ 外部API调用异常,第{attempt}次尝试: {e}", exc_info=True)
|
||||
|
||||
# 如果不是最后一次尝试,等待一段时间再重试
|
||||
if attempt < max_retries:
|
||||
await asyncio.sleep(2 ** attempt) # 指数退避
|
||||
|
||||
# 所有重试都失败了
|
||||
logger.error(f"❌ 外部API调用失败,已重试{max_retries}次")
|
||||
return False, max_retries
|
||||
|
||||
|
||||
async def process_external_api_call(
|
||||
db: AsyncSession,
|
||||
callback_data: dict,
|
||||
siteId: str,
|
||||
callback_log_id: int
|
||||
):
|
||||
"""
|
||||
异步调用外部接口
|
||||
|
||||
Args:
|
||||
db: 数据库会话
|
||||
callback_data: 回调数据
|
||||
siteId: 站点ID
|
||||
callback_log_id: 回调日志ID
|
||||
|
||||
Returns:
|
||||
None: 无返回值
|
||||
"""
|
||||
|
||||
# 提取所有手机号并去重
|
||||
phone_numbers_set = set()
|
||||
data_list = callback_data.get('data', [])
|
||||
for item in data_list:
|
||||
number_data = item.get('number_data', {})
|
||||
phone_number = number_data.get('number')
|
||||
if phone_number:
|
||||
phone_numbers_set.add(phone_number)
|
||||
|
||||
phone_numbers = list(phone_numbers_set)
|
||||
logger.info(f"📱 提取到的手机号列表(去重后): {phone_numbers}")
|
||||
|
||||
# 检查手机号列表是否为空
|
||||
if not phone_numbers or len(phone_numbers) == 0:
|
||||
logger.warning(f"⚠️ 手机号列表为空,直接退出异步处理")
|
||||
return
|
||||
|
||||
# 检查去重后的手机号是否超过 settings.count_threshold,超过直接返回
|
||||
for phone_number in phone_numbers:
|
||||
|
||||
# 检查手机号是否超过阈值
|
||||
exceeds_threshold, query_success = await check_phone_number_threshold(db, phone_number)
|
||||
|
||||
# 如果查询失败,直接退出异步处理
|
||||
if not query_success:
|
||||
logger.error(f"❌ 查询手机号 {phone_number} 失败,直接退出异步处理")
|
||||
return
|
||||
|
||||
if exceeds_threshold:
|
||||
logger.info(f"✅ 手机号 {phone_number} 出现次数超过阈值,跳过处理")
|
||||
continue
|
||||
|
||||
# 手机号出现的次数少于settings.count_threshold,调用外部API
|
||||
# 调用外部API并支持重试
|
||||
request_body = callback_data
|
||||
success, retry_count = await call_external_api_with_retry(
|
||||
db=db,
|
||||
request_body=request_body,
|
||||
max_retries=settings.external_api_retry_max,
|
||||
callback_logs_id=callback_log_id
|
||||
)
|
||||
|
||||
if success:
|
||||
logger.info(f"✅ 外部API调用成功,重试次数: {retry_count}")
|
||||
return
|
||||
else:
|
||||
logger.error(f"❌ 外部API调用失败,已重试{retry_count}次")
|
||||
return
|
||||
|
||||
# 如果所有手机号都被跳过,返回成功但未处理
|
||||
logger.info(f"📝 所有手机号都被跳过,返回成功但未处理")
|
||||
return
|
||||
Reference in New Issue
Block a user