监测 Celery 队列、worker 和任务事件

原文:Celery 文档贡献者,Monitoring and Management Guide。本文依据 2026-10-05 读取的 Celery 5.6.3 稳定文档全文翻译整理,并对原文的安全默认值、旧配置名称和事件处理示例作明确修订。命令与 Python 代码仅经静态审查,未连接 broker、worker 或 Flower,也未执行清空队列等管理操作。

监控 Celery 需要同时看三个位置:broker 中等待投递的消息,worker 已接收和正在执行的任务,以及系统发出的任务/worker 事件。它们描述的时间点和含义不同。一个队列长度很小的系统仍可能有很多任务卡在 worker 上;一次 inspect 没收到回复,也不足以判定 worker 已死亡。

Celery 观测示意:生产者把任务交给 broker,broker 投递给 worker 的 reserved、scheduled 和 active 状态;独立事件流供 Flower 或 Receiver 观察,不能把队列长度等同于执行中任务数量。
图:未完纪原创的观测层次示意,不是运行中集群截图。

从 inspect 命令读 worker 的当前状态

通过 celery --help 和具体子命令的 --help 查看安装版本支持的选项。下文沿用示例应用名 proj;运行 CLI 会导入这个应用,因此先确认目标应用和连接配置,不能在未知配置下试命令。

celery -A proj status
celery -A proj inspect active
celery -A proj inspect reserved
celery -A proj inspect scheduled
celery -A proj inspect registered
celery -A proj inspect revoked
celery -A proj inspect stats
命令或状态 观察对象 不能据此断言什么
status 在查询窗口内响应的活动节点。 未回复不等于进程一定死亡。
active worker 当前正在执行的任务。 不包括全部 broker 待投递消息。
reserved 已经预取、等待执行的任务,通常不含 ETA 任务。 不等于尚未被 worker 接收的队列积压。
scheduled worker 持有、带 ETA/countdown 的延时任务。 不是 Celery Beat 全部周期计划的列表。
registered worker 已注册的任务类型。 不代表这些任务都执行过。
revoked worker 记住的撤销任务标识。 不等于持久、完整的历史审计。
stats worker 统计。 不保证所有节点在同一瞬间采样。

若知道任务 ID,可以请求持有该任务的 worker 返回信息:

celery -A proj inspect query_task e9f6c8f0-fec9-4ae8-a8c6-cf8c8451d4f8
celery -A proj inspect query_task id1 id2 id3

源文说明,这个查询面向 active 或 reserved 等 worker 当前持有的任务。任务已经结束、转移,或节点没及时响应时,不一定能得到结果。若要读取结果后端,可用:

celery -A proj result -t tasks.add 4e196aa4-0141-4601-8138-7aa33db0f577

没有自定义结果后端时,原文说明可省略任务名。结果是否仍存在还取决于是否保存结果、后端可用性和过期策略,不能把监控事件存储当成结果后端。

inspect 和 control 默认面向所有 worker。可以用目标节点和超时缩小查询范围:

celery -A proj inspect --timeout=5 -d w1@example.com,w2@example.com reserved

无回复可能来自网络延迟、节点繁忙、控制通道不可达或超时过短,应结合 broker 状态和其他探针判断。超时增加只扩大等待窗口,不会把分布式查询变成强一致快照。

把查看状态和修改集群分开

celery shell 进入交互式 Python 环境,包含当前应用和已知任务;可选择 IPython、bpython 或普通 Python,并用 --without-tasks 避免自动加入任务名。这是管理入口,不是只读仪表盘,调用任务或应用函数可能产生业务副作用。

启用或关闭 worker 任务事件也是控制操作,会改变事件流量:

celery -A proj control -d w1@example.com enable_events
celery -A proj control -d w1@example.com disable_events

原指南还列出两类容易被误当成诊断步骤的命令,本文保留其语义,但不提供可顺手执行的破坏性示例:

  • purge 永久删除配置队列中的消息,没有撤销;-Q 可指定队列,-X 可排除队列。它不用于“修复监控”。原页使用旧配置名 CELERY_QUEUES,新代码应核对安装版本的 task_queues 配置。
  • migrate 在 broker 之间迁移任务,源文标为实验性,要求先备份。它是数据迁移,不是只读检查,应另行制定验证、重复处理和恢复方案。

