Python后端中间件专题15:一条坏消息堵住整条队列——有限重试、DLQ 与重放
Python后端中间件专题15一条坏消息堵住整条队列——有限重试、DLQ 与重放值班人员收到一条“通知队列持续失败”告警。如果 Worker 没有终点一条缺字段的 JSON 会被反复投递占用执行槽并刷屏日志。如果直接 ACK 并丢弃现场证据和可恢复性又消失。DLQ 的价值不是“另一个垃圾队列”而是给失败工作一个有限、可审计、可操作的去处。值班 runbook先隔离再定点重放按reason、schema version、producer 与时间窗口聚合 SQLdead_letters阻断仍在制造毒消息的来源保存原 body 和 error。对指定 dead-letter ID 核对修复条件不批量重放异质消息。条件 claim 必须有 lease expiry防两个操作员并发获得同一行。保留原 event ID 和 body附加重放来源 header以 confirmed publisher 发回主队列确认后才写replayed_at。发布失败释放 claim确认不确定仍可能重复由消费 receipt 收敛。再次操作同一 ID 应返回not_replayed并核对最终业务 effect 只有一份。若未运行真实 broker/SQL 集成这四步只是操作合同不是完成记录。runbook 的学习目标你将能按 reason 分类 poison/permanent/exhausted 死信保留原始 event body 并重置 delivery 计数设计“按死信 ID 修复后重放”的运维流程并说明重放本身为什么也要 confirm、claim 和幂等。值班前置与授权边界请了解 12 的有界重试和 14 的 confirm 不确定窗口。manifest selector 需远程 RabbitMQ、Worker、PostgreSQL 和真实任务提交本次禁止远程与 Docker因此不会运行它。上一课练习答案答案 EX-14-01 [code]在project/下保存dlq_replay_probe.pyfromticketflow.messaging.dead_letterimportDeadLetterMessage,replay_message originalDeadLetterMessage(bodyb{event_id:event-15},reasonretries_exhausted,attempts3,headers{x-tenant:tenant-a},)replayedreplay_message(original)assertreplayed.bodyoriginal.bodyassertreplayed.attempts0assertreplayed.headers[x-tenant]tenant-aassertreplayed.headers[x-replayed-from-dlq]trueassertreplayed.headers[x-dead-letter-reason]retries_exhaustedprint(replayed.body.decode())print(replayed.attempts,replayed.headers)运行命令$env:PYTHONPATHsrc; python dlq_replay_probe.py预期输出{event_id:event-15} 0 {x-tenant: tenant-a, x-replayed-from-dlq: true, x-dead-letter-reason: retries_exhausted}重放不改写 event ID 或 body因为它仍是同一个业务事件。它只重置本次 delivery 的计数并附加运维来源使后续日志能区分首次投递与人工重放。答案 EX-14-02 [prose]poison_message先保留原 JSON 和 parse error核对 producer schema/version只有在能恢复为完整 envelope 后才重放。permanent_failure先判断业务规则、租户边界或不支持 event type 是否真已修复如果事件本身不再合法应终结并审计不强行重放。retries_exhausted先确认依赖已恢复并有容量再分批重放否则会立即再次压垮依赖。不能全量一键重放因为 DLQ 里的 reason、schema 和故障时间可能不同旧毒消息会再次占满队列而突发重放还可能造成下游流量洪峰。为什么隔离不能直接画成“已修复”Worker 在 parse 失败、永久错误或 retry 耗尽后把 original primitive event、reason、attempts 和 error 持久化到 SQLdead_letters。RabbitMQ 的 dead queue 可承载一个可观测指针但 SQL row 才是本项目的可查询重放事实。这样即使发布 DLQ pointer 失败操作员仍能检索和修复原始记录。重放不应直接把 row 标记为 replayed 后再 publish。正确顺序是先用有过期时间的 claim 防止并发操作员重复发布通过 confirmed publisher 发回主队列成功后才设replayed_at。发布失败则释放 claimrow 仍可重试。即使两个操作员在极端 confirm 不确定窗口中仍产生重复原 event ID 未改变13 的 receipt 会将业务副作用收敛为一份。runbook 的证据分层本地可运行的 primitive replay 合同python -m pytest tests/unit/test_messaging_core.py::test_dlq_replay_preserves_envelope_and_resets_delivery_attempt_metadata -q本地可信输出. [100%] 1 passed本课累积 revision 另用 SQLAlchemy 2 SQLite 文件数据库验证原始 body/headers 持久化、原子条件 claim、租约过期后换主、旧 owner 不得 mark/release、确认后 mark、失败后 release 与二次重放 no-op这只证明本地 SQL 状态机。它不替代真实 RabbitMQ/Worker/PostgreSQL 运行manifest 远程 selectorpython -m pytest tests/integration/test_dead_letter.py::test_poison_notification_reaches_dlq_and_replay_applies_one_effect -q仍需验证 poison 持久化、人工修复、重放只产生一份 effect 且第二次重放返回not_replayed。本次远程检查保持PENDING未运行。runbook 中途失败时DLQ 持续增长按 reason/schema version/producer 聚合先停止新毒消息源而不是立即重放存量。replayed_at已有值但下游没有消息检查是否在 confirm 前就标记完成正确实现只能在 confirmed publish 后更新。两个操作员都说重放成功需要 DB 条件 claim而不是先 SELECT 后 UPDATE同时保留 consumer receipt 作最后防线。交接记录有限重试保护容量DLQ 保留证据按 ID 修复重放恢复业务。重放是一次新的传输尝试因此 confirm、claim 和下游幂等一个都不能少。本课练习练习 EX-15-01 [code]实现一个 in-memorySlaStore始终返回同一个(ticket_id, version)但mark_escalated只对首次条件转移返回 True。在同一now连续扫描两次断言 escalated 为 1 然后 0并打印结果。练习 EX-15-02 [prose]比较 Celery Beat “按 UTC 间隔每分钟触发扫描”与“为每张工单安排一个本地时区 ETA task”在时区/DST、丢调度、重复调度和补跑方面的差异给出 TicketFlow 的选择。下一篇预告16 不再等用户请求或消息到来Beat 会重复发起 SLA 扫描。我们会把“凌晨两点”转换成 UTC 事实时间并用条件 UPDATE 保证两次 tick 只升级一次。完整核心模块带 claim、确认与释放的死信重放DLQ payloads preserve the original JSON body for deliberate replay.from__future__importannotationsfromdataclassesimportdataclass,fieldfromdatetimeimportdatetime,timedelta,timezoneimportjsonfromtypingimportCallable,Dict,Protocolfromuuidimportuuid4fromsqlalchemyimportor_,updatefromsqlalchemy.ormimportSessionfromticketflow.storage.modelsimportDeadLetterModeldataclass(frozenTrue)classDeadLetterMessage:body:bytesreason:strattempts:intheaders:Dict[str,str]field(default_factorydict)dataclass(frozenTrue)classReplayMessage:body:bytesattempts:intheaders:Dict[str,str]defreplay_message(message:DeadLetterMessage)-ReplayMessage:Reset delivery accounting while retaining the immutable original event body.headersdict(message.headers)headers[x-replayed-from-dlq]trueheaders[x-dead-letter-reason]message.reasonreturnReplayMessage(bodymessage.body,attempts0,headersheaders)dataclass(frozenTrue)classClaimedDeadLetter:The durable row a single operator currently owns for replay.id:strmessage:DeadLetterMessage lease_owner:strclassDeadLetterStore(Protocol):defclaim_for_replay(self,*,dead_letter_id:str,lease_owner:str,now:datetime,lease_for:timedelta)-ClaimedDeadLetter|None:Atomically claim an unreplayed or expired-lease row....defmark_replayed(self,*,dead_letter_id:str,lease_owner:str,at:datetime)-bool:Record completion only for the active claim after publisher confirmation....defrelease_replay(self,*,dead_letter_id:str,lease_owner:str,reason:str,now:datetime)-None:Make a failed publication eligible for a later deliberate retry....classConfirmingReplayPublisher(Protocol):defpublish(self,*,body:bytes,headers:Dict[str,str])-bool:Return true only after the primary queue publisher confirmation arrives....classSqlAlchemyDeadLetterStore:Durable original body and atomic owner-checked replay lease.def__init__(self,*,sessions:Callable[[],Session])-None:self._sessionssessionsdefpersist(self,*,dead_letter_id:str,message:DeadLetterMessage)-None:withself._sessions()assession:withsession.begin():session.add(DeadLetterModel(iddead_letter_id,bodymessage.body,reasonmessage.reason,attemptsmessage.attempts,headers_jsonjson.dumps(message.headers,sort_keysTrue),))defclaim_for_replay(self,*,dead_letter_id:str,lease_owner:str,now:datetime,lease_for:timedelta)-ClaimedDeadLetter|None:ifnotlease_ownerorlease_fortimedelta(0):raiseValueError(owner and positive lease are required)withself._sessions()assession:withsession.begin():claimedsession.execute(update(DeadLetterModel).where(DeadLetterModel.iddead_letter_id,DeadLetterModel.replayed_at.is_(None),or_(DeadLetterModel.replay_claim_expires_at.is_(None),DeadLetterModel.replay_claim_expires_atnow),).values(replay_claim_ownerlease_owner,replay_claim_expires_atnowlease_for))ifclaimed.rowcount!1:returnNonerowsession.get(DeadLetterModel,dead_letter_id)messageDeadLetterMessage(row.body,row.reason,row.attempts,json.loads(row.headers_json))returnClaimedDeadLetter(row.id,message,lease_owner)defmark_replayed(self,*,dead_letter_id:str,lease_owner:str,at:datetime)-bool:withself._sessions()assession:withsession.begin():completedsession.execute(update(DeadLetterModel).where(DeadLetterModel.iddead_letter_id,DeadLetterModel.replay_claim_ownerlease_owner,DeadLetterModel.replay_claim_expires_atat,DeadLetterModel.replayed_at.is_(None)).values(replayed_atat,replay_claim_ownerNone,replay_claim_expires_atNone,last_replay_errorNone))returncompleted.rowcount1defrelease_replay(self,*,dead_letter_id:str,lease_owner:str,reason:str,now:datetime)-None:withself._sessions()assession:withsession.begin():session.execute(update(DeadLetterModel).where(DeadLetterModel.iddead_letter_id,DeadLetterModel.replay_claim_ownerlease_owner,DeadLetterModel.replayed_at.is_(None)).values(replay_claim_ownerNone,replay_claim_expires_atNone,last_replay_errorreason))classDeadLetterReplayer:Claim → confirmed publish → mark, with release on every publish failure.def__init__(self,*,store:DeadLetterStore,publisher:ConfirmingReplayPublisher,clock:Callable[[],datetime]lambda:datetime.now(timezone.utc),lease_for:timedeltatimedelta(minutes5))-None:iflease_fortimedelta(0):raiseValueError(replay lease must be positive)self._storestore self._publisherpublisher self._clockclock self._lease_forlease_fordefreplay(self,*,dead_letter_id:str,operator_id:str)-bool:nowself._clock()claim_ownerf{operator_id}:{uuid4()}claimedself._store.claim_for_replay(dead_letter_iddead_letter_id,lease_ownerclaim_owner,nownow,lease_forself._lease_for)ifclaimedisNone:returnFalsemessagereplay_message(claimed.message)try:confirmedself._publisher.publish(bodymessage.body,headersmessage.headers)ifnotconfirmed:raiseRuntimeError(publisher confirmation was not received)exceptExceptionaserror:self._store.release_replay(dead_letter_idclaimed.id,lease_ownerclaim_owner,reasontype(error).__name__,nowself._clock())raisereturnself._store.mark_replayed(dead_letter_idclaimed.id,lease_ownerclaim_owner,atself._clock())

