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

Celery中task.info是什么?能否修改它采集Worker更多运行数据?

回答

首先明确:完全可以修改或扩展task.info来采集更多Worker运行数据。task.info的数据实际上来自Worker发送的心跳包(heartbeat payload),Celery官方文档确实没详细提及这部分的自定义扩展,但通过重写Worker或Task类的相关方法就能实现需求。

两种实现方案

方案1:自定义Worker心跳包,全局扩展所有任务的info数据

Worker会定期向Broker发送心跳,默认包含hostname、pid等基础信息。我们可以重写Worker的heartbeat_payload方法,在心跳数据中加入运行时间、内存占用、CPU负载等指标:

  1. 先安装依赖库psutil(用于获取系统进程数据):
pip install psutil
  1. 编写自定义Worker和Celery应用:
from celery import Celery
from celery.worker import Worker
import psutil
import time

class CustomWorker(Worker):
    def heartbeat_payload(self):
        # 获取默认心跳数据
        payload = super().heartbeat_payload()
        # 获取当前Worker进程的实例
        worker_process = psutil.Process(self.pid)
        
        # 添加自定义监控指标
        payload.update({
            # Worker运行时长(秒)
            'uptime': round(time.time() - worker_process.create_time(), 2),
            # 内存占用(MB)
            'memory_usage_mb': round(worker_process.memory_info().rss / (1024 ** 2), 2),
            # CPU使用率(%)
            'cpu_percent': worker_process.cpu_percent(interval=0.1),
            # CPU累计耗时(用户态+内核态)
            'cpu_times': {
                'user': round(worker_process.cpu_times().user, 2),
                'system': round(worker_process.cpu_times().system, 2)
            }
        })
        return payload

# 初始化Celery应用
app = Celery('tasks', broker='redis://localhost:6379/0')
# 替换默认Worker类为自定义类
app.Worker = CustomWorker

# 示例任务
@app.task
def add(x, y):
    return x + y

启动Worker后,访问task.info就能看到新增的所有指标。

方案2:自定义Task类,为单个/部分任务扩展info数据

如果不需要全局所有任务都扩展,也可以单独重写Task类的info属性,仅针对特定任务生效:

from celery import Celery, Task
import psutil
import time

class CustomTask(Task):
    @property
    def info(self):
        # 获取默认的info基础数据
        base_info = super().info
        # 获取当前Worker的进程ID
        worker_pid = self.request.worker_pid
        worker_process = psutil.Process(worker_pid)
        
        # 扩展自定义指标
        base_info.update({
            'worker_uptime': round(time.time() - worker_process.create_time(), 2),
            'memory_rss_mb': round(worker_process.memory_info().rss / (1024 ** 2), 2),
            'cpu_load': worker_process.cpu_percent(interval=0.05)
        })
        return base_info

# 初始化Celery应用
app = Celery('tasks', broker='redis://localhost:6379/0')
# 设置默认Task类为自定义类(也可以在单个任务上指定base=CustomTask)
app.Task = CustomTask

# 示例任务
@app.task
def multiply(x, y):
    return x * y

注意事项

  • psutil是跨平台库,支持Windows/Linux/macOS,能稳定获取进程级别的系统数据。
  • CPU使用率采集的interval参数:设为0会返回最近的平均使用率,设置小数值(如0.1)会更实时,但会带来轻微的性能开销,可根据需求调整。
  • 确保自定义指标都是可序列化的基本类型(如数值、字符串、字典),避免直接返回psutil的对象,否则可能导致Broker序列化失败。
  • task.info的数据来自Worker的实时心跳,仅对运行中的Worker有效;如果Worker已停止,info可能无法获取到最新的扩展数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 14:17:37