From 5616d32024003a37183688c9ed58273ed18676c1 Mon Sep 17 00:00:00 2001 From: "mark.tian" Date: Thu, 11 Dec 2025 11:23:26 +0800 Subject: [PATCH] =?UTF-8?q?=E6=94=B9=E9=80=A0=E9=85=8D=E7=BD=AE=E6=96=87?= =?UTF-8?q?=E4=BB=B6=EF=BC=8C=E5=8E=BB=E6=8E=89=E6=95=8F=E6=84=9F=E4=BF=A1?= =?UTF-8?q?=E6=81=AF?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .env.example | 15 ++----- app/api_config.py | 89 +++++++++++++++++++++++------------------ app/callback_service.py | 2 +- app/celery_tasks.py | 45 ++++++++++++++++----- app/config.py | 10 +++++ app/redis_lock.py | 4 ++ 6 files changed, 104 insertions(+), 61 deletions(-) diff --git a/.env.example b/.env.example index 3d0ebbe..eac4728 100644 --- a/.env.example +++ b/.env.example @@ -16,21 +16,11 @@ FLOWER_BASIC_AUTH=admin:admin123 FLOWER_URL_PREFIX= FLOWER_URL=http://localhost:5555 -# Redis配置 (扩展配置,如果需要覆盖默认值) -REDIS_PASSWORD= -REDIS_MAX_CONNECTIONS=20 -REDIS_TIMEOUT=5 - -# Redis分布式锁配置 -REDIS_LOCK_TIMEOUT=300 -REDIS_LOCK_MAX_RETRIES=10 -REDIS_LOCK_RETRY_DELAY=0.5 - # 业务配置 COUNT_THRESHOLD=3 EXTERNAL_API_ENABLED=false EXTERNAL_API_RETRY_MAX=3 -EXTERNAL_API_URL=https://jeep-api.d2c.stlassac.com/api/openapi/customerApi/aiTaskResultFail +EXTERNAL_API_URL=your_external_api_url_here # 日志配置 LOG_LEVEL=INFO @@ -40,6 +30,9 @@ LOG_BACKUP_COUNT=5 LOG_FORMAT=%(asctime)s - %(name)s - %(levelname)s - %(message)s LOG_DATE_FORMAT=%Y-%m-%d %H:%M:%S +# API配置 +API_AUTHORIZATION_TOKEN=your_bearer_token_here + # 应用配置 APP_NAME=AI Talk Callback API ENVIRONMENT=production # 环境: production, test, development diff --git a/app/api_config.py b/app/api_config.py index b203ee2..a04b44c 100644 --- a/app/api_config.py +++ b/app/api_config.py @@ -1,46 +1,59 @@ """API接口配置文件""" # API配置 -API_CONFIG = { - # 关闭清洗任务 - 'task_1': { - 'url': 'http://agent-api.ai.telrobot.top/agent-api/new/user/8f6bae44-e02a-4393-a6cf-adb861bf7d8c/task_start/a3bdd5b6-1932-4bc9-95bf-dbc4bccb171f', - 'method': 'PUT', - 'headers': { - 'Content-Type': 'application/json', - 'User-Agent': 'Celery-Task/1.0', - 'Authorization': 'Bearer eyJ0eXAiOiJKV1QiLCJhbGciOiJSUzI1NiJ9.eyJhdWQiOiIxIiwianRpIjoiMTY4ZGY2NTNlNzBlMTlmNjlmMjQ1ODBmNTUxMzljNGQyODlkN2FiNWY0MmZhOGMxNzE0Y2Y4OTI5NjYxYjE1NGU0N2QyNzAyY2VmNTZiY2IiLCJpYXQiOjE3NjEwMDk4MDYsIm5iZiI6MTc2MTAwOTgwNiwiZXhwIjoxNzkyNTQ1ODA2LCJzdWIiOiIxMzQ1ZGY3MS1hZTdlLTRjMmYtYWIzMS1mMjhmODE2OTY2ODEiLCJzY29wZXMiOltdfQ.oAoKOmpo6-KBjLPh7p1YIQbtzQAB11xrooQon8ekj8rK8RC5UpsW79hrhG2dOZY2ZY4OURlYTD-rI5UDeXtNIlSyF1Rf1s1d-mOw4Tgrq8pPDh5oIvi1mKuWZWj2E-a8HUp2Eg3c2Hx56rTv9G4xCSCUm6fWghTdpa7gJQZgtBOWowuAad0ErsvrlO4R7CotFfVG4-hTZEK7eZOJIQHX2e-tfeHbRDm2Qoou9uwNl-3_LDpfyqrXbUrF-Rtlqq0aOoSV9HDJfR0YSFzk6uj0yVNt00G_7EdvpUveqd1xDQWy9qaJzxn771QE0M6aNiFCXYzAq8F9AJAvEl92U8xfsvM24xLBAcjR_FxOQordNeLn_xtDW9-fcNlNfV-ngf_BWpNsjFFT9T3QAjcXkiWP1eoE7x69NWJDYAqKKInSF8-5-md5wsMDm80VZYcB4BOrC2t7LaZhlziErib_4SW21DvKdrdAhIVZDHFqcWITMYSIddG52f7VA1jKqkssMKHfvaWDtefbzFUjhp-45C3rN7oiQ9sDgaye3VjFrLE0tFIOdXhFUcgn98G5SxVR9Rs72sccOxuYTdZwey0DZxX6uzwlLigi4-Zt0LQOW4TSRgfu9NtuYjzR_MqWoJatcu_3X1Lhc5OwIyTsms8rkz0JogjbF-4jodhSzXuhBrp--iA' - }, - 'timeout': 30, - 'body': { - } +''' +# 关闭清洗任务 +'task_1': { + 'url': 'http://agent-api.ai.telrobot.top/agent-api/new/user/8f6bae44-e02a-4393-a6cf-adb861bf7d8c/task_start/a3bdd5b6-1932-4bc9-95bf-dbc4bccb171f', + 'method': 'PUT', + 'headers': { + 'Content-Type': 'application/json', + 'User-Agent': 'Celery-Task/1.0' }, - # Ice跟进 - 'task_2': { - 'url': 'http://agent-api.ai.telrobot.top/agent-api/new/user/8f6bae44-e02a-4393-a6cf-adb861bf7d8c/task_start/db17ba5f-666a-4dc0-a635-1d4d0059338d', - 'method': 'PUT', - 'headers': { - 'Content-Type': 'application/json', - 'User-Agent': 'Celery-Task/1.0', - 'Authorization': 'Bearer eyJ0eXAiOiJKV1QiLCJhbGciOiJSUzI1NiJ9.eyJhdWQiOiIxIiwianRpIjoiMTY4ZGY2NTNlNzBlMTlmNjlmMjQ1ODBmNTUxMzljNGQyODlkN2FiNWY0MmZhOGMxNzE0Y2Y4OTI5NjYxYjE1NGU0N2QyNzAyY2VmNTZiY2IiLCJpYXQiOjE3NjEwMDk4MDYsIm5iZiI6MTc2MTAwOTgwNiwiZXhwIjoxNzkyNTQ1ODA2LCJzdWIiOiIxMzQ1ZGY3MS1hZTdlLTRjMmYtYWIzMS1mMjhmODE2OTY2ODEiLCJzY29wZXMiOltdfQ.oAoKOmpo6-KBjLPh7p1YIQbtzQAB11xrooQon8ekj8rK8RC5UpsW79hrhG2dOZY2ZY4OURlYTD-rI5UDeXtNIlSyF1Rf1s1d-mOw4Tgrq8pPDh5oIvi1mKuWZWj2E-a8HUp2Eg3c2Hx56rTv9G4xCSCUm6fWghTdpa7gJQZgtBOWowuAad0ErsvrlO4R7CotFfVG4-hTZEK7eZOJIQHX2e-tfeHbRDm2Qoou9uwNl-3_LDpfyqrXbUrF-Rtlqq0aOoSV9HDJfR0YSFzk6uj0yVNt00G_7EdvpUveqd1xDQWy9qaJzxn771QE0M6aNiFCXYzAq8F9AJAvEl92U8xfsvM24xLBAcjR_FxOQordNeLn_xtDW9-fcNlNfV-ngf_BWpNsjFFT9T3QAjcXkiWP1eoE7x69NWJDYAqKKInSF8-5-md5wsMDm80VZYcB4BOrC2t7LaZhlziErib_4SW21DvKdrdAhIVZDHFqcWITMYSIddG52f7VA1jKqkssMKHfvaWDtefbzFUjhp-45C3rN7oiQ9sDgaye3VjFrLE0tFIOdXhFUcgn98G5SxVR9Rs72sccOxuYTdZwey0DZxX6uzwlLigi4-Zt0LQOW4TSRgfu9NtuYjzR_MqWoJatcu_3X1Lhc5OwIyTsms8rkz0JogjbF-4jodhSzXuhBrp--iA' - }, - 'timeout': 30, - 'body': { - } - }, - # 战败清洗跟进 - 'task_3': { - 'url': 'http://agent-api.ai.telrobot.top/agent-api/new/user/8f6bae44-e02a-4393-a6cf-adb861bf7d8c/task_start/2121d91e-c7d8-47cf-8a67-fa2f22e087ae', - 'method': 'PUT', - 'headers': { - 'Content-Type': 'application/json', - 'User-Agent': 'Celery-Task/1.0', - 'Authorization': 'Bearer eyJ0eXAiOiJKV1QiLCJhbGciOiJSUzI1NiJ9.eyJhdWQiOiIxIiwianRpIjoiMTY4ZGY2NTNlNzBlMTlmNjlmMjQ1ODBmNTUxMzljNGQyODlkN2FiNWY0MmZhOGMxNzE0Y2Y4OTI5NjYxYjE1NGU0N2QyNzAyY2VmNTZiY2IiLCJpYXQiOjE3NjEwMDk4MDYsIm5iZiI6MTc2MTAwOTgwNiwiZXhwIjoxNzkyNTQ1ODA2LCJzdWIiOiIxMzQ1ZGY3MS1hZTdlLTRjMmYtYWIzMS1mMjhmODE2OTY2ODEiLCJzY29wZXMiOltdfQ.oAoKOmpo6-KBjLPh7p1YIQbtzQAB11xrooQon8ekj8rK8RC5UpsW79hrhG2dOZY2ZY4OURlYTD-rI5UDeXtNIlSyF1Rf1s1d-mOw4Tgrq8pPDh5oIvi1mKuWZWj2E-a8HUp2Eg3c2Hx56rTv9G4xCSCUm6fWghTdpa7gJQZgtBOWowuAad0ErsvrlO4R7CotFfVG4-hTZEK7eZOJIQHX2e-tfeHbRDm2Qoou9uwNl-3_LDpfyqrXbUrF-Rtlqq0aOoSV9HDJfR0YSFzk6uj0yVNt00G_7EdvpUveqd1xDQWy9qaJzxn771QE0M6aNiFCXYzAq8F9AJAvEl92U8xfsvM24xLBAcjR_FxOQordNeLn_xtDW9-fcNlNfV-ngf_BWpNsjFFT9T3QAjcXkiWP1eoE7x69NWJDYAqKKInSF8-5-md5wsMDm80VZYcB4BOrC2t7LaZhlziErib_4SW21DvKdrdAhIVZDHFqcWITMYSIddG52f7VA1jKqkssMKHfvaWDtefbzFUjhp-45C3rN7oiQ9sDgaye3VjFrLE0tFIOdXhFUcgn98G5SxVR9Rs72sccOxuYTdZwey0DZxX6uzwlLigi4-Zt0LQOW4TSRgfu9NtuYjzR_MqWoJatcu_3X1Lhc5OwIyTsms8rkz0JogjbF-4jodhSzXuhBrp--iA' - }, - 'timeout': 30, - 'body': { - } + 'timeout': 30, + 'body': { } +}, +# Ice跟进 +'task_2': { + 'url': 'http://agent-api.ai.telrobot.top/agent-api/new/user/8f6bae44-e02a-4393-a6cf-adb861bf7d8c/task_start/db17ba5f-666a-4dc0-a635-1d4d0059338d', + 'method': 'PUT', + 'headers': { + 'Content-Type': 'application/json', + 'User-Agent': 'Celery-Task/1.0' + }, + 'timeout': 30, + 'body': { + } +}, +# 战败清洗跟进 +'task_3': { + 'url': 'http://agent-api.ai.telrobot.top/agent-api/new/user/8f6bae44-e02a-4393-a6cf-adb861bf7d8c/task_start/2121d91e-c7d8-47cf-8a67-fa2f22e087ae', + 'method': 'PUT', + 'headers': { + 'Content-Type': 'application/json', + 'User-Agent': 'Celery-Task/1.0' + }, + 'timeout': 30, + 'body': { + } +} +''' +API_CONFIG = { + 'task_test': { + 'url': 'https://httpbin.org/post', + 'method': 'POST', + 'headers': { + 'Content-Type': 'application/json', + 'User-Agent': 'Celery-Task/1.0' + }, + 'timeout': 30, + 'body': { + 'test': 'data', + 'timestamp': '2025-12-11T10:00:00Z', + 'source': 'api_config_test' + } + }, } # 重试配置 diff --git a/app/callback_service.py b/app/callback_service.py index 7bb7ce9..f995506 100644 --- a/app/callback_service.py +++ b/app/callback_service.py @@ -119,7 +119,7 @@ async def save_callback_data_items( if existing_data_list: logger.warning(f"⚠️ callback_log_id {callback_failure_log_id} 已存在 {len(existing_data_list)} 条数据,跳过保存") - return False + return True for item in callback_data_items: # 获取手机号 diff --git a/app/celery_tasks.py b/app/celery_tasks.py index d9a2e1e..84bc99c 100644 --- a/app/celery_tasks.py +++ b/app/celery_tasks.py @@ -4,7 +4,6 @@ Celery任务定义 from datetime import datetime from functools import wraps -from sqlalchemy import create_engine import json from app.celery_app import celery_app from app.config import settings @@ -16,7 +15,7 @@ from app.callback_service import ( mark_callback_log_completed, get_related_records_by_unique_data_list ) -from app.redis_lock import redis_manager, async_redis_manager, distributed_lock +from app.redis_lock import async_redis_manager, distributed_lock import requests from app.api_config import API_CONFIG, RETRY_CONFIG @@ -55,7 +54,6 @@ def push_data_to_dtc_task(self): 自动获取一条未完成的回调请求进行处理 """ import asyncio - from functools import partial async def async_task(): task_name = 'push_data_to_dtc' @@ -84,7 +82,7 @@ def push_data_to_dtc_task(self): state='PROGRESS', meta={'current': 0, 'total': 100, 'status': f'正在获取下一条回调请求日志...'} ) - logger.info("2") + # 获取一条未完成的回调请求(按创建时间取最小值) query_success, callback_log_data = await get_uncompleted_callback_log(db) if not query_success: @@ -93,7 +91,7 @@ def push_data_to_dtc_task(self): if not query_success: logger.error("❌ 再次查询未完成的回调请求失败,停止任务执行") raise Exception("查询未完成的回调请求失败,任务停止执行") - logger.info("3") + if not callback_log_data: logger.info("📋 没有找到未完成的回调请求") return {"status": "skipped", "message": "没有找到未完成的回调请求"} @@ -227,8 +225,16 @@ def push_data_to_dtc_task(self): finally: # 分布式锁会通过上下文管理器自动释放 logger.debug(f"🔓 任务 {task_name} 的分布式锁已通过上下文管理器处理") + + # 清理Redis连接 + try: + await async_redis_manager.close() + logger.debug(f"🔌 Redis连接已清理") + except Exception as e: + logger.warning(f"⚠️ 清理Redis连接时出现警告: {e}") # 在同步的Celery任务中运行异步代码 + loop = None try: # 创建新的事件循环 loop = asyncio.new_event_loop() @@ -238,12 +244,22 @@ def push_data_to_dtc_task(self): logger.error(f"❌ 异步任务执行失败: {e}", exc_info=True) return {"status": "error", "message": str(e)} finally: - # 清理事件循环 - try: - if 'loop' in locals(): + # 确保所有异步任务完成后再关闭事件循环 + if loop and not loop.is_closed(): + try: + # 等待所有待处理的任务完成 + pending = asyncio.all_tasks(loop) + if pending: + logger.debug(f"⏳ 等待 {len(pending)} 个异步任务完成...") + loop.run_until_complete(asyncio.gather(*pending, return_exceptions=True)) + + # 关闭事件循环 loop.close() - except: - pass + logger.debug(f"🔌 事件循环已正确关闭") + except Exception as e: + logger.warning(f"⚠️ 关闭事件循环时出现警告: {e}") + else: + logger.debug(f"🔌 事件循环已关闭或未创建") @celery_app.task(bind=True, name='call_api', max_retries=None) @@ -267,10 +283,17 @@ def execute_call_api_task(self): call_url = call_config['url'] method = call_config['method'].upper() - headers = call_config['headers'] + headers = call_config['headers'].copy() # 复制headers避免修改原配置 timeout = call_config['timeout'] payload = call_config['body'] + # 从环境变量获取Authorization并添加到headers + if settings.api_authorization_token: + headers['Authorization'] = settings.api_authorization_token + else: + logger.warning(f"⚠️ API Authorization Token 未配置,跳过API调用: {api_key}") + continue + try: logger.info(f"正在调用任务API: {method} {call_url}") diff --git a/app/config.py b/app/config.py index 0d1cd1e..9471f7d 100644 --- a/app/config.py +++ b/app/config.py @@ -53,6 +53,16 @@ class Settings(BaseSettings): log_format: str = "%(asctime)s - %(name)s - %(levelname)s - %(message)s" log_date_format: str = "%Y-%m-%d %H:%M:%S" + # API配置 + _api_authorization_token: str = "" # 内部存储API Authorization Token + + @property + def api_authorization_token(self) -> str: + """API Authorization Token - 只在生产环境下返回""" + if self.environment.lower() == "production": + return self._api_authorization_token + return "" + # 应用配置 app_name: str = "AI Talk Callback API" environment: str = "production" # 环境: production, test, development diff --git a/app/redis_lock.py b/app/redis_lock.py index f024850..7e99bc3 100644 --- a/app/redis_lock.py +++ b/app/redis_lock.py @@ -190,6 +190,10 @@ class AsyncRedisManager: await self.redis_pool.wait_closed() logger.info("🔴 Redis连接已关闭") + async def close(self): + """断开Redis连接 (disconnect方法的别名)""" + await self.disconnect() + async def create_lock(self, key: str, timeout: int = None) -> AsyncRedisLock: """创建分布式锁""" if not self.redis_pool: