Files
task-executer/call_api_task.py
2025-12-09 19:09:29 +08:00

132 lines
6.3 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""呼叫任务API函数"""
from celery_app import app
from datetime import datetime
import logging
import requests
from api_config import API_CONFIG, RETRY_CONFIG
from logging_config import get_logger
# 获取配置好的logger
logger = get_logger('call_api_task')
@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
}