增加执行外部api的定时任务
This commit is contained in:
@@ -5,6 +5,9 @@ DATABASE_URL=postgresql+asyncpg://user:password@localhost:5432/ai_talk_callback_
|
|||||||
CELERY_BROKER_URL=redis://localhost:6379/0
|
CELERY_BROKER_URL=redis://localhost:6379/0
|
||||||
CELERY_RESULT_BACKEND=redis://localhost:6379/0
|
CELERY_RESULT_BACKEND=redis://localhost:6379/0
|
||||||
|
|
||||||
|
# 任务配置
|
||||||
|
TASK_ENABLED=true
|
||||||
|
|
||||||
# Flower监控配置
|
# Flower监控配置
|
||||||
FLOWER_ENABLED=true
|
FLOWER_ENABLED=true
|
||||||
FLOWER_PORT=5555
|
FLOWER_PORT=5555
|
||||||
|
|||||||
51
app/api_config.py
Normal file
51
app/api_config.py
Normal file
@@ -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]
|
||||||
|
}
|
||||||
@@ -34,6 +34,10 @@ celery_app.conf.update(
|
|||||||
'task': 'push_data_to_dtc',
|
'task': 'push_data_to_dtc',
|
||||||
'schedule': 60.0, # 每60秒执行一次(1分钟)
|
'schedule': 60.0, # 每60秒执行一次(1分钟)
|
||||||
},
|
},
|
||||||
|
'daily-morning-task': {
|
||||||
|
'task': 'call_api',
|
||||||
|
'schedule': '120.0'
|
||||||
|
},
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|||||||
@@ -2,6 +2,8 @@
|
|||||||
Celery任务定义
|
Celery任务定义
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
from datetime import datetime
|
||||||
|
from functools import wraps
|
||||||
from sqlalchemy import create_engine
|
from sqlalchemy import create_engine
|
||||||
import json
|
import json
|
||||||
from app.celery_app import celery_app
|
from app.celery_app import celery_app
|
||||||
@@ -15,9 +17,22 @@ from app.callback_service import (
|
|||||||
get_related_records_by_unique_data_list
|
get_related_records_by_unique_data_list
|
||||||
)
|
)
|
||||||
from app.redis_lock import redis_manager
|
from app.redis_lock import redis_manager
|
||||||
|
import requests
|
||||||
|
from app.api_config import API_CONFIG, RETRY_CONFIG
|
||||||
|
|
||||||
logger = get_celery_tasks_logger()
|
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任务
|
# 创建同步数据库连接用于Celery任务
|
||||||
engine = create_engine(
|
engine = create_engine(
|
||||||
settings.database_url,
|
settings.database_url,
|
||||||
@@ -26,6 +41,7 @@ engine = create_engine(
|
|||||||
)
|
)
|
||||||
|
|
||||||
@celery_app.task(bind=True, name='push_data_to_dtc')
|
@celery_app.task(bind=True, name='push_data_to_dtc')
|
||||||
|
@conditional_task(enabled=settings.task_enabled)
|
||||||
def push_data_to_dtc_task(self):
|
def push_data_to_dtc_task(self):
|
||||||
"""
|
"""
|
||||||
推送数据给DTC的Celery任务
|
推送数据给DTC的Celery任务
|
||||||
@@ -204,3 +220,125 @@ def push_data_to_dtc_task(self):
|
|||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"❌ 释放任务 {task_name} 的分布式锁失败: {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
|
||||||
|
}
|
||||||
@@ -1,6 +1,5 @@
|
|||||||
from pydantic_settings import BaseSettings
|
from pydantic_settings import BaseSettings
|
||||||
|
|
||||||
|
|
||||||
class Settings(BaseSettings):
|
class Settings(BaseSettings):
|
||||||
@property
|
@property
|
||||||
def redis_url(self) -> str:
|
def redis_url(self) -> str:
|
||||||
@@ -19,6 +18,9 @@ class Settings(BaseSettings):
|
|||||||
celery_timezone: str = "Asia/Shanghai"
|
celery_timezone: str = "Asia/Shanghai"
|
||||||
celery_enable_utc: bool = False
|
celery_enable_utc: bool = False
|
||||||
|
|
||||||
|
# 任务配置
|
||||||
|
task_enabled: bool = False # 是否启用任务
|
||||||
|
|
||||||
# Flower监控配置
|
# Flower监控配置
|
||||||
flower_enabled: bool = True # 是否启用Flower监控
|
flower_enabled: bool = True # 是否启用Flower监控
|
||||||
flower_url: str = "http://localhost:5555" # Flower访问URL
|
flower_url: str = "http://localhost:5555" # Flower访问URL
|
||||||
|
|||||||
@@ -13,3 +13,4 @@ httpx>=0.28.1
|
|||||||
python-dotenv>=1.2.1
|
python-dotenv>=1.2.1
|
||||||
pytest>=9.0.1
|
pytest>=9.0.1
|
||||||
pytest-asyncio>=1.3.0
|
pytest-asyncio>=1.3.0
|
||||||
|
requests==2.31.0
|
||||||
Reference in New Issue
Block a user