相关新闻

数据库学生成绩管理系统课程设计:从E-R图到SQL实现全攻略

数据库学生成绩管理系统课程设计:从E-R图到SQL实现全攻略

简介:一份以SQL Server 2008与VC6.0为开发环境的学生成绩管理系统数据库课程设计报告,适合计算机、软件工程等专业学生在数据库课程设计、毕业设计或实训中参考。报告从课题背景与需求分析出发,完整覆盖概念设计(E-R模型&#xff…

2026/10/12 3:48:18 阅读更多 →
Python GIL 深度解析:从多线程翻车到并行方案

Python GIL 深度解析:从多线程翻车到并行方案

我第一次在 Python 里正儿八经写并发的时候,一度怀疑是电脑坏了。任务很简单:8 个线程分别跑一段纯 CPU 计算,按道理就算不跑满 8 核,至少也比单线程快个三四倍吧。结果 8 个线程跑完的时间不但没变短,反而比串行还慢了…

2026/10/12 3:48:18 阅读更多 →
Dev Container 并行生命周期脚本执行:object 语法原理、规范与实战

Dev Container 并行生命周期脚本执行:object 语法原理、规范与实战

开发工具 【免费下载链接】spec Development Containers: Use a container as a full-featured development environment. 项目地址: https://gitcode.com/gh_mirrors/spec2/spec 点击查看 免费下载 导读 本文围绕 Development Container Specification&#xff0…

2026/10/12 3:48:18 阅读更多 →

最新新闻

量子开发者人才缺口百万?入门技能图谱与实操路径全解析

量子开发者人才缺口百万?入门技能图谱与实操路径全解析

一份“2030年量子开发人才缺口达百万”的预测,最近在朋友圈被转得很猛。我第一反应是:这数字靠不靠谱先放一边,“量子开发者”到底是个什么工种,多数人其实说不清楚。作为写过几年经典软件、又花了不少时间钻进量子计算这个交叉领…

2026/10/12 4:44:49 阅读更多 →
闲置PS5变身标准媒体终端与测试机:AnyPS5配置指南

闲置PS5变身标准媒体终端与测试机:AnyPS5配置指南

把“AnyPS5”这个名字扔上来的时候,估计会有人下意识往越狱、固化那类方向想。我先把话说清楚:我这里的AnyPS5不是破解工具,也不是某个隐藏系统,而是我这段时间在工作室里把几台PS5翻来覆去折腾之后,沉淀出来的一套“任…

2026/10/12 4:44:49 阅读更多 →
Excel 关键指标 Top-N 高亮导出实战:SenseNova-Skills top-value-coloring 技能深度解析

Excel 关键指标 Top-N 高亮导出实战:SenseNova-Skills top-value-coloring 技能深度解析

AI 技能人工智能深度研究数据分析媒体生成 【免费下载链接】SenseNova-Skills Modular SenseNova skills for building AI-powered office assistants and productivity workflows 项目地址: https://gitcode.com/gh_mirrors/se/SenseNova-Skills 点击查看 免费下载…

2026/10/12 4:44:49 阅读更多 →
OpenCore Legacy Patcher 完整指南:2007—2017 老 Mac 免费安装 macOS Big Sur 到 Sequoia 新系统

OpenCore Legacy Patcher 完整指南:2007—2017 老 Mac 免费安装 macOS Big Sur 到 Sequoia 新系统

OpenCore Legacy Patcher 完整指南:2007—2017 老 Mac 免费安装 macOS Big Sur 到 Sequoia 新系统 【免费下载链接】OpenCore-Legacy-Patcher Experience macOS just like before 项目地址: https://gitcode.com/GitHub_Trending/op/OpenCore-Legacy-Patcher …

2026/10/12 4:44:49 阅读更多 →
Claude Code + MCP 实操:从空文件夹到可玩 Unity 游戏的全自动开发闭环

Claude Code + MCP 实操:从空文件夹到可玩 Unity 游戏的全自动开发闭环

开年这几个月,AI 辅助编程的玩法算是彻底变天了。以前大家讨论的是"AI 能不能帮你写代码",现在的问题已经变成"AI 能不能直接替你完成一个独立开发者的全套工作流"。我今天想聊的,就是我最近反复折腾又实测过好几轮的完整…

2026/10/12 4:44:49 阅读更多 →
OpenClaw技能合集:从Clawdbot到Moltbot的Agent Skill精选与TaoToken接入实践

OpenClaw技能合集:从Clawdbot到Moltbot的Agent Skill精选与TaoToken接入实践

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/12 4:43:48 阅读更多 →

日新闻

复古胶片颗粒感噪点合成器:Canvas ImageData 像素高斯杂色注入算法

复古胶片颗粒感噪点合成器:Canvas ImageData 像素高斯杂色注入算法

在数码相机、高清显示屏与现代矢量图形技术高度发达的今天,画面可以做到绝对的锐利、平滑与无瑕。然而,当一张秋日手账插画或拍立得照片过于“平整无瑕”时,往往会散发出一种冰冷生硬的“数码塑料感(Digital Plasticity&#xff0…

2026/10/12 0:00:59 阅读更多 →
活字印刷古籍线装排版:Canvas 竖排文字与栏线自适应算法

活字印刷古籍线装排版:Canvas 竖排文字与栏线自适应算法

在现代网页与移动端设计中,横排(Horizontal Layout)早已经成为了绝对的主流。然而,当我们翻开泛黄的线装古籍、宋版木刻诗集,或是欣赏一张茶道雅集的手写便签时,那种**自上而下纵向书写、自右向左逐列铺展&…

2026/10/12 0:00:59 阅读更多 →
周日晚间的“精神松绑减震器”:无压力情绪倾倒箱与温和轻声陪伴

周日晚间的“精神松绑减震器”:无压力情绪倾倒箱与温和轻声陪伴

每到周日的晚上八点到十点,很多人心里都会悄悄亮起一盏警示灯。 在心理学上,这种现象有一个专门的称谓——“周日夜晚焦虑症(Sunday Scaries)”。明天又是周一,闹钟又要重新在七点响彻卧房;脑海里仿佛有一个…

2026/10/12 0:00:59 阅读更多 →

周新闻

流感时间序列预测实战:ARIMA/LSTM全流程拆解与避坑指南

流感时间序列预测实战:ARIMA/LSTM全流程拆解与避坑指南

简介:基于 ARIMA、LSTM、Transformer 等模型的流感时间序列预测 Python 源码,面向计算机相关专业课程设计与期末大作业学生,以及项目实战学习者。内容覆盖预处理、平稳性检验、定阶、残差分析、多模型对比预测的完整时序建模流程,…

2026/10/12 0:16:30 阅读更多 →
影刀RPA新手教程:键盘模拟输入实战——输入文本与模拟按键的区别

影刀RPA新手教程:键盘模拟输入实战——输入文本与模拟按键的区别

影刀RPA新手教程:键盘模拟输入实战——输入文本与模拟按键的区别 做影刀RPA自动化,十个新手有八个栽在"往输入框里填东西"这件事上:要么填不进去,要么填了一半,要么直接把原来内容追加在后面。这背后的根因&…

2026/10/12 0:16:38 阅读更多 →
影刀RPA新手教程:阅文起点小说数据采集实战——书籍信息与章节内容

影刀RPA新手教程:阅文起点小说数据采集实战——书籍信息与章节内容

影刀RPA新手教程:阅文起点小说数据采集实战——书籍信息与章节内容 1. 认识影刀:什么场景该用RPA采小说数据 起点中文网的页面结构相对稳定——分类榜单、书籍详情、章节内容三块独立页面,跳转链路清晰。这种场景非常适合影刀自动化&#x…

2026/10/12 0:16:43 阅读更多 →

月新闻

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/11 10:45:37 阅读更多 →
Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证

Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/11 14:36:53 阅读更多 →
黑夜航拍船只数据集训练YOLOV5模型全流程解析

黑夜航拍船只数据集训练YOLOV5模型全流程解析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/11 14:36:54 阅读更多 →