You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Python RQ任务队列分任务分Worker日志实现方案咨询

一、Python代码实现方案

1. 本地独立文件日志

不要在任务函数开头重复调用logging.basicConfig(多次调用会导致配置混乱),应该为每个任务创建独立的日志器实例,绑定专属文件处理器,确保日志隔离:

import logging
from rq import get_current_job
import os

# 提前创建日志目录
os.makedirs("task_logs", exist_ok=True)

def your_task():
    job = get_current_job()
    # 用任务ID作为日志器唯一标识
    logger = logging.getLogger(f"rq_task_{job.id}")
    logger.setLevel(logging.INFO)
    
    # 避免重复添加处理器导致日志重复输出
    if not logger.handlers:
        file_handler = logging.FileHandler(f"task_logs/{job.id}.log")
        formatter = logging.Formatter("%(asctime)s - %(levelname)s - %(message)s")
        file_handler.setFormatter(formatter)
        logger.addHandler(file_handler)
    
    # 后续直接用该logger输出日志
    logger.info("任务启动")
    # ... 任务业务逻辑
    logger.error("执行出错:XXXXX")

核心要点:

  • 依托RQ的get_current_job()获取任务唯一ID,作为日志文件名和日志器的核心标识
  • 检查日志器是否已绑定处理器,防止重复输出
  • 提前创建日志目录,避免文件写入失败

2. 网络日志(集中收集)

核心是给每条日志打上任务唯一标识(如task_id),让收集器能基于该标识聚合同一任务的所有日志。以下是两种实用配置:

方案1:SysLogHandler(适配通用日志收集系统)

import logging
from logging.handlers import SysLogHandler
from rq import get_current_job

def your_task():
    job = get_current_job()
    logger = logging.getLogger(f"rq_task_{job.id}")
    logger.setLevel(logging.INFO)
    
    if not logger.handlers:
        # 指向集中收集器的IP和端口
        syslog_handler = SysLogHandler(address=("192.168.1.100", 514))
        # 日志格式强制携带task_id、worker_id等聚合字段
        formatter = logging.Formatter(
            "%(asctime)s - %(levelname)s - task_id=%(task_id)s - worker_id=%(worker_id)s - %(message)s"
        )
        syslog_handler.setFormatter(formatter)
        logger.addHandler(syslog_handler)
    
    # 打日志时通过extra参数传入标识字段
    logger.info("任务开始处理", extra={"task_id": job.id, "worker_id": job.worker_name})
    # ... 任务逻辑
    logger.warning("资源占用过高", extra={"task_id": job.id, "worker_id": job.worker_name})

方案2:自定义HTTP Handler(适配结构化日志收集)

如果需要发送JSON格式的结构化日志,可自定义HTTP处理器:

import logging
import requests
from rq import get_current_job

class HTTPLogHandler(logging.Handler):
    def __init__(self, endpoint):
        super().__init__()
        self.endpoint = endpoint
    
    def emit(self, record):
        try:
            log_data = {
                "task_id": record.task_id,
                "worker_id": record.worker_id,
                "level": record.levelname,
                "timestamp": self.formatTime(record),
                "message": record.getMessage()
            }
            requests.post(self.endpoint, json=log_data)
        except Exception:
            self.handleError(record)

def your_task():
    job = get_current_job()
    logger = logging.getLogger(f"rq_task_{job.id}")
    logger.setLevel(logging.INFO)
    
    if not logger.handlers:
        handler = HTTPLogHandler(endpoint="http://your-collector:8080/logs")
        logger.addHandler(handler)
    
    logger.info("任务执行中", extra={"task_id": job.id, "worker_id": job.worker_name})

核心要点:

  • 必须将task_id作为日志的核心元数据,收集器端需基于该字段建立索引或分组
  • 可选添加worker_id、机器IP等字段,辅助排查跨机器Worker的问题
二、集中式日志展示方案推荐

以下方案支持直接按task_id查询并展示完整任务日志,而非逐行事件搜索:

1. Loki + Grafana(推荐)

  • Loki负责日志收集与存储,天然支持按标签(如task_id)聚合日志
  • Grafana可创建专用面板,输入task_id即可一键展示该任务的所有日志按时间顺序拼接的完整文本
  • 资源占用远低于Elasticsearch,配置简单,适合中小规模场景

2. Fluentd + 自定义Web界面

  • Fluentd作为收集器,接收Worker日志后按task_id分组存储(如写入对应目录的文件或数据库)
  • 开发轻量Web服务,根据task_id查询对应日志文件/数据库记录,直接输出完整文本
  • 完全自定义展示样式,适合有特定界面需求的场景

3. Rsyslog + Logrotate + 简易查询页面

  • 用rsyslog收集网络日志,通过配置模板按task_id拆分存储到独立文件
  • 配合logrotate管理日志文件大小,避免磁盘溢出
  • 开发简单Web页面,输入task_id后读取对应日志文件内容展示
  • 最轻量化的方案,适合初期快速落地

内容的提问来源于stack exchange,提问作者persson

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.09 03:25:28