Python后端中间件专题12:Worker 在 ACK 前死掉——重投、重试与时间限制
故障日志里出现两条很像的记录:redelivered=True和retry in 10s。前者是 broker 对未 ACK delivery 的恢复,后者是 task 对已分类异常的主动重试。把两者混为“系统会再试”,就无法设置上限,也无法解释为什么一条消息会执行两次。
先分诊两行日志
| 线索 | 触发者与计数 | 下一步检查 |
|---|---|---|
redelivered=True | channel/Worker 消失,broker 重新交付未 ACK 消息;不增加 Celery retry count | 查 ACK 前崩溃点和业务 receipt |
retry in 10s | handler 分类为RetryableTaskError后主动安排下一次 task;受max_retries约束 | 查依赖故障、countdown 和耗尽去向 |
poison_message | envelope 无法解析,等待不能修复 schema | 保留原 body 与错误并隔离 |
诊断的主键应是业务 event ID,不能靠 delivery tag 判断是不是同一次业务效果。
本次分诊要得到什么
本课要建立一张可执行的失败分类表:进程丢失由 late ACK +task_reject_on_worker_lost导致 redelivery;短暂依赖失效由RetryableTaskError进入指数退避;永久错误和耗尽重试进入失败汇。
日志从哪里来
要先理解 11 的 late ACK 和 prefetch。本课 manifest selector 标记为remote-integration:它需要真 RabbitMQ、Celery Worker、PostgreSQL,还会 kill/restart Worker 并暂停 Elasticsearch。本次明确禁止远程和 Docker,因此不运行该 selector,不把任何本地 unit 结果写成远程 PASS。
上一课练习答案
答案 EX-11-01 [code]
在project/下保存retry_policy_probe.py:
fromticketflow.workers.tasksimportBoundedRetryPolicy,RetryableTaskError policy=BoundedRetryPolicy(max_retries=2,base_delay_seconds=5)cases=[(RetryableTaskError("temporary"),0,(True,5,False)),(RetryableTaskError("temporary"),1,(True,10,False)),(RetryableTaskError("temporary"),2,(False,None,True)),(ValueError("invalid"),0,(False,None,True)),]forerror,retries,expectedincases:decision=policy.decision(error,retries_so_far=retries)observed=(decision.retry,decision.countdown_seconds,decision.dead_letter)assertobserved==expectedprint(type(error).__name__,retries,*observed)运行命令:$env:PYTHONPATH="src"; python retry_policy_probe.py
预期输出:
RetryableTaskError 0 True 5 False RetryableTaskError 1 True 10 False RetryableTaskError 2 False None True ValueError 0 False None True重试计数表示已经发生的 retry;达到 2 时不再生成第三个 retry。非RetryableTaskError不做时间退避,因为等待不会修复格式错误或业务不变量。
答案 EX-11-02 [prose]
- handler 开始前 Worker 被 kill:late ACK 尚未发生,broker 在 channel 丢失后把 delivery 重投;这不消耗 Celery task 内部的 retry 计数。
- 副作用提交后、ACK 前被 kill:broker 仍会重投,因此 handler 必须使用 event receipt 让第二次变为 duplicate/no-op。
- handler 抛出可重试错误:task 主动调用
retry(countdown=...),生成受max_retries约束的新执行尝试。这是应用层决策,不是 broker 对消费者失联的推断。
三种时序都可以导致相同 event 再次到达,所以幂等不能只在“调用 retry 的分支”打补丁。
分诊后才选择恢复路径
BoundedRetryPolicy是纯函数式决策:输入异常类型和已重试次数,输出 retry/countdown/dead-letter。这使得延迟表可单元测试,也避免在各个 handler 中散落“遇错就重试”。
当 envelope 无法解析时,task 将其标记为poison_message;当非短暂错误抛出时标记permanent_failure;可重试错误达到上限后标记retries_exhausted。这三个 reason 给运维不同的修复入口,而不是把所有错误丢进同一个无限循环。
第三类故障:长任务的 soft/hard 时间边界
重试次数限制与任务执行时间限制是两套配置。本课将soft_time_limit=15与time_limit=20注册到通知任务。可信远端 selector 暂停受控依赖后,先观测任务特定的 soft-limit 日志;另一条任务必须先同时出现目标 event 的 started 日志与 queue unacked,再暂停该 Worker 的唯一 prefork 子进程,使 soft signal 无法执行 Python cleanup,由 Celery 主进程在 20 秒边界记录任务特定Hard time limit (20s) exceeded并杀子进程。直接 kill Worker 也先记录 running/healthy baseline,并在恢复后做有界 readiness 验证。随后仍核对 dead letter、receipt 与 effect 收敛。未实际远端执行前只算 PENDING。
分诊证据:本地策略与远程时序
本地可运行的策略契约:python -m pytest tests/unit/test_messaging_core.py::test_retry_policy_dead_letters_after_bounded_retry_exhaustion -q
本地预期并可复现的输出:
. [100%] 1 passed课程 manifest 的正式 selector 是:python -m pytest tests/integration/test_celery_delivery.py::test_worker_loss_redelivers_then_bounded_retry_dead_letters -q。它要求授权远程 Compose,本次结果是SKIP(未运行),不是 PASS。真实验收同时要求任务特定 soft/hard 日志、Worker loss 后redelivered=True、每条失败事件一份 dead letter,以及每个 event 的通知 receipt/effect 仍恰一。
排除误判
- 把所有
ConnectionError都无限 retry:必须有次数上限、退避与最终可观测去处。 - 把
redelivered=True当成 Celery retry 计数:它们来自不同机制,调查时同时记录 event ID、delivery info 和 retry 次数。 - 标题写“时间限制”就假设已有 hard time limit:先查 task 注册参数,再分别验证 soft cleanup 与 hard kill 后重投;没有证据时保留待验收。
分诊结论
broker redelivery 修复未 ACK delivery,task retry 处理已识别的短暂失败,两者都会重复执行。有界退避只防止无限重试,不消除副作用去重责任。
本课练习
练习 EX-12-01 [code]
使用 SQLite in-memory engine、Base.metadata.create_all、SqlAlchemyEffectTransaction和一个EventEnvelope,连续两次调用apply_notification_once。断言返回值为True, False,并查询 receipt 与 business task 数量都是 1。
练习 EX-12-02 [prose]
解释为什么先插入 receipt 并提交、再写通知副作用会丢通知;再解释为什么先写副作用、后单独插 receipt 会重复通知。
下一篇预告
13 会把“最终会重复”收敛为一个 SQL 事务:同一个 event 可到达多次,但 notification receipt 和可审计业务任务只能一起成功一次。
完整核心模块:有界重试与失败分类
"""Bounded retry classification for the notification worker."""from__future__importannotationsfromdataclassesimportdataclassimportjsonfromtypingimportCallablefromticketflow.messaging.outboximportEventEnvelopeclassRetryableTaskError(RuntimeError):pass@dataclass(frozen=True)classRetryDecision:retry:boolcountdown_seconds:int|Nonedead_letter:boolclassBoundedRetryPolicy:def__init__(self,*,max_retries:int=3,base_delay_seconds:int=5):ifmax_retries<0orbase_delay_seconds<1:raiseValueError("retry bounds must be non-negative and delay positive")self._max_retries=max_retries self._base_delay_seconds=base_delay_secondsdefdecision(self,error:Exception,*,retries_so_far:int)->RetryDecision:ifnotisinstance(error,RetryableTaskError):returnRetryDecision(retry=False,countdown_seconds=None,dead_letter=True)ifretries_so_far>=self._max_retries:returnRetryDecision(retry=False,countdown_seconds=None,dead_letter=True)returnRetryDecision(retry=True,countdown_seconds=self._base_delay_seconds*(2**retries_so_far),dead_letter=False,)defregister_tasks(app,*,notification_handler:Callable[[EventEnvelope],object]|None=None,dead_letter_handler:Callable[...,str]|None=None,)->None:"""Register a late-ACK task with bounded retry and an explicit failure sink."""ifnotification_handlerisNone:raiseValueError("notification_handler is required")defsend_to_dead_letter(**details:object)->str:ifdead_letter_handlerisNone:raiseValueError("dead_letter_handler is required for failed delivery")returndead_letter_handler(**details)@app.task(name="ticketflow.workers.tasks.deliver_notification",bind=True,acks_late=True,max_retries=3,soft_time_limit=15,time_limit=20,)defdeliver_notification(task,event:dict)->dict:try:envelope=EventEnvelope.from_json_bytes(json.dumps(event).encode("utf-8"))except(TypeError,ValueError,json.JSONDecodeError)aserror:dead_letter_id=send_to_dead_letter(event=event,reason="poison_message",attempts=task.request.retries+1,error=str(error),cause=error,)event_id=event.get("event_id","unknown")ifisinstance(event,dict)else"unknown"return{"event_id":str(event_id),"status":"dead_lettered","dead_letter_id":dead_letter_id,}try:applied=notification_handler(envelope)exceptExceptionaserror:decision=BoundedRetryPolicy(max_retries=task.max_retries).decision(error,retries_so_far=task.request.retries)ifdecision.retry:raisetask.retry(exc=error,countdown=decision.countdown_seconds,max_retries=task.max_retries,)reason=("retries_exhausted"ifisinstance(error,RetryableTaskError)else"permanent_failure")dead_letter_id=send_to_dead_letter(event=event,reason=reason,attempts=task.request.retries+1,error=str(error),cause=error,)return{"event_id":envelope.event_id,"status":"dead_lettered","dead_letter_id":dead_letter_id,}return{"event_id":envelope.event_id,"status":"dispatched"ifappliedisnotFalseelse"duplicate",}