diff --git a/.env.example b/.env.example index e72e6d6..c43e467 100644 --- a/.env.example +++ b/.env.example @@ -5,6 +5,9 @@ DATABASE_URL=postgresql+asyncpg://user:password@localhost:5432/ai_talk_callback_ CELERY_BROKER_URL=redis://localhost:6379/0 CELERY_RESULT_BACKEND=redis://localhost:6379/0 +# 任务配置 +TASK_ENABLED=true + # Flower监控配置 FLOWER_ENABLED=true FLOWER_PORT=5555 diff --git a/app/api_config.py b/app/api_config.py new file mode 100644 index 0000000..c7981a8 --- /dev/null +++ b/app/api_config.py @@ -0,0 +1,51 @@ +"""API接口配置文件""" + +# API配置 +API_CONFIG = { + # 关闭清洗任务 + 'task_01': { + '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': { + } + }, + # Ice跟进 + 'task_01': { + '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_01': { + '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': { + } + } +} + +# 重试配置 +RETRY_CONFIG = { + 'max_retries': 3, + 'retry_delay': 60, # 重试间隔(秒) + 'retry_on_status': [500, 502, 503, 504, 429] +} \ No newline at end of file diff --git a/app/celery_app.py b/app/celery_app.py index b0115e0..a37adcc 100644 --- a/app/celery_app.py +++ b/app/celery_app.py @@ -34,6 +34,10 @@ celery_app.conf.update( 'task': 'push_data_to_dtc', 'schedule': 60.0, # 每60秒执行一次(1分钟) }, + 'daily-morning-task': { + 'task': 'call_api', + 'schedule': '120.0' + }, }, ) diff --git a/app/celery_tasks.py b/app/celery_tasks.py index 5767028..bfacb40 100644 --- a/app/celery_tasks.py +++ b/app/celery_tasks.py @@ -2,6 +2,8 @@ Celery任务定义 """ +from datetime import datetime +from functools import wraps from sqlalchemy import create_engine import json from app.celery_app import celery_app @@ -15,9 +17,22 @@ from app.callback_service import ( get_related_records_by_unique_data_list ) from app.redis_lock import redis_manager +import requests +from app.api_config import API_CONFIG, RETRY_CONFIG logger = get_celery_tasks_logger() +def conditional_task(enabled=True): + def decorator(task_func): + @wraps(task_func) + def wrapper(*args, **kwargs): + if not enabled: + logger.info(f"Task {task_func.__name__} is disabled") + return None + return task_func(*args, **kwargs) + return wrapper + return decorator + # 创建同步数据库连接用于Celery任务 engine = create_engine( settings.database_url, @@ -26,6 +41,7 @@ engine = create_engine( ) @celery_app.task(bind=True, name='push_data_to_dtc') +@conditional_task(enabled=settings.task_enabled) def push_data_to_dtc_task(self): """ 推送数据给DTC的Celery任务 @@ -204,3 +220,125 @@ def push_data_to_dtc_task(self): except Exception as e: logger.error(f"❌ 释放任务 {task_name} 的分布式锁失败: {e}") + +@celery_app.task(bind=True, name='call_api', max_retries=None) +def execute_call_api_task(self): + """遍历API配置文件中的所有API信息并调用接口""" + current_time = datetime.now().strftime('%Y-%m-%d %H:%M:%S') + logger.info(f"开始遍历API配置文件执行调用任务 - 当前时间: {current_time}") + + # 获取重试配置 + max_retries = RETRY_CONFIG['max_retries'] + retry_delay = RETRY_CONFIG['retry_delay'] + retry_on_status = RETRY_CONFIG['retry_on_status'] + + # 遍历所有API信息 + results = [] + success_count = 0 + + for api_key, call_config in API_CONFIG.items(): + logger.info(f"开始处理API配置: {api_key}") + + call_url = call_config['url'] + method = call_config['method'].upper() + headers = call_config['headers'] + timeout = call_config['timeout'] + payload = call_config['body'] + + try: + logger.info(f"正在调用任务API: {method} {call_url}") + + # 根据method选择请求方式 + if method == 'POST': + response = requests.post(call_url, headers=headers, json=payload, timeout=timeout) + elif method == 'GET': + response = requests.get(call_url, headers=headers, params=payload, timeout=timeout) + elif method == 'PUT': + response = requests.put(call_url, headers=headers, json=payload, timeout=timeout) + elif method == 'DELETE': + response = requests.delete(call_url, headers=headers, timeout=timeout) + else: + raise ValueError(f"不支持的HTTP方法: {method}") + + # 检查响应状态 + if response.status_code in [200, 201, 204]: # 成功状态码 + try: + result_data = response.json() + except: + result_data = {'response': response.text} + + logger.info(f"任务API调用成功,响应数据: {result_data}") + print(f"[{current_time}] ✅ 任务API成功 ({api_key}) - {method} {call_url}") + success_count += 1 + + results.append({ + 'call_api_key': api_key, + 'status': 'success', + 'message': f'任务API调用成功 ({method})', + 'response': result_data, + 'timestamp': current_time, + 'retry_count': getattr(self.request, 'retries', 0) + }) + elif response.status_code in retry_on_status: + # 需要重试的状态码 + current_retry = getattr(self.request, 'retries', 0) + if current_retry < max_retries: + logger.warning(f"任务API需要重试,当前重试次数: {current_retry + 1}/{max_retries}") + print(f"[{current_time}] 🔄 任务API失败,正在重试 ({current_retry + 1}/{max_retries}) - {api_key} - 状态码: {response.status_code}") + raise self.retry(countdown=retry_delay, exc=Exception(f"任务API失败,状态码: {response.status_code}")) + else: + logger.error(f"任务API失败,已达到最大重试次数: {max_retries}") + print(f"[{current_time}] ❌ 任务API失败,已达到最大重试次数 - {api_key} - 状态码: {response.status_code}") + results.append({ + 'call_api_key': api_key, + 'status': 'failed', + 'message': f'任务API失败,已达到最大重试次数: {max_retries}', + 'response': response.text, + 'timestamp': current_time, + 'retry_count': current_retry + 1 + }) + else: + # 其他失败状态码,不重试 + logger.error(f"任务API失败,状态码: {response.status_code}, 响应: {response.text}") + print(f"[{current_time}] ❌ 任务API失败 - {api_key} - 状态码: {response.status_code}") + results.append({ + 'call_api_key': api_key, + 'status': 'failed', + 'message': f'任务API失败,状态码: {response.status_code}', + 'response': response.text, + 'timestamp': current_time, + 'retry_count': getattr(self.request, 'retries', 0) + }) + + except requests.exceptions.RequestException as e: + # 网络异常,需要重试 + current_retry = getattr(self.request, 'retries', 0) + if current_retry < max_retries: + logger.warning(f"任务API异常,需要重试,当前重试次数: {current_retry + 1}/{max_retries}") + print(f"[{current_time}] 🔄 任务API异常,正在重试 ({current_retry + 1}/{max_retries}) - {api_key} - 异常: {str(e)}") + raise self.retry(countdown=retry_delay, exc=e) + else: + logger.error(f"任务API异常,已达到最大重试次数: {max_retries}") + print(f"[{current_time}] ❌ 任务API异常,已达到最大重试次数 - {api_key} - 异常: {str(e)}") + results.append({ + 'call_api_key': api_key, + 'status': 'error', + 'message': f'任务API异常,已达到最大重试次数: {max_retries}', + 'error': str(e), + 'timestamp': current_time, + 'retry_count': current_retry + 1 + }) + + # 汇总结果 + total_count = len(results) + logger.info(f"遍历API配置文件完成 - 成功: {success_count}/{total_count}") + print(f"[{current_time}] 📊 任务汇总 - 成功: {success_count}/{total_count}") + + return { + 'status': 'completed', + 'total_apis': total_count, + 'success_count': success_count, + 'failed_count': total_count - success_count, + 'results': results, + 'timestamp': current_time + } \ No newline at end of file diff --git a/app/config.py b/app/config.py index d585a18..10fd561 100644 --- a/app/config.py +++ b/app/config.py @@ -1,6 +1,5 @@ from pydantic_settings import BaseSettings - class Settings(BaseSettings): @property def redis_url(self) -> str: @@ -19,6 +18,9 @@ class Settings(BaseSettings): celery_timezone: str = "Asia/Shanghai" celery_enable_utc: bool = False + # 任务配置 + task_enabled: bool = False # 是否启用任务 + # Flower监控配置 flower_enabled: bool = True # 是否启用Flower监控 flower_url: str = "http://localhost:5555" # Flower访问URL diff --git a/requirements.txt b/requirements.txt index dc9248e..ecb2092 100644 --- a/requirements.txt +++ b/requirements.txt @@ -12,4 +12,5 @@ python-multipart>=0.0.20 httpx>=0.28.1 python-dotenv>=1.2.1 pytest>=9.0.1 -pytest-asyncio>=1.3.0 \ No newline at end of file +pytest-asyncio>=1.3.0 +requests==2.31.0 \ No newline at end of file