This commit is contained in:
mark.tian
2025-12-13 13:43:04 +08:00

View File

@@ -1,10 +1,16 @@
""" """
Celery应用配置 Celery应用配置
""" """
from celery import Celery from celery import Celery
from celery.schedules import crontab from celery.schedules import crontab
from app.config import settings from app.config import settings
from app.logger import get_celery_logger, get_celery_beat_logger, get_celery_worker_logger, LoggerManager from app.logger import (
get_celery_logger,
get_celery_beat_logger,
get_celery_worker_logger,
LoggerManager,
)
# 确保日志系统初始化 # 确保日志系统初始化
LoggerManager.setup_logging() LoggerManager.setup_logging()
@@ -17,7 +23,7 @@ celery_app = Celery(
"ai_talk_callback", "ai_talk_callback",
broker=settings.celery_broker_url, broker=settings.celery_broker_url,
backend=settings.celery_result_backend, backend=settings.celery_result_backend,
include=['app.celery_tasks'] include=["app.celery_tasks"],
) )
# Celery配置 # Celery配置
@@ -32,21 +38,15 @@ celery_app.conf.update(
task_soft_time_limit=25 * 60, # 25分钟软超时 task_soft_time_limit=25 * 60, # 25分钟软超时
worker_prefetch_multiplier=1, worker_prefetch_multiplier=1,
worker_max_tasks_per_child=1000, worker_max_tasks_per_child=1000,
# Worker日志配置 # Worker日志配置
worker_log_format='[%(asctime)s: %(levelname)s/%(processName)s] %(message)s', worker_log_format="[%(asctime)s: %(levelname)s/%(processName)s] %(message)s",
worker_task_log_format='[%(asctime)s: %(levelname)s/%(processName)s][%(task_name)s(%(task_id)s)] %(message)s', worker_task_log_format="[%(asctime)s: %(levelname)s/%(processName)s][%(task_name)s(%(task_id)s)] %(message)s",
# Beat 调度配置 # Beat 调度配置
# crontab(hour=9, minute=0) # 每天上午9点执行 # crontab(hour=9, minute=0) # 每天上午9点执行
beat_schedule={ beat_schedule={
'push-data-to-dtc-every-minute': { "daily-morning-task": {
'task': 'push_data_to_dtc', "task": "call_api",
'schedule': 120.0, # 每120秒执行一次(2分钟) "schedule": crontab(hour=9, minute=0), # 每天上午9点执行
},
'daily-morning-task': {
'task': 'call_api',
'schedule': crontab(hour=9, minute=0), # 每天上午9点执行
}, },
}, },
) )
@@ -63,5 +63,7 @@ worker_logger.info("⚙️ Worker 配置参数:")
worker_logger.info(f" - 任务超时: {celery_app.conf.task_time_limit}秒") worker_logger.info(f" - 任务超时: {celery_app.conf.task_time_limit}秒")
worker_logger.info(f" - 软超时: {celery_app.conf.task_soft_time_limit}秒") worker_logger.info(f" - 软超时: {celery_app.conf.task_soft_time_limit}秒")
worker_logger.info(f" - 预取倍数: {celery_app.conf.worker_prefetch_multiplier}") worker_logger.info(f" - 预取倍数: {celery_app.conf.worker_prefetch_multiplier}")
worker_logger.info(f" - 每个子进程最大任务数: {celery_app.conf.worker_max_tasks_per_child}") worker_logger.info(
worker_logger.info(f" - 包含模块: {celery_app.conf.include}") f" - 每个子进程最大任务数: {celery_app.conf.worker_max_tasks_per_child}"
)
worker_logger.info(f" - 包含模块: {celery_app.conf.include}")