diff --git a/app/api_config.py b/app/api_config.py index 2247eff..7e018bb 100644 --- a/app/api_config.py +++ b/app/api_config.py @@ -1,6 +1,5 @@ """API接口配置文件""" -''' # API配置 API_CONFIG = { # 关闭清洗任务 @@ -40,8 +39,8 @@ API_CONFIG = { } } } -''' +''' API_CONFIG = { 'task_test': { 'url': 'http://localhost:8000/health', @@ -55,6 +54,7 @@ API_CONFIG = { } }, } +''' # 重试配置 RETRY_CONFIG = { diff --git a/app/celery_app.py b/app/celery_app.py index 40503c4..325c92d 100644 --- a/app/celery_app.py +++ b/app/celery_app.py @@ -4,9 +4,13 @@ Celery应用配置 from celery import Celery from celery.schedules import crontab from app.config import settings -from app.logger import get_celery_logger +from app.logger import get_celery_logger, get_celery_beat_logger, get_celery_worker_logger, LoggerManager +# 确保日志系统初始化 +LoggerManager.setup_logging() logger = get_celery_logger() +beat_logger = get_celery_beat_logger() +worker_logger = get_celery_worker_logger() # 创建Celery应用实例 celery_app = Celery( @@ -29,6 +33,10 @@ celery_app.conf.update( worker_prefetch_multiplier=1, worker_max_tasks_per_child=1000, + # Worker日志配置 + 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', + # Beat 调度配置 # crontab(hour=9, minute=0) # 每天上午9点执行 beat_schedule={ @@ -43,4 +51,17 @@ celery_app.conf.update( }, ) -logger.info("🌿 Celery应用配置完成") \ No newline at end of file +logger.info("🌿 Celery应用配置完成") +beat_logger.info("📅 Celery Beat 调度器配置完成") +beat_logger.info("📋 定时任务列表:") +for task_name, task_config in celery_app.conf.beat_schedule.items(): + beat_logger.info(f" - {task_name}: {task_config['schedule']}秒") + +# Worker 配置日志 +worker_logger.info("🔧 Celery Worker 配置完成") +worker_logger.info("⚙️ Worker 配置参数:") +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.worker_prefetch_multiplier}") +worker_logger.info(f" - 每个子进程最大任务数: {celery_app.conf.worker_max_tasks_per_child}") +worker_logger.info(f" - 包含模块: {celery_app.conf.include}") \ No newline at end of file diff --git a/app/logger.py b/app/logger.py index 72d792e..c34d4dc 100644 --- a/app/logger.py +++ b/app/logger.py @@ -115,30 +115,73 @@ celery_tasks_logger = get_logger("celery_tasks") callback_service_logger = get_logger("callback_service") redis_logger = get_logger("redis") +# 通用文件日志器设置函数 +def setup_celery_file_logger(logger_name: str, log_file: str, level: str = "INFO"): + """为Celery组件设置专用的文件日志器,同时保持控制台输出""" + # 创建日志目录 + log_dir = Path(log_file).parent + log_dir.mkdir(parents=True, exist_ok=True) + + # 获取日志器 + logger = get_logger(logger_name) + logger.setLevel(getattr(logging, level.upper())) + + # 清除现有文件处理器(避免重复添加) + for handler in logger.handlers[:]: + if isinstance(handler, logging.handlers.TimedRotatingFileHandler): + logger.removeHandler(handler) + + # 创建文件处理器 + file_handler = logging.handlers.TimedRotatingFileHandler( + filename=log_file, + when='midnight', + interval=1, + backupCount=7, + encoding='utf-8', + ) + file_handler.suffix = "%Y%m%d" + file_handler.setLevel(getattr(logging, level.upper())) + + # 格式化器 + formatter = logging.Formatter( + fmt='%(asctime)s - %(name)s - %(levelname)s - [%(filename)s:%(lineno)d] - %(message)s', + datefmt='%Y-%m-%d %H:%M:%S' + ) + file_handler.setFormatter(formatter) + + # 添加文件处理器 + logger.addHandler(file_handler) + logger.propagate = True + + return logger + +# 各模块专用日志器映射 +CELERY_LOGGERS = { + "celery_tasks": ("logs/celery_tasks.log", "INFO"), + "celery": ("logs/celery_app.log", "INFO"), + "celery.beat": ("logs/celery_beat.log", "WARNING"), + "celery.worker": ("logs/celery_worker.log", "WARNING"), +} + +# 初始化所有Celery文件日志器 +_celery_file_loggers = {} +for logger_name, (log_file, level) in CELERY_LOGGERS.items(): + _celery_file_loggers[logger_name] = setup_celery_file_logger(logger_name, log_file, level) + # 向后兼容的默认日志器 logger = app_logger -# 便捷获取各模块日志器的函数 -def get_main_logger(): - """获取main模块日志器""" - return main_logger +# 通用获取器函数 +def get_celery_file_logger(logger_name: str): + """获取指定名称的Celery文件日志器""" + return _celery_file_loggers.get(logger_name, get_logger(logger_name)) -def get_routes_logger(): - """获取routes模块日志器""" - return routes_logger - -def get_celery_logger(): - """获取celery模块日志器""" - return celery_logger - -def get_celery_tasks_logger(): - """获取celery_tasks模块日志器""" - return celery_tasks_logger - -def get_callback_service_logger(): - """获取callback_service模块日志器""" - return callback_service_logger - -def get_redis_logger(): - """获取redis模块日志器""" - return redis_logger \ No newline at end of file +# 简化的便捷函数 +def get_main_logger(): return main_logger +def get_routes_logger(): return routes_logger +def get_celery_logger(): return get_celery_file_logger("celery") +def get_celery_tasks_logger(): return get_celery_file_logger("celery_tasks") +def get_callback_service_logger(): return callback_service_logger +def get_redis_logger(): return redis_logger +def get_celery_beat_logger(): return get_celery_file_logger("celery.beat") +def get_celery_worker_logger(): return get_celery_file_logger("celery.worker") \ No newline at end of file