Flower:既是监控界面,也是管理接口

Flower 使用 Celery Events 展示任务进度与历史、参数、开始时间、运行时长、图表和统计。它还能查看 worker 配置,调整进程池大小和自动伸缩,改变消费队列、限速和时间限制,撤销或终止任务,以及关闭或重启 worker。HTTP API 也覆盖查询、执行任务和多项管理行为。因此,暴露 Flower 不只是暴露一个状态页。

原文用 pip install flower 安装,再通过 celery -A proj flower 启动。它给出浏览器访问地址 http://localhost:5555,但这不代表服务只监听本机。Flower 配置文档说明 address 默认是空字符串,可监听全部接口;认证默认也需要另行启用。

下面是本文补充的本机观察配置示例。认证值由部署环境安全注入,缺失时配置加载失败;示例没有硬编码口令。实际部署仍要核对 Flower 版本并配置可靠的访问控制:

# flowerconfig.py;配置文件本身是 Python,必须来自可信来源。
import os

address = '127.0.0.1'
port = 5555
basic_auth = [os.environ['FLOWER_BASIC_AUTH']]
# FLOWER_BASIC_AUTH 由受保护的部署配置提供一个 user:password 值。
celery -A proj flower --conf=flowerconfig.py

不要把真实 broker 密码写在命令行 URL、shell 历史或文章示例中。原文的 guest:guest 只是本机示例凭据,不是生产配置;本文改用已配置的 proj 应用。需要远程访问时,应经受控反向代理或隧道,并配置 TLS、认证、网络限制与最小权限。HTTP Basic 凭据本身不提供加密。

Flower 中的参数、返回值、异常和配置可能含秘密或个人数据;认证不是数据脱敏的替代。根据业务设置可见字段、日志保留和访问权限。原文提到 OpenID,当前支持的认证提供方与相关选项应查所用 Flower 版本。配置文件和插件都是可执行 Python,也必须按代码来源管理。

终端事件监控与快照入口

celery events 是基于 curses 的终端监视器,可以查看任务、worker、结果和 traceback,也包括限速、关闭 worker 等管理操作。原文说明它始于概念验证,通常优先使用 Flower。命令行还有两个有用入口:

celery -A proj events
celery -A proj events --dump
celery -A proj events --help

--dump 会把事件写到标准输出,可能泄露任务参数与结果;只能在明确的数据边界内使用,不能默认把完整流量贴到日志平台或工单。Camera 则允许周期性观察聚合后的状态,后文给出例子。

RabbitMQ:区分 ready 与 unacknowledged

RabbitMQ 自带的 rabbitmqctl 能查询队列、交换机、绑定、长度、内存和消费者。下面保留源文的队列检查命令;若不用默认虚拟主机,应增加正确的 -p:

rabbitmqctl list_queues name messages messages_ready messages_unacknowledged
rabbitmqctl list_queues name consumers
rabbitmqctl list_queues name memory
# 自定义虚拟主机示例:
rabbitmqctl list_queues -p my_vhost name messages_ready messages_unacknowledged

messages_ready 是已经发送但尚待投递的消息数;messages_unacknowledged 是已投递却尚未确认的消息数;messages 是二者之和。未确认消息可能已预取或正在处理,但确认时机受任务和 worker 配置影响,不能直接等同于 active 数量。早确认任务执行期间可能已不在 broker 的未确认计数中。

consumers 是 broker 看到的消费者数量,不宜直接当成 Celery 机器或进程数量;一个 worker 的连接和消费者结构可能不同。memory 表示队列内存占用。-q 可减少部分输出,便于受控解析,但字段与权限仍应按 RabbitMQ 版本检查。

Redis:先知道队列键,再观察长度

使用 Redis 作为 broker 时,可以查询已知队列列表长度:

redis-cli -h HOST -p PORT -n DATABASE_NUMBER LLEN QUEUE_NAME

默认队列名通常是 celery。空列表的键会被 Redis 删除,所以键不存在时 LLEN 返回 0;这不意味着 worker 没有执行任务,也不能仅凭一个键否定优先级队列或其他路由队列中的积压。

安全修订:原文用 KEYS * 列出所有键。该命令会遍历数据库,大库中可能阻塞服务;日常监控应优先使用已知队列名,确需发现键时再作有界、受控的增量扫描:

