恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
Python后端中间件专题12:Worker 在 ACK 前死掉——重投、重试与时间限制
首页
资讯中心
/
Python后端中间件专题12:Worker 在 ACK 前死掉——重投、重试与时间限制
Python后端中间件专题12:Worker 在 ACK 前死掉——重投、重试与时间限制
发布时间:2026/10/11 1:41:46
Python后端中间件专题12Worker 在 ACK 前死掉——重投、重试与时间限制故障日志里出现两条很像的记录redeliveredTrue和retry in 10s。前者是 broker 对未 ACK delivery 的恢复后者是 task 对已分类异常的主动重试。把两者混为“系统会再试”就无法设置上限也无法解释为什么一条消息会执行两次。先分诊两行日志线索触发者与计数下一步检查redeliveredTruechannel/Worker 消失broker 重新交付未 ACK 消息不增加 Celery retry count查 ACK 前崩溃点和业务 receiptretry in 10shandler 分类为RetryableTaskError后主动安排下一次 task受max_retries约束查依赖故障、countdown 和耗尽去向poison_messageenvelope 无法解析等待不能修复 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.pyfromticketflow.workers.tasksimportBoundedRetryPolicy,RetryableTaskError policyBoundedRetryPolicy(max_retries2,base_delay_seconds5)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:decisionpolicy.decision(error,retries_so_farretries)observed(decision.retry,decision.countdown_seconds,decision.dead_letter)assertobservedexpectedprint(type(error).__name__,retries,*observed)运行命令$env:PYTHONPATHsrc; 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 被 killlate ACK 尚未发生broker 在 channel 丢失后把 delivery 重投这不消耗 Celery task 内部的 retry 计数。副作用提交后、ACK 前被 killbroker 仍会重投因此 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_limit15与time_limit20注册到通知任务。可信远端 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 后redeliveredTrue、每条失败事件一份 dead letter以及每个 event 的通知 receipt/effect 仍恰一。排除误判把所有ConnectionError都无限 retry必须有次数上限、退避与最终可观测去处。把redeliveredTrue当成 Celery retry 计数它们来自不同机制调查时同时记录 event ID、delivery info 和 retry 次数。标题写“时间限制”就假设已有 hard time limit先查 task 注册参数再分别验证 soft cleanup 与 hard kill 后重投没有证据时保留待验收。分诊结论broker redelivery 修复未 ACK deliverytask 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):passdataclass(frozenTrue)classRetryDecision:retry:boolcountdown_seconds:int|Nonedead_letter:boolclassBoundedRetryPolicy:def__init__(self,*,max_retries:int3,base_delay_seconds:int5):ifmax_retries0orbase_delay_seconds1:raiseValueError(retry bounds must be non-negative and delay positive)self._max_retriesmax_retries self._base_delay_secondsbase_delay_secondsdefdecision(self,error:Exception,*,retries_so_far:int)-RetryDecision:ifnotisinstance(error,RetryableTaskError):returnRetryDecision(retryFalse,countdown_secondsNone,dead_letterTrue)ifretries_so_farself._max_retries:returnRetryDecision(retryFalse,countdown_secondsNone,dead_letterTrue)returnRetryDecision(retryTrue,countdown_secondsself._base_delay_seconds*(2**retries_so_far),dead_letterFalse,)defregister_tasks(app,*,notification_handler:Callable[[EventEnvelope],object]|NoneNone,dead_letter_handler:Callable[...,str]|NoneNone,)-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(nameticketflow.workers.tasks.deliver_notification,bindTrue,acks_lateTrue,max_retries3,soft_time_limit15,time_limit20,)defdeliver_notification(task,event:dict)-dict:try:envelopeEventEnvelope.from_json_bytes(json.dumps(event).encode(utf-8))except(TypeError,ValueError,json.JSONDecodeError)aserror:dead_letter_idsend_to_dead_letter(eventevent,reasonpoison_message,attemptstask.request.retries1,errorstr(error),causeerror,)event_idevent.get(event_id,unknown)ifisinstance(event,dict)elseunknownreturn{event_id:str(event_id),status:dead_lettered,dead_letter_id:dead_letter_id,}try:appliednotification_handler(envelope)exceptExceptionaserror:decisionBoundedRetryPolicy(max_retriestask.max_retries).decision(error,retries_so_fartask.request.retries)ifdecision.retry:raisetask.retry(excerror,countdowndecision.countdown_seconds,max_retriestask.max_retries,)reason(retries_exhaustedifisinstance(error,RetryableTaskError)elsepermanent_failure)dead_letter_idsend_to_dead_letter(eventevent,reasonreason,attemptstask.request.retries1,errorstr(error),causeerror,)return{event_id:envelope.event_id,status:dead_lettered,dead_letter_id:dead_letter_id,}return{event_id:envelope.event_id,status:dispatchedifappliedisnotFalseelseduplicate,}