适用场景
使用 Celery + Redis/RabbitMQ 承担异步任务。队列中仍有大量待执行任务,但监控显示少数 Worker 忙碌、其余 Worker 空闲;或者短任务被长任务“压住”,整体延迟持续升高。本文以默认的 prefork 并发池和 RabbitMQ 为例,Redis broker 的排查思路相同。
现象描述
一次批量导入触发了数万条任务,其中约 5% 会访问慢接口,执行时间为 2~5 分钟,其余任务通常在 1 秒内完成。发布后出现以下现象:
- 队列深度持续增加,业务侧看到任务排队数分钟;
- 8 个 Worker 中有 2 个 CPU 较高,其余进程大多空闲;
- Flower 中部分 Worker 的
reserved数量很大,但active数量不多; - 重启 Worker 后短暂恢复,随后问题再次出现。
这类问题很容易被误判为 broker 性能不足。实际常见原因是 Worker 预取了超出其当前处理能力的任务,任务已从队列取走却还未执行,其他空闲 Worker 因此拿不到任务。
原理:预取为什么会造成“假积压”
Celery 的预取数量近似为:
prefetch_count = worker_concurrency × worker_prefetch_multiplier
默认 worker_prefetch_multiplier 通常为 4。一个并发度为 16 的 Worker 可能一次保留 64 个任务;当它先拿到一批慢任务后,后续短任务也可能被保留在该 Worker 内存中。它们不再可被其他 Worker 消费,表面上就形成“有 Worker 空闲、任务却不动”的不均衡。
reserved 是已从 broker 取出、等待执行的任务;active 是正在运行的任务。若 reserved 长期远大于 active,且各 Worker 差异明显,应优先检查预取。
排查步骤
1. 确认队列深度与消费者分布
RabbitMQ 可以先查看队列和消费者数量:
rabbitmqctl list_queues name messages_ready messages_unacknowledged consumers
messages_ready:仍留在 broker、尚未被取走的任务;messages_unacknowledged:已投递给消费者但还未确认的任务;consumers:正在消费该队列的连接数。
如果 messages_ready 不高而 messages_unacknowledged 很高,说明任务大概率已被 Worker 预取或正在执行;不能只看队列长度判断吞吐。
2. 对比各 Worker 的 active 与 reserved
在任意能访问 broker 的环境执行:
celery -A config inspect active
celery -A config inspect reserved
celery -A config inspect stats
将输出按 hostname 对比。排障时重点观察:某个 Worker 是否有大量 reserved,而另一个 Worker 的 active/reserved 都接近 0。stats 中的 pool.max-concurrency 则用于核对真实并发度,避免把容器副本数误当成单 Worker 并发。
3. 核对实际生效配置
不要只检查 Git 中的配置文件。容器启动参数、环境变量及 Helm values 都可能覆盖默认值:
celery -A config report | grep -E 'worker_prefetch_multiplier|task_acks_late|worker_concurrency'
ps -ef | grep '[c]elery worker'
如果启动命令包含 --prefetch-multiplier,以命令行值为准。还要记录 task_acks_late:它决定确认时机,和预取一起影响故障时任务会不会重新投递。
定位示例
事故现场有两个并发度为 8 的 Worker,预取倍数为 4。理论上每个 Worker 可保留 32 条任务。第一个 Worker 先取得 30 条慢任务及少量短任务;第二个 Worker 完成手头任务后空闲。虽然 broker 中已几乎没有 messages_ready,但并不代表任务已完成:大量任务仍处于第一个 Worker 的 reserved 列表中。
此时盲目扩容 Worker 不一定有效,因为新 Worker 也无法取得已被预取的任务。重启第一个 Worker 会使未确认任务回到队列,因此会出现“重启后恢复”的假象,但这不是根治方案。
修复方案
方案一:长短任务混跑时,将预取倍率设为 1
在 Celery 配置中显式设置:
# config/celery.py
app.conf.update(
worker_prefetch_multiplier=1,
task_acks_late=True,
task_reject_on_worker_lost=True,
)
worker_prefetch_multiplier=1 让每个执行槽最多预留一个任务,显著改善公平性。task_acks_late=True 表示任务执行成功后才确认;配合 task_reject_on_worker_lost=True,Worker 异常退出时未完成任务可重新投递。任务实现必须具备幂等性,例如用业务唯一键去重,避免重试造成重复扣款或重复发信。
也可用启动参数快速验证:
celery -A config worker -l INFO --concurrency=8 --prefetch-multiplier=1
方案二:按耗时拆分队列与 Worker
预取倍率设为 1 会降低极短任务的批量获取效率。更稳妥的长期方案是隔离任务类型:
app.conf.task_routes = {
'orders.tasks.sync_partner': {'queue': 'slow'},
'orders.tasks.send_notice': {'queue': 'fast'},
}
# 慢任务优先公平调度
celery -A config worker -Q slow -n slow@%h --concurrency=4 --prefetch-multiplier=1
# 短任务可适当提高预取,提升吞吐
celery -A config worker -Q fast -n fast@%h --concurrency=12 --prefetch-multiplier=4
隔离后,外部接口抖动不会占用通知、缓存刷新等短任务的执行槽。调整前先用运行时长的 P95/P99 数据划分任务,而不是只根据函数名称分类。
方案三:平滑发布与验证
先滚动重启一小部分 Worker,观察 15~30 分钟:
celery -A config inspect reserved
rabbitmqctl list_queues name messages_ready messages_unacknowledged consumers
验证标准包括:各 Worker 的 reserved 数量接近、短任务排队时间下降、messages_unacknowledged 不再长期单边集中。确认稳定后再完成全量滚动发布,避免同时中断所有消费者。
预防措施与注意事项
- 为任务记录排队耗时、执行耗时和重试次数;仅有成功/失败计数无法发现饥饿。
- 以队列、Worker hostname 和任务名三个维度监控
active、reserved与 unacknowledged。 - 慢任务设置软、硬时间限制,并对外部请求设置连接和读取超时,避免永久占槽。
- 使用延迟确认时,所有有副作用的任务都要设计幂等键或去重表。
- 将
worker_prefetch_multiplier、并发度和队列路由纳入发布变更单;这三项共同决定调度行为。
总结
当 Celery 出现“任务积压但 Worker 空闲”时,先区分任务是在 broker 中等待,还是已被某个 Worker 预取。通过队列的 messages_unacknowledged 与 Celery 的 reserved 输出,可以快速定位任务分配不均。对耗时差异大的任务,将预取倍率设为 1,并按队列隔离长短任务,通常比单纯扩容更有效、更可预测。
Discussion
评论