# 从返回的游标继续分页;COUNT 是提示值,不是总量上限。
redis-cli -h HOST -p PORT -n DATABASE_NUMBER SCAN 0 MATCH 'celery*' COUNT 100

这个模式只是教学例子,不能覆盖所有自定义队列;SCAN 也不是一致性快照,可能重复返回键。不要把扫描结果逐个盲目执行 LLEN,其他用途的键可能不是列表。专用 Redis 实例或数据库可以减少键名混淆,但 Redis Pub/Sub 事件不按数据库号隔离,所以不同逻辑数据库不等于不同 Flower 事件域或安全边界。

Prometheus 与原文中的 Munin 插件

Prometheus 并非 Celery 核心内置监控,但可以通过 Flower 暴露指标,再结合其提供的 Grafana 仪表盘观察任务、worker 等统计。具体启用方式见 Flower 的 Prometheus 集成文档。指标端点也需要网络与访问控制,不能因没有任务正文就默认公开。

原文还列出 rabbitmq-munin、celery_tasks 和 celery_tasks_states。后两者注明依赖 celerymon,属于文档沿用的历史生态信息。本文保留链接,不把这些插件宣称为已验证兼容当前 Celery 的推荐部署。

Camera 保存的是状态快照,不是每一条事件历史

一个 worker 也可能产生大量事件。app.events.State 在内存中把事件合并成任务与 worker 状态;继承 Polaroid 的 Camera 可以按固定间隔处理快照,比如写入数据库。快照可以帮助保存状态变化,但没有保证保留期间每一条原始事件,断线丢失、内存淘汰和快照之间的中间状态都要单独考虑。

原文的 DumpCam 打印全部 worker 和 task 信息。以下修订版只输出聚合计数,避免在示例中默认记录参数、结果和 traceback:

# myapp.py
import json
from celery.events.snapshot import Polaroid

class CountCam(Polaroid):
    clear_after = True

    def on_shutter(self, state):
        if not state.event_count:
            return
        print(json.dumps({
            'event_count': state.event_count,
            'task_count': state.task_count,
            'worker_records': len(state.workers),
        }))
celery -A proj events -c myapp.CountCam --frequency=2.0

worker_records 是状态对象里保存的 worker 记录数,不是保证在线的 worker 数。clear_after=True 会在刷新后按 Camera/State 的清理语义处理状态,因此不能拿这里的计数当成永不丢失的累计审计。需要程序化启动时,结构如下:

from proj import app
from myapp import CountCam

def main(app, frequency=2.0):
    state = app.events.State()
    with app.connection() as connection:
        receiver = app.events.Receiver(
            connection, handlers={'*': state.event}
        )
        with CountCam(state, freq=frequency):
            receiver.capture(limit=None, timeout=None)

if __name__ == '__main__':
    main(app)

这是常驻事件消费者示意,没有加入重连、持久化事务、退出信号、背压或完整监控自监测。生产实现需要另外设计,不能仅把无限捕获循环交给后台运行就认为可靠。

实时处理:失败事件需要此前的任务信息

实时监控由 Receiver、事件处理器,以及可选 State 组成。State 能根据后续事件合并任务字段,并跟踪心跳与时间信息。原文第一个失败告警示例通过 '*': state.event 接收其他事件;第二个示例只接收 task-failed,却仍假定已经知道任务名。这是一个真实的示例缺口:任务名通常来自 task-received,不能靠失败事件凭空补出。

下面的修订版保留兜底接收,防御缺失字段,并只输出任务 ID 和名称。JSON 编码也避免任务名中的换行直接伪造日志行。即便接收全部相关事件,监控器中途启动或断线后仍可能缺失 received 事件,因此名称允许为 null:

import json
from proj import app

def monitor(app):
    state = app.events.State()

    def failed(event):
        state.event(event)
        task_id = event.get('uuid')
        task = state.tasks.get(task_id) if task_id else None
        print(json.dumps({
            'event': 'task-failed',
            'task_id': task_id,
            'task_name': getattr(task, 'name', None),
        }, ensure_ascii=False))

    with app.connection() as connection:
        receiver = app.events.Receiver(connection, handlers={
            'task-failed': failed,
            '*': state.event,
        })
        receiver.capture(limit=None, timeout=None, wakeup=False)

