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负载等指标:
- 先安装依赖库
psutil(用于获取系统进程数据):
pip install psutil
- 编写自定义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
相关产品推荐
相关产品推荐

