改造配置文件,去掉敏感信息
This commit is contained in:
@@ -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}")
|
||||
|
||||
|
||||
Reference in New Issue
Block a user