From f1f694bd63df9f09e4f160836f433cce91da3078 Mon Sep 17 00:00:00 2001 From: "mark.tian" Date: Wed, 3 Dec 2025 21:58:16 +0800 Subject: [PATCH] =?UTF-8?q?=E7=8B=AC=E7=AB=8B=E5=87=BA=E8=B0=83=E7=94=A8?= =?UTF-8?q?=E5=A4=96=E9=83=A8=E6=8E=A5=E5=8F=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app/external_api_processor.py | 243 ++++++++++++++++++++++++++++++++++ 1 file changed, 243 insertions(+) create mode 100644 app/external_api_processor.py diff --git a/app/external_api_processor.py b/app/external_api_processor.py new file mode 100644 index 0000000..be137f8 --- /dev/null +++ b/app/external_api_processor.py @@ -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 \ No newline at end of file