if __name__ == '__main__':
    monitor(app)

原文 wakeup=True 会向 worker 发信号,促使其发送心跳,从而更快显示节点;这里改成 False,避免示例默认广播唤醒。需要即时发现时可在已授权的运维场景调整。State 仍可能在内存里持有事件中的敏感字段,应限制监控进程和转储的访问。上述代码没有硬编码 broker 凭据,也不执行或反序列化来自事件字段的 Python 表达式;这不构成对实际消息序列化配置或整个服务的安全认证。

任务事件字段参考

以下保留原指南列出的主要字段,便于设计事件接收器;它们是事件数据,不是可直接调用的业务函数:

事件 字段 含义
task-sent uuid, name, args, kwargs, retries, eta, expires, queue, exchange, routing_key, root_id, parent_id 生产者发布任务时产生,需启用 task_send_sent_event。
task-received uuid, name, args, kwargs, retries, eta, hostname, timestamp, root_id, parent_id worker 收到任务,通常包含后续关联所需的名称。
task-started uuid, hostname, timestamp, pid worker 即将执行任务。
task-succeeded uuid, result, runtime, hostname, timestamp 成功完成;runtime 的口径是送入进程池到结果处理回调,不是从发布开始的端到端耗时。
task-failed uuid, exception, traceback, hostname, timestamp 执行失败;异常和堆栈可能含敏感内容。
task-rejected uuid, requeue 被 worker 拒绝,可能重新入队或进入死信队列。
task-revoked uuid, terminated, signum, expired 被撤销,可能由多个 worker 报告;terminated 表示进程被终止,signum 是信号,expired 表示到期。
task-retried uuid, exception, traceback, hostname, timestamp 此次失败后将重试,不等于最终业务失败。

实际统计要按任务 ID、重试和生命周期去重,处理乱序、重复及缺失事件。若需要审计级事实,应把业务状态和持久存储纳入设计,不能仅靠短期内存 State 推断完整执行史。

worker 事件与心跳

worker-online 与 worker-offline 使用 hostname, timestamp, freq, sw_ident, sw_ver, sw_sys 等字段;分别表示接入 broker 与离线。worker-heartbeat 还包括 active 和 processed,用于描述正在执行任务数和处理总数。

freq 的单位是秒,sw_ident 是软件标识,sw_ver 是版本,sw_sys 是操作系统。原页仍保留“每分钟发送、两分钟无心跳视为离线”的叙述。它不能当成所有安装版的固定间隔;应核对 worker 实际配置、事件携带的 freq 与监控器的过期策略。节点停顿、网络分区和监控端自身积压都会影响心跳判断。

高级配置:Mailbox 的持久与独占选项

Celery 内部使用 kombu.pidbox.Mailbox 发送控制与广播命令。原文新增了 Kombu 5.6.0 的配置说明:durable 默认 False,启用后相关交换实体按其持久语义跨 broker 重启保留;exclusive 默认 False,用于独占语义。两者不能同时为 True,文档说明这个组合会报错。

原页链接到 event_queue_durable 与 event_queue_exclusive 配置。部署时应分别核对事件队列、控制 Mailbox 和所用 broker 的实体语义,不要因为名称相近就把所有消息都理解为已经持久化,也不要把独占当成用户授权机制。

来源、许可与改动记录

Celery User Manual 的版权页署名 Ask Solem,Copyright © 2009–2016, Ask Solem;现由 Celery 文档贡献者维护。原文文档采用 CC BY-SA 4.0,Celery 软件本身另采用 BSD 3-Clause。本文作为中文翻译整理作品按 CC BY-SA 4.0 提供,保留来源和改动标记;本篇原创配图也按 CC BY-SA 4.0 提供。

本文的改动包括:重排观测流程,修正只订阅失败事件却期待任务名的缺口,用受控 SCAN 替代原文的全库 KEYS 建议,移除默认口令示例,补充 Flower 监听/认证、敏感字段、分布式快照和过时心跳说明,并将破坏性管理命令单独解释。参考:原指南、文档版权页、Flower 配置。静态审查未发现这些修订示例中的注入执行入口,不代表相关系统没有其他漏洞。

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容