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

如何统计Celery活跃Worker数量?现有统计结果异常求优化

正确统计Celery活跃Worker数量的方法与思路

我用以下命令启动了10个Celery Worker:

celery -A worker.celery worker -l info -c 10

需求是统计活跃Worker的总数:当活跃数量不超过10时处理新任务,否则等待Worker空闲。我原本用这段代码统计:

import json
import subprocess

def get_celery_worker():
    bash_command = "celery -A worker inspect active -j"
    process = subprocess.Popen(bash_command.split(), stdout=subprocess.PIPE)
    output, error = process.communicate()
    output_string = output.decode("utf-8")
    output_json = json.loads(output_string)
    number_of_celery_worker = 0
    if len(list(output_json.values())[0]) == 0:
        pass
    else:
        for value in list(output_json.values())[0]:
            for v in value.values():
                if v == 'run_task':  # Here run_task is the worker name.
                    number_of_celery_worker += 1 
    return int(number_of_celery_worker / 2)  # Every task contains two run_task

但实际每秒启动1个任务(共10个)时,统计结果为0, 0, 1, 1, 2, 3, 4, 5, 6, 6,明显不符合预期。以下是正确的实现方法和思路:


方法一:直接使用Celery官方API(推荐)

避免调用命令行,直接通过Celery实例的控制接口获取活跃任务数,每个活跃任务对应一个正在占用的Worker进程:

from celery import Celery

# 初始化Celery实例,和你的worker配置一致
app = Celery('worker')
app.config_from_object('path.to.your.celeryconfig')  # 替换为你的配置加载方式

def get_active_workers_count():
    # 创建inspect对象
    inspect = app.control.inspect()
    # 获取所有节点的活跃任务
    active_tasks = inspect.active()
    
    if not active_tasks:
        return 0
    
    # 统计所有节点的活跃任务总数,即活跃Worker数量
    total_active = sum(len(tasks_list) for tasks_list in active_tasks.values())
    return total_active

方法二:优化命令行解析逻辑

如果必须用命令行方式,修正JSON解析逻辑——celery inspect active -j返回的每个任务条目对应一个活跃Worker,无需额外字段匹配和除法:

import json
import subprocess

def get_celery_worker():
    bash_command = "celery -A worker inspect active -j"
    # 捕获stderr避免错误输出干扰
    process = subprocess.Popen(
        bash_command.split(),
        stdout=subprocess.PIPE,
        stderr=subprocess.PIPE
    )
    output, error = process.communicate()
    
    # 处理命令执行失败的情况(比如Worker未响应)
    if process.returncode != 0:
        return 0
    
    output_string = output.decode("utf-8").strip()
    if not output_string:
        return 0
    
    output_json = json.loads(output_string)
    # 统计所有Worker节点下的活跃任务总数
    total_active = sum(len(tasks) for tasks in output_json.values())
    return total_active

问题分析

原代码统计结果异常的原因是对celery inspect active返回的JSON结构理解错误:

  • 每个任务条目是独立的对象,对应一个正在运行的任务(即一个被占用的Worker)
  • 原代码遍历任务的所有字段匹配run_task,导致重复计数,再除以2的操作进一步扭曲了统计结果

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 16:46:13