Python后端中间件专题13:“执行两次”不能“通知两次”——幂等消费者

发布时间:2026/10/11 1:41:45
Python后端中间件专题13:“执行两次”不能“通知两次”——幂等消费者
Python后端中间件专题13“执行两次”不能“通知两次”——幂等消费者现在把时间定格在最危险的一行通知已经写入ACK 还没发出Worker 突然退出。下一个 Worker 看到的仍是合法消息不能靠redelivered直接丢弃因为第一次也可能根本没有提交。唯一可靠的问题是数据库中是否已存在这个 event 的这种 effect事务现场的三份证物第一份是不可变 event ID第二份是 SQL 中(event_id, effect_name)唯一 receipt第三份是同库business_tasks通知记录。只看到第一份说明“有工作机会”只看到第二份却看不到第三份说明事务边界可能被拆开两条 SQL 事实必须一起提交或一起回滚。图里的 receipt 因此与业务任务同属数据库事实不是绿色临时缓存。取证目标你将能设计(event_id, effect_name)唯一约束说明为什么 receipt 和本地副作用必须共用一个数据库事务并能区分“已处理”与“已交付真实外部服务”。取证工具与边界请先掌握 12 的 redelivery/retry 区分了解 SQLAlchemysession.begin()和唯一约束。本课使用 SQLite 验证 SQL 事务合同它不是 PostgreSQL 并发时序的运行时证据。上一课练习答案答案 EX-12-01 [code]将下面保存为idempotent_effect_probe.py在project/下运行。代码使用StaticPool确保多个 Session 共享同一个 in-memory SQLite 连接。fromdatetimeimportdatetime,timezonefromsqlalchemyimportcreate_engine,func,selectfromsqlalchemy.ormimportsessionmakerfromsqlalchemy.poolimportStaticPoolfromticketflow.messaging.outboximportEventEnvelopefromticketflow.messaging.sql_storeimportSqlAlchemyEffectTransactionfromticketflow.storage.modelsimportBase,BusinessTaskModel,EffectReceiptModel enginecreate_engine(sqlitepysqlite:///:memory:,connect_args{check_same_thread:False},poolclassStaticPool,)Base.metadata.create_all(engine)sessionssessionmaker(bindengine,expire_on_commitFalse)eventEventEnvelope(event_idevent-13,ticket_idticket-13,tenant_idtenant-a,version1,trace_context{},occurred_atdatetime(2026,9,22,tzinfotimezone.utc),payload{event_type:ticket.created},)effectsSqlAlchemyEffectTransaction(sessionssessions)asserteffects.apply_notification_once(event)isTrueasserteffects.apply_notification_once(event)isFalsewithsessions()assession:receiptssession.scalar(select(func.count()).select_from(EffectReceiptModel))taskssession.scalar(select(func.count()).select_from(BusinessTaskModel))assert(receipts,tasks)(1,1)print(ffirstTrue secondFalse receipts{receipts}tasks{tasks})engine.dispose()运行命令$env:PYTHONPATHsrc; python idempotent_effect_probe.py预期输出firstTrue secondFalse receipts1 tasks1第二次不是用SELECT先查出存在再返回而是让唯一约束仲裁并捕获特定 duplicateIntegrityError。这避免两个 Worker 同时“查无”后都执行。答案 EX-12-02 [prose]如果 receipt 先单独提交随后写通知失败重投看到 receipt 就会跳过于是通知永久丢失。如果先写通知并提交再单独写 receipt两步之间崩溃会让重投再写一份通知。TicketFlow 把 receipt insert 与business_tasksinsert 放在同一session.begin()中。副作用失败则两者都回滚副作用成功则两者一起提交重投冲突时新事务整体不产生第二份 effect。对唯一键做反证只用event_id去重会过于宽泛。同一ticket.created可能要产生 notification、search projection 和 audit 等不同效果如果 notification 的 receipt 阻止了 search就把幂等做成了丢数据。因此唯一键是(event_id, effect_name)。也不能只用ticket_id。一张工单有多个 version 与多个业务事件按 ticket 去重会把后续合法通知全部吞掉。event ID 代表一次不可变的业务事件effect name 划定某一个消费者副作用。证物三不是外部短信BusinessTaskModel是一个可审计本地 sink用于确定性验收“通知任务产生一次”。如果真实副作用是发送邮件或短信外部 API 不能与本地 DB 参加同一 ACID 事务。此时应把“待发送通知”持久化为本地任务再由下一层发送器使用提供商幂等 key 或对账机制。所以本课证明的是“本地业务任务只产生一次”不是宣称任意外部世界都能 exactly-once。用事务测试复现现场在project/运行python -m pytest tests/unit/test_messaging_sql_store.py::test_sql_effect_transaction_persists_one_notification_task_for_duplicate_delivery -q本地 Python 3.11 可信输出. [100%] 1 passed该本地测试经 SQLAlchemy model 和事务断言一份 receipt、一份 business task。可信远端 selector 会经 RabbitMQ/Celery 提交同一 event ID 两次并观察 Worker 两次接收随后查询 PostgreSQL 的 receipt/effect 各一还用缺失 ticket 外键制造 effect insert 失败要求 receipt 同事务回滚补齐 ticket 后重投成功一次。未执行前仍为 PENDING。异常现场的反查路径duplicate delivery 仍产生两份通知检查 unique constraint 是否真正在 DB而不只是 Python 先查后写。receipt 存在但任务缺失检查两次 insert 是否同 Session、同 transaction以及 effect 失败是否被错误吞掉。所有IntegrityError都当 duplicate只能吞掉目标唯一约束冲突其他完整性错误必须向上抛出。取证结论幂等消费不是“看到 redelivered 就不做”而是让事实库唯一约束决定某 event/effect 是否首次提交。receipt 与本地副作用必须同生共死。本课练习练习 EX-13-01 [code]构造一个内存OutboxStore和可配置 confirm 结果的 publisher。对同一OutboxEvent分别验证 confirmTrue 时只调用mark_publishedconfirmFalse 时只调用release_for_retry并打印两个RelayResult。练习 EX-13-02 [prose]分析“先提交 ticket再发布消息”和“先发布消息再提交 ticket”各自的崩溃窗口说明 Transactional Outbox 为什么把事件行与 ticket 放在同一 DB 事务。下一篇预告14 会移到 producer 一端工单已提交时 RabbitMQ 可能不可用而 publisher confirm 也可能丢失。Outbox 不会消除重复但会把丢失变成可恢复状态。完整核心模块幂等效果 SQL 事务SQLAlchemy transaction boundary for an idempotent notification effect.from__future__importannotationsfromtypingimportCallablefromuuidimportuuid4fromsqlalchemy.excimportIntegrityErrorfromsqlalchemy.ormimportSessionfromticketflow.messaging.outboximportEventEnvelopefromticketflow.storage.modelsimportBusinessTaskModel,EffectReceiptModel SessionFactoryCallable[[],Session]classSqlAlchemyEffectTransaction:Receipt insertion and local fake notification task share one transaction.def__init__(self,*,sessions:SessionFactory)-None:self._sessionssessionsdefapply_once(self,*,event_id:str,effect:str,apply:Callable[[],None])-bool:try:withself._sessions()assession:withsession.begin():session.add(EffectReceiptModel(event_idevent_id,effect_nameeffect))session.flush()apply()returnTrueexceptIntegrityErroraserror:detailstr(error).lower()ifeffect_receipts.event_idindetailandeffect_nameindetailoruq_effect_event_nameindetail:returnFalseraisedefapply_notification_once(self,event:EventEnvelope)-bool:try:withself._sessions()assession:withsession.begin():session.add(EffectReceiptModel(event_idevent.event_id,effect_namenotification))session.flush()session.add(BusinessTaskModel(idstr(uuid4()),ticket_idevent.ticket_id,tenant_idevent.tenant_id,task_typenotification,statuscompleted,))returnTrueexceptIntegrityErroraserror:detailstr(error).lower()ifeffect_receipts.event_idindetailandeffect_nameindetailoruq_effect_event_nameindetail:returnFalseraise