136 lines
6.6 KiB
Python
136 lines
6.6 KiB
Python
from celery_app import app
|
||
from datetime import datetime
|
||
import logging
|
||
import requests
|
||
from api_config import API_CONFIG, RETRY_CONFIG
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
@app.task(bind=True, max_retries=None)
|
||
def daily_morning_task(self):
|
||
"""每天定时运行配置中的API接口并支持重试机制"""
|
||
current_time = datetime.now().strftime('%Y-%m-%d %H:%M:%S')
|
||
logger.info(f"开始执行,当前时间: {current_time}")
|
||
|
||
# 获取重试配置
|
||
max_retries = RETRY_CONFIG['max_retries']
|
||
retry_delay = RETRY_CONFIG['retry_delay']
|
||
retry_on_status = RETRY_CONFIG['retry_on_status']
|
||
|
||
# 遍历配置文件中的所有API接口
|
||
results = []
|
||
for api_name, api_config in API_CONFIG.items():
|
||
logger.info(f"开始处理API配置: {api_name}")
|
||
|
||
api_url = api_config['url']
|
||
headers = api_config['headers']
|
||
timeout = api_config['timeout']
|
||
payload = api_config['body']
|
||
method = api_config['method'].upper()
|
||
|
||
try:
|
||
logger.info(f"正在调用API接口: {method} {api_url}")
|
||
|
||
# 根据method选择请求方式
|
||
if method == 'POST':
|
||
response = requests.post(api_url, headers=headers, json=payload, timeout=timeout)
|
||
elif method == 'GET':
|
||
response = requests.get(api_url, headers=headers, params=payload, timeout=timeout)
|
||
elif method == 'PUT':
|
||
response = requests.put(api_url, headers=headers, json=payload, timeout=timeout)
|
||
elif method == 'DELETE':
|
||
response = requests.delete(api_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}")
|
||
|
||
# 构建成功消息
|
||
if response.status_code == 201:
|
||
print(f"[{current_time}] ✅ API调用成功 ({api_name}) - ID: {result_data.get('id')}")
|
||
message = 'API调用成功'
|
||
response_id = result_data.get('id')
|
||
else:
|
||
print(f"[{current_time}] ✅ API调用成功 ({api_name}) - {method}")
|
||
message = f'API调用成功 ({method})'
|
||
response_id = None
|
||
|
||
results.append({
|
||
'api_name': api_name,
|
||
'status': 'success',
|
||
'message': message,
|
||
'response': result_data,
|
||
'response_id': response_id,
|
||
'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_name} - 状态码: {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_name} - 状态码: {response.status_code}")
|
||
results.append({
|
||
'api_name': api_name,
|
||
'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_name} - 状态码: {response.status_code}")
|
||
results.append({
|
||
'api_name': api_name,
|
||
'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_name} - 异常: {str(e)}")
|
||
raise self.retry(countdown=retry_delay, exc=e)
|
||
else:
|
||
logger.error(f"API调用异常,已达到最大重试次数: {max_retries}")
|
||
print(f"[{current_time}] ❌ API调用异常,已达到最大重试次数 - {api_name} - 异常: {str(e)}")
|
||
results.append({
|
||
'api_name': api_name,
|
||
'status': 'error',
|
||
'message': f'API调用异常,已达到最大重试次数: {max_retries}',
|
||
'error': str(e),
|
||
'timestamp': current_time,
|
||
'retry_count': current_retry + 1
|
||
})
|
||
|
||
# 汇总结果
|
||
success_count = len([r for r in results if r['status'] == 'success'])
|
||
total_count = len(results)
|
||
logger.info(f"定时任务执行完成 - 成功: {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
|
||
} |