第27章:Celery、Django / FastAPI 集成与事务一致性
0. 上一章思考题参考答案思考题 1以max_retries3为例任务失败 3 次后最终失败——task_retry触发 3 次每次重试前task_failure只触发 1 次第 4 次执行最终失败时task_postrun每次都触发含重试中的成功/失败共 4 次。所以耗时统计放在 task_postrun 会把重试也计入每次执行都记一次——「总执行时间」与「单次执行时间」要区分开前者累加、后者按request.retries过滤retries0 跳过或单独打点。思考题 2不能。Celery 的信号是观察型fire-and-observe返回值不被消费、不能中止发布before_task_publish只能修改 headers/body 的可变内容如注入 trace_id不能「取消这条消息」。真正要拦截投递必须在业务代码层判断或改路由让消息进死信队列——信号是「看」不是「拦」。1. 项目背景订单平台接入 FastAPI 后上线第一周就出了一个「幽灵单」事故用户在 App 上看到「支付成功」但客服查订单库没有这条订单——而短信、审计任务都执行了。排查发现下单接口的代码是「先建订单写库→ 立刻发任务」——数据库事务还没提交任务就先发出去了任务执行时读库读到的是「未提交的订单」不更糟——事务随后回滚了订单没了但任务已经消费执行完短信照发、审计照记。「库有单、任务没发出」和「库没单、任务已执行」是同一枚硬币的两面根子都是任务投递时机没有和数据库事务对齐。经典时序问题 BEGIN INSERT INTO orders ... # 订单入库未提交 send_order_sms.delay(...) # ❌ 任务已投递事务还没提交 COMMIT # 若这里回滚 → 库没单、短信已发本章目标FastAPI 下单接口实现「先提交事务、再投递任务」再补一条Outbox 扫描兜底——即使投递失败也能靠扫描重发杜绝「库有单、任务没发出」。2. 项目设计场景幽灵单事故复盘小周把「先建单再发任务」的代码放大到屏幕。小胖这有啥难的事务提交完再delay()不就行了代码顺序调一下两分钟的事你们非要搞什么 Outbox又是新名词怕不是要加钱小白小胖你说得轻巧——「调顺序」在两个地方会翻车① 提交后投递也可能失败Broker 抖动第 16 章思考题② 多个任务要跟事务联动时如「下单 扣库存 审计」三件套第 16 章顺序调了但失败语义没变。我查了 Django 生态有transaction.on_commitFastAPI 没有内置事务钩子——所以我想先问on_commit 和「普通顺序」到底差在哪大师差在语义正确性。on_commit(callback)Django是「事务真正提交后才执行回调」——它在事务成功提交的瞬间触发回滚则完全不触发。而「代码顺序」是假的你以为写库和发任务有先后但事务的 COMMIT 发生在函数返回之后由 ORM/框架统一提交你代码里的「先」根本没对齐「提交时刻」。所以正确姿势① Django 用transaction.on_commit(lambda: send_order_sms.delay(...))② FastAPI 手写「提交点」——接口层拿到db.commit()成功之后再投递任务③ 更严谨Outbox 模式——把「待发消息」和业务数据同一个事务写进 outbox 表事务提交后由独立投递器转发。技术映射事务 点餐的「最后确认付款」on_commit 「付款成功小票打出」之后的事普通顺序 「先喊厨师做菜、后付款」——菜做了钱没付回滚菜就白做幽灵单。小白那 Outbox 完整流程是什么「扫描兜底」具体兜什么大师Outbox 三步① 写——业务事务里同时INSERT outbox(event_type, payload, status)与订单同事务原子② 发——事务提交后投递器把 outbox 记录转成任务消息投递投递成功标记sent③ 扫——定时任务扫描「未 sent 且超时」的记录重投幂等键保证不重复投递。它兜的是投递瞬间 Broker 挂了 / 投递器进程崩溃——反正 outbox 记录和订单在一个事务里订单在记录就在扫描总能补发。这是「库有单、任务没发出」的最终解药也是第 40 章自研平台的核心件。小胖那测试呢我听说测试里task_always_eagerTrue能让任务同步执行是不是测试就稳了大师task_always_eager是双刃剑celery/contrib/testing/也在帮你管理它好处是测试里任务同步执行、能直接断言返回值幻觉是——① eager 模式不经过 Broker消息序列化/路由/队列订阅的错误全部测不到② eager 模式不经过 Worker重试、晚确认、崩溃重投这些运行时语义全是假的③ eager 模式任务在调用方进程跑跟真实「跨进程」差了十万八千里。所以正确分层单元测试用 eager/mock 测逻辑集成测试必须真实 Broker Worker 进程第 28 章完整方法论。技术映射task_always_eager 考试「开卷但题目都是背诵题」——方便但测不出真实能力真本事跨进程/序列化/重试还得靠「闭卷实战」真实 Broker。3. 项目实战3.1 环境准备沿用环境Redis Broker Backend。新增依赖fastapi已装uvicornOutbox 演示用 sqlite。3.2 分步实现步骤 1FastAPI 下单接口——事务提交后再投递目标消灭「先发任务、后提交」的时序问题。# web_order.pyimportsqlite3,time,threadingfromfastapiimportFastAPIfromorder_tasksimportapp,send_order_sms,audit_log apiFastAPI()ORDER_DBorders.dbdef_init_db():connsqlite3.connect(ORDER_DB)conn.execute(CREATE TABLE IF NOT EXISTS orders (id INTEGER PRIMARY KEY AUTOINCREMENT, mobile TEXT, created_at TEXT DEFAULT CURRENT_TIMESTAMP))conn.commit()conn.close()_init_db()api.post(/api/orders)defcreate_order(mobile:str13800000000):# —— 业务事务建订单 ——connsqlite3.connect(ORDER_DB)try:curconn.execute(INSERT INTO orders(mobile) VALUES (?),(mobile,))order_idcur.lastrowid conn.commit()# ★ 事务提交成功才进入投递阶段exceptException:conn.rollback()raisefinally:conn.close()# —— 提交后投递此刻订单在库里「看得见」 ——r1send_order_sms.apply_async(args[order_id,mobile])r2audit_log.apply_async(args[order_id,create])return{order_id:order_id,tasks:[r1.id,r2.id]}运行结果文字描述接口创建订单 → 事务提交 → 任务投递模拟「回滚」场景在 commit 前 raise任务不会投递——幽灵单的根因在接口层被消灭。步骤 2Outbox 表 同事务写入目标把「待发消息」与订单放进同一个事务。# outbox.pyimportsqlite3,json DBoutbox.dbdef_init():connsqlite3.connect(DB)conn.execute(CREATE TABLE IF NOT EXISTS outbox (id INTEGER PRIMARY KEY AUTOINCREMENT, event_type TEXT, payload TEXT, status TEXT DEFAULT pending, created_at TEXT DEFAULT CURRENT_TIMESTAMP))conn.commit()conn.close()defadd(event_type:str,payload:dict)-None:与业务数据同事务调用保证『库有单、记录必有』。connsqlite3.connect(DB)try:conn.execute(INSERT INTO outbox(event_type, payload) VALUES (?, ?),(event_type,json.dumps(payload)))conn.commit()finally:conn.close()defclaim_batch(limit50):领取一批待发记录投递器消费。connsqlite3.connect(DB)rowsconn.execute(SELECT id, event_type, payload FROM outbox WHERE statuspending ORDER BY id LIMIT ?,(limit,)).fetchall()returnrowsdefmark_sent(record_id:int)-None:connsqlite3.connect(DB)conn.execute(UPDATE outbox SET statussent WHERE id?,(record_id,))conn.commit()conn.close()步骤 3Outbox 投递器 扫描兜底任务目标投递失败不丢单扫描任务兜底重发。# outbox_deliverer.pyimportjson,timefromceleryimportCeleryfromoutboximportclaim_batch,mark_sent,_init appCelery(outbox_deliver,brokerredis://localhost:6379/0)app.task(nameoutbox.deliver,bindTrue)defdeliver_batch(self):投递器把 pending 记录转发成任务消息。_init()forrid,etype,payloadinclaim_batch():payloadjson.loads(payload)try:ifetypeorder_created:app.send_task(orders.send_order_sms,args[payload[order_id],payload[mobile]])mark_sent(rid)# 投递成功才标记exceptExceptionasexc:print(f[outbox] record{rid}投递失败:{exc}等下轮扫描重试)returndeliveredapp.task(nameoutbox.scan,bindTrue)defscan_pending(self):扫描兜底每 30 秒扫一次未 sent 的自动重投幂等键保证不重复。_init()deliver_batch.delay()returnscanned# 启动Worker Beat每 30 秒扫一次 outboxcelery-Aoutbox_deliverer worker--loglevelinfo--poolsolo celery-Aoutbox_deliverer beat--loglevelinfo--scheduleoutbox_beat:scan_every_30s# beat_schedule 配置见代码清单运行结果文字描述下单接口写订单 写 outbox同一事务投递器把 outbox 记录转成短信任务手动模拟 Broker 挂掉再恢复——pending 记录在 Broker 恢复后由扫描任务自动重投status从 pending → sent全程「库有单 → 任务必发出」。步骤 3.5Outbox 的幂等投递——防「重扫重投」目标扫描兜底会重复投递同一条记录消息侧必须幂等第 28 章思考题 1 的落地。# outbox_deliverer.py 中改进投递消息带 event_id 幂等键app.task(nameoutbox.deliver,bindTrue)defdeliver_batch(self):_init()forrid,etype,payloadinclaim_batch():payloadjson.loads(payload)try:# 消息带 event_idoutbox 记录 ID——消费方按它去重app.send_task(orders.send_order_sms,args[payload[order_id],payload[mobile]],headers{event_id:foutbox-{rid}})mark_sent(rid)exceptExceptionasexc:print(f[outbox] record{rid}投递失败:{exc}等下轮扫描重试)returndelivered说明headers[event_id]就是第 11 章「幂等键」在消息层的载体——即使同一记录被扫两遍、发两条消息消费方按 event_id 去重后只发一次短信before_task_publish信号也可统一注入第 26 章。「扫描兜底」与「幂等投递」是一对只做前者会把丢消息变成重复消息。步骤 4Django 集成要点对照附代码目标给 Django 团队一份「开箱即用」的对照。# Djangosettings.py 配好 app 后tasks.py 用 shared_task第 5 章fromceleryimportshared_taskshared_task(nameorders.send_order_sms)defsend_order_sms(order_id,mobile):...# Django事务提交后再投递on_commitfromdjango.dbimporttransactiontransaction.atomicdefplace_order(request):orderOrder.objects.create(...)# 写库未提交transaction.on_commit(lambda:send_order_sms.delay(order.id,order.mobile))# ★ 提交后才投递returnorder运行结果文字描述on_commit里投递——事务回滚时回调不执行无幽灵任务提交成功才发任务。注意 Django 需保证django.setup()与app.autodiscover_tasks()先于 Worker 启动执行celery/fixups/django.py会自动处理大部分但首次配置记得manage.py侧完成初始化。3.3 可能遇到的坑及解决方法坑现象解决先发任务后提交回滚后任务照发幽灵任务事务提交后再投递步骤 1/4提交后投递失败Broker 抖动丢任务Outbox 兜底重扫步骤 2/3eager 模式测不出路由错误测试全绿、上线即挂集成测试用真实 Broker第 28 章FastAPI 里每个请求 new Celery()连接池暴涨app 模块级单例步骤 1 写法shared_task 任务找不到模块没被 autodiscoverDjango 配置CELERY_IMPORTS或 autodiscover 路径Outbox 重复投递扫描兜底把同一条记录发两遍消息带 event_id 幂等键步骤 3.53.4 完整代码清单与测试验证清单web_order.py事务提交后投递、outbox.pyoutbox 表、outbox_deliverer.py投递器扫描、Django 对照代码。事务一致性决策表沉淀 Wiki场景方案说明Django 单任务transaction.on_commit原生钩子最轻FastAPI/Flask 单任务接口层提交后投递手写提交点关键链路订单/支付Outbox 模式同事务记录 扫描兜底测试环境eager mock只测逻辑别信序列化/路由测试验证# tests/test_web_outbox.pyimportsqlite3fromweb_orderimportcreate_orderfromoutboximportclaim_batch,mark_sentdeftest_order_created_and_tasks_published():rcreate_order(mobile13900000000)assertr[order_id]0assertlen(r[tasks])2deftest_outbox_record_pending_then_sent():fromoutboximportadd,claim_batch,mark_sent add(order_created,{order_id:1,mobile:138})rowsclaim_batch()assertlen(rows)1androws[0][1]order_createdmark_sent(rows[0][0])assertclaim_batch()[]# 已标记后不再被领取deftest_rollback_does_not_publish():# 事务回滚场景抛异常路径不投递模拟fromunittestimportmockwithmock.patch(web_order.sqlite3.connect,side_effectRuntimeError(tx rollback)):try:create_order(mobilex)exceptRuntimeError:passpython-mpytest tests/test_web_outbox.py-v# 3 passed4. 项目总结4.1 优点 缺点维度Outbox同事务记录 扫描提交后直接投递先发任务后提交原子性记录与业务同事务无记录无投递可靠性扫描兜底不丢Broker 抖动即丢丢错复杂度中表投递器扫描低低一致性最强中最弱4.2 适用场景适用① Django 生态on_commit 原生② FastAPI/Flask 的订单/支付关键链路手写提交点 Outbox③ 多任务与事务联动的场景三件套④ 对「库有单任务必发」有 SLA 的团队⑤ 需要「投递可审计、可重放」的合规场景outbox 表即审计底账。不适用① 纯通知类非关键任务Outbox 是过度设计提交后投递足够② 无状态无事务的旁路流程如纯缓存刷新③ 已经有可靠消息中间件 事务消息的团队Outbox 是其自研替代品。4.3 注意事项task_always_eager只配测试环境进生产配置让所有任务变同步调用性能与语义双崩。FastAPI 的 Celery 实例要模块级单例每个请求 new 一个 连接池/注册表爆炸。Outbox 投递要幂等同一条记录重复扫描不能重复投递投递前按 event_id 去重步骤 3.5。Django 升级 Celery 版本时fixups/django的行为变化要回归on_commit 支持情况随版本演进。outbox 表要「留痕可查」sent记录别删保留一段窗口供对账第 23 章审计思维pending量进监控。4.4 常见踩坑经验3 个生产故障故障幽灵单——库里没单、短信已发。根因先发任务后提交事务回滚后任务照跑。对策提交点后移 on_commit/Outbox本章落地。教训投递时机是「事务提交」的函数不是「代码行序」的函数。故障FastAPI 启动 10 分钟后连接池打满。根因每个请求Celery(xxx)新建实例。对策模块级单例。教训框架集成第一课app 是单例不是请求级对象。故障测试全绿上线订单任务全 NotRegistered。根因测试用 eager 模式序列化/路由错误没暴露。对策集成测试真实 Broker第 28 章。教训eager 的绿是「模拟题」的绿。故障Outbox 扫描兜底把「投递成功但没标记」的记录又发了一遍用户收到两条短信。根因重扫重投无幂等。对策消息带 event_id 幂等键步骤 3.5。教训兜底机制本身也要幂等否则兜底就是新的故障源。4.5 思考题Outbox 的「投递」与「标记 sent」不是原子的先投递后标记投递器崩溃在两者之间会发生什么怎么保证不重复投递提示幂等键与事件 IDDjango 的transaction.on_commit在「嵌套事务」atomic 内再开 atomic里回调什么时候触发答案见第 28 章开头的「上一章思考题参考答案」。事务一致性是「库与消息」的交接礼仪——下一章把「验证交接没有漏洞」的方法论测试体系补齐。延伸阅读与资源Dify 从入门到进阶LLM 应用平台实战修炼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 实战修炼与源码剖析