0. 上一章思考题参考答案
思考题 1:ignore_result=True是源头拦截——Worker 执行完根本不做「写结果」这一步,不产生键、不占网络与内存;result_expires=1800是事后回收——结果照样写,只是带 1800 秒 TTL,到期由 Redis 自动删除。所以「不需要结果」优先用 ignore_result,「需要结果但要防堆积」用 expires,两者可叠加。
思考题 2:chord 的 body 必须「精确触发一次」,实现上依赖计数器:每个 header 子任务完成后到 Backend 执行原子递增(backends/base.py的on_chord_part_return),最后一个完成者触发 body。计数器没地方放(无 Backend)→ chord 直接报错;Backend 在子任务执行期间挂掉 → 计数中断,body 无法被触发,chord_join_timeout超时后触发失败路径。这就是 chord 对 Backend 可用性的硬依赖(第 36 章源码细讲)。
1. 项目背景
三个月过去,任务量翻了三倍:短信、订单处理、报表导出、支付回调、邮件推送……全挤在默认的celery一条队列里。问题开始集中爆发:财务月报任务(单条跑 8 分钟)发起后,后面的短信任务被堵在队尾排 10 分钟,用户付完款收不到确认短信,客诉又上来了;运营群发邮件把队列打到 12 万条积压,支付回调被邮件淹没,退款迟了 1 小时,投诉升级到消协。
小周排查时发现一个诡异现象:有三个任务一直在PENDING,队列里根本没人消费——原来隔壁组起了个 Worker 专门跑报表,但生产者的路由没改,消息还是发进了celery队列,报表 Worker 订的却是report队列:生产发到东街,消费在西街等着,消息成了黑洞。
症结:所有任务共享一条 celery 队列 celery 队列 [短信|订单|报表|支付|邮件|...] ▲ 问题 1:慢任务堵队头,快任务饿死(队头阻塞) ▲ 问题 2:无法按业务分别扩容/限速 ▲ 问题 3:生产者路由与 Worker 订阅各写各的,错位成黑洞本章目标:拆出sms、order、report三条队列,让不同业务隔离消费、独立扩容、互不影响,并用压测证明隔离效果。
2. 项目设计
场景:客诉周会上,小周把「队头阻塞」画了出来。
小胖:一个队列多省事啊,现在要拆成三个,是不是又要改一堆代码?我觉得就是 Worker 开少了,多开几台机器就行了!
小白:小胖你想想食堂:只有一条队伍,前头的人点 20 个菜打包,后头买馒头的人也得等着——这叫队头阻塞。加机器没用,队还是一条,慢任务照样堵在所有人前面。我想问的是:Celery 里「消息进哪条队列」是谁决定的?生产者怎么声明队列?Worker 怎么知道该吃哪条队?
大师:两个问题,一个答案——路由是「生产者声明 + 消费者订阅」的成对契约。生产者侧,用task_routes把任务名映射到队列(celery/app/routes.py的路由规则),或者声明Queue/Exchange对象(celery/app/amqp.py):
task_routes={'orders.send_order_sms':{'queue':'sms'},'orders.close_order':{'queue':'order'},'orders.export_statement':{'queue':'report'},}消费者侧,Worker 启动时用-Q声明订阅:celery -A order_tasks worker -Q sms -c 4。缺一半就是灾难:生产路由改了、Worker 没加-Q,消息进了sms队列却没人消费,永远 PENDING——这就是你们看到的那三个「黑洞任务」。
技术映射:队列 = 食堂分窗口(套餐窗口、面食窗口、甜点窗口);task_routes = 点餐系统自动把单派到对应窗口;Worker
-Q= 后厨师傅站岗的窗口。系统派单了、师傅站错窗口,饭就永远出不来。
小白:那 Queue 对象和 Exchange 是什么关系?我听说有 direct、topic 好几种交换机,是不是越高档越好?
大师:Queue 是队列,Exchange 是「消息分发站」:生产者把消息交给 Exchange,Exchange 按 routing_key 把消息转进绑定的队列。Celery 里每个 Queue 默认配一个 direct 类型的 Exchange(名字与队列名相同),routing_key 也默认等于队列名——这就是「按队列名直投」,90% 的业务场景足够。topic 交换机适合「多字段模式匹配」(比如orders.cancel.*匹配所有取消类消息),fanout 适合广播。什么时候才上 topic?当你有多个 Worker 组要按同一批消息的不同维度订阅时才值得;为了 topic 而 topic,纯属自找麻烦。先 direct,压测证明不够了再 topic。
小胖:那我还有个馊主意——既然担心漏订阅,Worker 直接-Q sms,order,report全订了不就行了?一个都跑不掉。
大师(笑):这恰恰是「拆队列」最典型的返祖操作:全订等于没拆,慢任务还是会堵快任务,而且每个 Worker 都持有三个队列的连接与预取,资源翻倍,排障时根本说不清任务被谁吃了。正确姿势是按业务配并发:sms -c 4(短平快,4 并发够用)、order -c 8(核心链路)、report -c 1(重活串行,别互相抢)。隔离的意义就在于每个队列可以独立调并发、独立限速、独立告警。
技术映射:全订队列 = 一个人同时排三个窗口(手忙脚乱还堵别人);分工 = 每个窗口配数量合适的师傅。
3. 项目实战
3.1 环境准备
沿用环境(Redis Broker)。三队列隔离用 Redis 同样生效;RabbitMQ 上只是队列实体更直观(第 7 章已对比)。
3.2 分步实现
步骤 1:定义三条队列 + 路由表
目标:声明队列实体并绑定路由规则。
# order_tasks.pyfromceleryimportCelery app=Celery('order_tasks',broker='redis://localhost:6379/0')app.config_from_object('celeryconfig')app.conf.task_routes={'orders.send_order_sms':{'queue':'sms'},'orders.close_order':{'queue':'order'},'orders.gen_invoice':{'queue':'order'},'orders.export_statement':{'queue':'report'},}步骤 2:按队列声明 Worker(三个终端)
目标:三个队列分别由不同并发的 Worker 消费。
# 终端 A:短信队列(短任务,4 并发)celery-Aorder_tasks worker-Qsms-c4--loglevel=info--pool=solo-nsms-w1# 终端 B:订单队列(核心链路,8 并发)celery-Aorder_tasks worker-Qorder-c8--loglevel=info--pool=solo-norder-w1# 终端 C:报表队列(重任务,1 并发串行)celery-Aorder_tasks worker-Qreport-c1--loglevel=info--pool=solo-nreport-w1注意:不写
-Q的 Worker 只订默认celery队列——以后谁「忘了 -Q」,谁就只吃空气(第 15 章的经典故障)。
步骤 3:投递验证路由生效
目标:确认消息各自归位。
# route_demo.pyfromorder_tasksimportapp r1=app.send_task('orders.send_order_sms',args=[1])# → sms 队列r2=app.send_task('orders.close_order',args=[100])# → order 队列r3=app.send_task('orders.export_statement',args=['2026-08-23'])# → report 队列print(r1.id,r2.id,r3.id)运行结果(文字描述):三个任务分别被sms-w1、order-w1、report-w1打印的 received 日志接走,互不串线。
步骤 4:压测验证隔离效果
目标:报表队列灌 50 条 8 秒重任务,证明短信任务不受影响。
# flood_report.pyfromorder_tasksimportappforiinrange(50):app.send_task('orders.export_statement',args=[f'2026-08-{i%28+1}'])print("已灌入 50 条报表任务")python flood_report.py# 立即投递短信任务并计时python-c"import time; from order_tasks import app; t=time.time(); \ app.send_task('orders.send_order_sms', args=[9]); print('投递耗时', time.time()-t)"运行结果(文字描述):报表队列积压 50 条(report-w1 每秒只消化 1/8 条),但短信任务 2 秒内完成执行——慢任务堵在自己的队列里,快任务丝滑通过。对比拆队列前,短信要排在 50 条报表后面约 6 分钟。
步骤 5:演示「漏 -Q」的黑洞故障
目标:亲手复现一次路由错位,学会排障。
# 终端 A 停掉 order-w1 后,用默认命令重启(漏了 -Q order):celery-Aorder_tasks worker--loglevel=info--pool=solo-norder-w1-bug# 投递订单任务,观察其永远 PENDING排障三板斧:①celery -A order_tasks inspect registered确认 Worker 活着;② 看 Worker 启动日志的queues:行(没有 order);③ 补-Q order重启,任务被接走。排障套路固定化后,这个坑 10 分钟闭环。
步骤 6:在 RabbitMQ 上验证「队列—交换机—绑定」三实体
目标:路由不止停留在 Celery 配置,Broker 里要有实体可查(生产排查路由问题的最终依据)。
# 拆队列后(broker 切到 RabbitMQ),查看实体:dockerexecdocker-rabbit-1 rabbitmqctl list_queues name messages_readydockerexecdocker-rabbit-1 rabbitmqctl list_exchanges nametypedockerexecdocker-rabbit-1 rabbitmqctl list_bindings source_name destination_name routing_key运行结果(文字描述):能看到sms、order、report三条队列各自的消息数;每个队列对应一个同名 direct 交换机;绑定表里三条celery@... → 队列名的绑定,routing_key 即队列名。以后「路由对不对」的第一站不是看代码,而是看这张绑定表——代码改来改去,Broker 实体不会骗人。
3.3 可能遇到的坑及解决方法
| 坑 | 现象 | 解决 |
|---|---|---|
| 任务永远 PENDING | 生产路由指向新队列,Worker 没加-Q | 检查 Worker 启动日志queues:行;补-Q重启 |
| 拆队后全订队列(-Q 全选) | 隔离失效,慢任务依旧堵快任务 | 按业务配独立 Worker 与并发 |
task_routes规则不生效 | 路由写错或任务名不匹配 | 规则支持前缀通配:'orders.*': {'queue': 'order'};改后重启 |
| Redis 上看不到队列名 | 用 LLEN 查队列查不到 | Redis 的队列就是celery键名变化,用KEYS或 RabbitMQ 管理台 |
| 报表 Worker 忙死,其余队列正常 | 报表任务堆积 | report 队列单发压测 + 限速(rate_limit,第 21 章) |
3.4 完整代码清单与测试验证
清单:order_tasks.py(路由表)、route_demo.py、flood_report.py+ 三个 Worker 启动脚本。队列规划表(沉淀 Wiki):
| 队列 | 业务 | 并发 | 时限要求 | 告警积压阈值 |
|---|---|---|---|---|
| sms | 短信通知 | 4 | 秒级 | 1000 |
| order | 订单核心 | 8 | 秒级 | 5000 |
| report | 报表导出 | 1 | 分钟级 | 100 |
队列规划表本身就是路由的「权威字典」:以后新增任务,第一件事不是写代码,而是来这张表认领队列;没有合适队列就开会加行,队列的增减是可治理的变更,不是随手起名。表随代码库一起版本化,路由表变更即表变更。
测试验证:
# tests/test_routes.pyfromorder_tasksimportappdeftest_route_table_covers_all_orders():routes=app.conf.task_routesassertroutes['orders.send_order_sms']['queue']=='sms'assertroutes['orders.close_order']['queue']=='order'assertroutes['orders.export_statement']['queue']=='report'deftest_route_target_queue_defined():# 路由指向的队列必须实际声明,避免「路由到空气」queues={q.nameforqinapp.amqp.queues.values()}assert{'sms','order','report'}<=queuespython-mpytest tests/test_routes.py-v# 2 passed4. 项目总结
4.1 优点 & 缺点
| 维度 | 多队列隔离(本章方案) | 单队列大杂烩 |
|---|---|---|
| 队头阻塞 | 快慢任务互不影响 | 慢任务堵死全体 |
| 弹性伸缩 | 按队列独立扩缩容 | 只能整体加机器 |
| 可观测 | 按队列积压/速率分别告警 | 一团数据无法定位 |
| 运维成本 | 多队列需规划路由与订阅 | 零配置 |
| 缺点 1 | 路由与 -Q 错位产生黑洞任务 | 无此问题 |
| 缺点 2 | 队列划分不合理时频繁调整 | —— |
4.2 适用场景
- 适用:① 快慢任务混合(通知 vs 报表);② 不同业务需独立限速/扩容;③ 核心链路与边缘任务要隔离保护(支付回调 vs 邮件营销);④ 多团队共用一个 Celery 集群(每个团队自己的队列与 Worker,互不打扰)。
- 不适用:① 任务量小且同质(一个队列够了);② 队列数量失控(20 条队列 3 个人维护 = 过度设计);③ 需要跨队列全局排序的场景(队列隔离天然放弃全局顺序);④ 消息量极小却追求「一个 Worker 吃所有队列」的极简部署。
4.3 注意事项
- 路由规则与
-Q订阅是成对契约,改一侧必须改另一侧;代码评审双人签字。 - 队列名要有命名规范(业务域小写 + 用途),乱起名三个月后没人知道
q1是啥。 - 先 direct 后 topic:模式匹配只有真实需求出现时才引入。
- 每个队列配置独立的积压告警阈值——隔离的目的就是让告警「指哪打哪」。
- 队列实体(Broker 侧)与配置(代码侧)要定期对账:RabbitMQ 绑定表与 task_routes 逐条比对,防「配置删了、实体还在」的僵尸队列。
4.4 常见踩坑经验(3 个生产故障)
- 故障:支付回调延迟 1 小时,客诉升级。根因:营销邮件把单队列打到 12 万积压,回调排在后面。对策:拆分
pay队列并隔离 Worker。教训:核心链路必须有自己的队列。 - 故障:三个报表任务永远 PENDING。根因:生产者路由改了
report队列,Worker 启动命令没加-Q。对策:启动脚本模板化 +inspect巡检。教训:路由成对变更,缺一即黑洞。 - 故障:拆队列后短信反而更慢。根因:短信 Worker 顺手
-Q sms,report全订,报表任务抢占了短信 Worker。对策:严格一 Worker 一队列组。教训:全订等于没拆,隔离是纪律不是配置。
4.5 思考题
- 生产者把消息发进了
sms队列,但没有任何 Worker 订阅它——这条消息的最终命运是什么?它会在什么时候被「发现」? - 三条队列的 Worker 都配置了
worker_prefetch_multiplier=4,报表 Worker 并发 1、短信 Worker 并发 4,各自会预取多少条?预取数对「饿死其他 Worker」有什么影响?(提示:第 18 章)
答案见第 10 章开头的「上一章思考题参考答案」。
延伸阅读与资源
Java 工程师进阶:从 JVM 生产排障到OpenJDK原理
NumPy 从入门到生产落地:全链路实战指南(科学计算/向量化)
Redis 8 实战精讲:从 CRUD 到源码,构建高可用缓存系统
Redis 实战修炼与原理进阶
Python 3实战精进:从脚本到高并发订单引擎
python入门:Rquests从菜鸟脚本到企业级SDK的网络实战圣经
Milvus向量数据库实战修炼:从 0 到 1精通向量检索与生产落地
MongoDB 实战进阶与内核修炼
后端工程师的 AI 转型第一课:Ollama 与私有化大模型实战
10倍开发者的 Dify 魔法书:从零构建全栈 AI 应用
后端工程师转型AI第一课-Ollama 与私有化大模型实战
大型语言模型(LLM) vLLM 高性能推理落地实战
Agent开发之LlamaIndex 实战修炼与源码进阶
大语言模型Transformers 实战修炼与源码剖析