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

如何为Celery生成自定义应用指标并接入Prometheus监控

问题描述

我正尝试在Celery中生成自定义应用指标,并将这些指标采集到Prometheus中。目前使用celery-exporter导出Celery任务指标,这部分功能可开箱即用。但始终无法找到生成自定义应用指标(例如内部方法调用次数、方法执行延迟等)的方案。此前已在Flask应用中成功使用标准python prometheus client实现指标采集,但无法在Celery中复现该效果。
我了解到Celery监控能力与celery events绑定,但发布自定义事件后,Prometheus exporter中并未展示任何新增指标。
实现代码如下:

#!/usr/bin/env python3
import os

from celery import Celery
from celery.events import EventDispatcher
from celery.utils.log import get_task_logger
from kombu import Queue
from prometheus_client import Counter

QUEUE_INPUT = 'test_queue'
TASK_NAME = 'test_task'
QUEUE_ARGS = {'x-queue-mode': 'lazy'}
QUEUE_CONN_STR = os.environ['RMQ_MASTER_CONN_STR']

__LOG = get_task_logger(__LOG__)

# Define our consumer object
app = Celery(TASK_NAME)

# Define routes
app.conf.update({
    'broker_url': "{0}".format(QUEUE_CONN_STR),
    'task_routes': {"test_celery.test_task": {"queue": QUEUE_INPUT}},
    'task_serializer': 'json',
    'task_send_sent_event': True,
    'worker_send_task_events': True,
    'result_serializer': 'json',
    "broker_pool_limit": 20,
    "task_always_eager": False,
    "result_expires": 2,
    'worker_enable_remote_control': False,
    'worker_prefetch_multiplier': 1
})

app.conf.task_queues = [
    Queue(QUEUE_INPUT, queue_arguments=QUEUE_ARGS),
]

c = Counter('test_counter', 'Number of hits')
d = Counter('test_counter_1', 'Number of msgs')
dispatcher = EventDispatcher(app.connection())
print("Hi")
dispatcher.send("C_INIT")
c.inc()


@app.task(name="test_task", bind=True, max_retries=2)
def test_task(self, payload: dict):
    d.inc()
    dispatcher.send("START_PROCESSING")
    print(payload)
    dispatcher.send("END_PROCESSING")
核心问题

当前实现存在两个本质误区,导致自定义指标无法被采集:

  • celery-exporter默认仅解析Celery原生内置的任务事件(任务发送、启动、成功、失败、重试等官方定义的事件类型),不会自动识别自定义发送的C_INIT/START_PROCESSING/END_PROCESSING事件,更不会自动将这些事件转换为Prometheus指标,因此发了自定义事件看不到对应指标是预期行为。
  • Celery默认采用多进程prefork模型运行worker,每个worker进程内存空间完全隔离。直接在代码中初始化的Prometheus Counter存储在单个进程内存中,既没有解决多进程下指标聚合的问题,也没有暴露可供Prometheus抓取的HTTP端点,自然无法采集到对应指标。
解决方案

方案1:Worker内置指标端点(推荐,与Flask使用逻辑一致)

该方案不需要依赖celery events,和在Flask中使用prometheus client的逻辑完全一致,直接在worker进程中启动独立HTTP服务暴露指标即可,步骤如下:

  1. 配置prometheus client适配Celery多进程环境:
    • 启动worker前先创建一个有读写权限的空目录,比如/tmp/prometheus_multiproc
    • 配置环境变量PROMETHEUS_MULTIPROC_DIR指向该目录,解决多进程下指标统计错乱、丢失的问题
  2. 移除自定义EventDispatcher相关逻辑,直接在业务代码中定义需要的Counter、Histogram等指标,在对应逻辑触发时调用inc()、observe()等方法更新指标值。
  3. 绑定Celery worker初始化信号,在每个worker进程启动时,以守护线程的方式启动Prometheus HTTP服务,暴露/metrics端点供Prometheus抓取。

修改后的参考代码:

#!/usr/bin/env python3
import os
import threading
from prometheus_client import Counter, start_http_server, REGISTRY
from prometheus_client.multiprocess import MultiProcessCollector

from celery import Celery
from celery.utils.log import get_task_logger
from kombu import Queue

QUEUE_INPUT = 'test_queue'
TASK_NAME = 'test_task'
QUEUE_ARGS = {'x-queue-mode': 'lazy'}
QUEUE_CONN_STR = os.environ['RMQ_MASTER_CONN_STR']
# 自定义指标暴露端口,替换为机器上未被占用的端口即可
METRICS_PORT = 9102

__LOG = get_task_logger(__name__)

app = Celery(TASK_NAME)

app.conf.update({
    'broker_url': f"{QUEUE_CONN_STR}",
    'task_routes': {"test_celery.test_task": {"queue": QUEUE_INPUT}},
    'task_serializer': 'json',
    'task_send_sent_event': True,
    'worker_send_task_events': True,
    'result_serializer': 'json',
    "broker_pool_limit": 20,
    "task_always_eager": False,
    "result_expires": 2,
    'worker_enable_remote_control': False,
    'worker_prefetch_multiplier': 1
})

app.conf.task_queues = [
    Queue(QUEUE_INPUT, queue_arguments=QUEUE_ARGS),
]

# 定义自定义指标
c = Counter('test_counter', 'Number of hits')
d = Counter('test_counter_1', 'Number of msgs')

# 初始化多进程指标收集器
MultiProcessCollector(REGISTRY)

# 绑定worker启动后信号,启动指标HTTP服务
@app.on_after_configure.connect
def run_metrics_server(**kwargs):
    threading.Thread(
        target=start_http_server,
        args=(METRICS_PORT,),
        daemon=True
    ).start()
    # 初始化计数器
    c.inc()


@app.task(name="test_task", bind=True, max_retries=2)
def test_task(self, payload: dict):
    # 业务逻辑执行时更新对应指标
    d.inc()
    print(payload)

配置完成后,在Prometheus中新增抓取任务,目标地址为Celery worker所在IP加配置的METRICS_PORT,抓取路径为/metrics,即可采集到所有自定义指标,使用方式和Flask应用完全一致。

方案2:自定义扩展celery-exporter(不推荐)

如果一定要走celery events采集路径,需要自行二次开发celery-exporter:

  • 新增自定义事件监听逻辑,捕获发送的非原生事件类型
  • 为每个自定义事件定义对应的Prometheus指标结构
  • 在事件回调中更新对应指标值
    该方案需要自行维护定制化的exporter镜像,后续官方版本升级时需要同步合并代码,维护成本高,无特殊需求不建议使用。
注意事项
  • 使用prefork模式启动Celery worker时必须配置PROMETHEUS_MULTIPROC_DIR环境变量,否则会出现指标丢失、统计值不准的问题,每次重启worker前建议清空该目录下的残留文件。
  • 不要在模块全局初始化EventDispatcher,Celery启动时会fork多个子进程,全局初始化的连接在fork后会出现异常,导致事件发送失败。
  • 指标暴露端口不要对公网开放,避免业务监控信息泄露。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.01 23:03:53