代码语言

知识点思维导图

17 个知识节点

Kafka(02) - 消息顺序、重复与丢失

读完后,你应能完成以下任务:

  • 绘制“Kafka(02) - 消息顺序、重复与丢失 / 先定义“成功”发生在哪一层”的关键对象与数据流,解释“必须分别定义生产确认、Broker 持久化、消费处理和业务副作用四个检查点。”,并用源码位置、日志或 Trace 标注证据。
  • 为“Kafka(02) - 消息顺序、重复与丢失 / 顺序为何会被打乱”设计正常与异常输入,验证“Kafka 只保证单 Partition 日志顺序。”,输出首个偏差位置与回归测试结果。
  • 实现“Kafka(02) - 消息顺序、重复与丢失 / 重复与丢失如何产生”的最小代码或配置,检验“Offset 先提交、业务后失败,会丢失处理机会。”,输出命令、结果与 Diff,并说明不适用边界。

一、先建立全局:消息顺序、重复与丢失 是什么?

理解“消息顺序、重复与丢失”,先要把标题中的对象放进同一条处理链:它接收什么输入,经过哪些状态变化,最终用什么证据判断结果。下表不另造概念,只把作者正文已经解释的章节按依赖顺序连起来。

“消息顺序、重复与丢失”的第一个核心判断是:必须分别定义生产确认、Broker 持久化、消费处理和业务副作用四个检查点。。先弄清这个判断中的对象和输入输出,后面的实现、故障和验收才有共同语境。

顺序 章节 读完本节应抓住的结论
1 先定义“成功”发生在哪一层 必须分别定义生产确认、Broker 持久化、消费处理和业务副作用四个检查点。
2 顺序为何会被打乱 Kafka 只保证单 Partition 日志顺序。
3 重复与丢失如何产生 Offset 先提交、业务后失败,会丢失处理机会。
4 验证与排障 故障注入要覆盖 Producer 超时重试、Broker 重启、消费者业务提交后崩溃、Offset 提交失败、重复和乱序事件。
5 Producer send 成功只表示消息达到配置要求的确认级别 Producer send 成功只表示消息达到配置要求的确认级别,
6 不表示消费者已经完成业务 不表示消费者已经完成业务。

1.1 核心对象之间怎样衔接

flowchart LR
  S1["先定义“成功”发生在哪一层"] --> S2
  S2["顺序为何会被打乱"] --> S3
  S3["重复与丢失如何产生"] --> S4
  S4["验证与排障"] --> S5
  S5["Producer send 成功只表示消息达到配置要求的确认级别"]

这张图只表达本文的讲解顺序,不替代正文机制。判断“消息顺序、重复与丢失”是否真正掌握,需要能从最后一个结果沿图回到前面每个章节的输入、状态变化和证据。

1.2 再看失败:问题最早会出现在哪一步?

在“消息顺序、重复与丢失”的对象和顺序已经明确后,再看可观察的失败:计划退化、死锁、热点击穿、消息重复丢失或恢复后数据不一致。定位时不从最后一条错误猜原因,而是沿上图找第一个偏离正文结论的节点。

二、先定义“成功”发生在哪一层

Producer send 成功只表示消息达到配置要求的确认级别, 不表示消费者已经完成业务。 消费者提交 Offset 只表示下次从哪里继续,不表示外部数据库事务一定成功。 必须分别定义生产确认、Broker 持久化、消费处理和业务副作用四个检查点。

Producer 常用 acks=all, 并让 Topic 的 min.insync.replicas 与副本数共同控制容错; 启用幂等 Producer 可避免部分重试导致的日志重复。 若配置 acks=0/1 或在 ISR 不足时仍接受写入,故障窗口的数据风险会扩大。

三、顺序为何会被打乱

Kafka 只保证单 Partition 日志顺序。 消息使用不同 Key、增加 Partition、Producer 并发重试或消费者把同一 Partition 的记录无序提交到线程池, 都可能打乱业务观察顺序。 对状态事件增加实体版本,消费者只接受比当前版本新的变化,可以降低乱序覆盖。

UPDATE order_projection
SET status = :status, event_version = :version
WHERE order_id = :orderId
  AND event_version < :version;

影响行数为零表示重复或旧事件,不应继续触发通知等副作用。 全局顺序代价极高,通常应把需求收敛为单实体或单账户顺序。

四、重复与丢失如何产生

业务成功后、Offset 提交前进程崩溃,会重复消费; Offset 先提交、业务后失败,会丢失处理机会。 因此通常在业务成功后提交 Offset,并让消费者幂等。 使用 event_id 唯一表、业务唯一键或状态机去重, 去重记录与业务写入放在同一数据库事务。

INSERT INTO consumed_event(event_id, consumed_at)
VALUES (:eventId, NOW()); -- event_id 唯一

-- 同一事务内执行真正业务更新;唯一冲突即视为已处理

非法消息不要无限重试阻塞 Partition。 区分临时错误与永久错误, 有限退避后进入死信, 并保存原消息、错误、次数和 Trace 供修复重放。

五、验证与排障

故障注入要覆盖 Producer 超时重试、Broker 重启、消费者业务提交后崩溃、Offset 提交失败、重复和乱序事件。 验证目标不是“没有重复”,而是重复不会产生重复业务结果,且任何失败都可追踪和恢复。

发现业务缺数据时, 从 Outbox/Producer 日志、Topic Offset、消费组位点、消费者日志和业务去重表逐层核对。 只看消费者“处理成功”日志无法证明数据库提交和 Offset 状态。

验收清单

  • Producer 确认、重试和幂等配置与数据风险一致。
  • 消费者在业务提交后再提交 Offset,副作用具备幂等键。
  • 顺序边界落实为稳定 Key 与实体版本,而非口头保证。
  • 死信可修复、可重放且重放仍保持幂等。

六、动手验证:先跑通 消息顺序、重复与丢失,再改变一个变量

前面的章节已经建立问题、概念和机制。现在把“消息顺序、重复与丢失”放进同一套基线中运行;本节不再引入新术语,只验证前文结论能否被复现。

6.1 基线与候选只允许一个变量不同

验证“消息顺序、重复与丢失”时,先固定数据快照、并发条件、客户端配置、拓扑和故障注入点。候选方案只能改变本次要验证的变量;如果同时更换数据、依赖和配置,即使结果改善,也不能知道是哪一项产生作用。

执行“消息顺序、重复与丢失”时,动作是:执行正常读写与故障场景,记录查询计划、锁、复制或消费状态。原始结果不能只保留截图或汇总分数,必须同步保存:执行计划、慢日志、锁等待、Offset、复制延迟、指标和数据校验,使下一次复查可以在同一输入上重放。

实验要素 本文要求
固定条件 固定数据快照、并发条件、客户端配置、拓扑和故障注入点
唯一变量 本次候选方案与基线之间的一项明确差异
原始证据 执行计划、慢日志、锁等待、Offset、复制延迟、指标和数据校验
通过阈值 一致性与性能满足正文约束,故障恢复后没有丢失或重复副作用
立即停止 计划退化、死锁、热点击穿、消息重复丢失或恢复后数据不一致

6.2 执行前先排除不可比较条件

“消息顺序、重复与丢失”开始前先确认下面四项;任一项不成立,都应先修复实验条件,而不是解释结果。

  • 基线能够在“消息顺序、重复与丢失”的当前环境重复运行。
  • 候选只改变一个与“消息顺序、重复与丢失”结论直接相关的条件。
  • “消息顺序、重复与丢失”的基线和候选使用同一批输入、同一版本依赖与同一通过阈值。
  • “消息顺序、重复与丢失”的原始输出和失败现场不会被重试、格式化或汇总覆盖。

6.3 执行后先核对证据完整性

结果出来后先检查证据,再讨论“消息顺序、重复与丢失”是否通过。缺少中间状态时,最终输出只能说明现象,不能证明机制。

检查项 当前文章的判定
输入可追溯 固定数据快照、并发条件、客户端配置、拓扑和故障注入点
过程可回放 执行正常读写与故障场景,记录查询计划、锁、复制或消费状态
结果可审计 执行计划、慢日志、锁等待、Offset、复制延迟、指标和数据校验

“消息顺序、重复与丢失”的一次合格基线对照按以下顺序执行:

  1. 保存“消息顺序、重复与丢失”基线版本及输入摘要,确认基线本身可以重复运行。
  2. 写下“消息顺序、重复与丢失”候选方案唯一变化的变量,以及它预期影响的指标。
  3. 在同一环境执行“消息顺序、重复与丢失”:执行正常读写与故障场景,记录查询计划、锁、复制或消费状态。
  4. 为“消息顺序、重复与丢失”保存:执行计划、慢日志、锁等待、Offset、复制延迟、指标和数据校验。
  5. 使用“消息顺序、重复与丢失”预登记条件判断:一致性与性能满足正文约束,故障恢复后没有丢失或重复副作用。
  6. 如果“消息顺序、重复与丢失”未通过,不修改第二个变量,先恢复基线并保留失败现场。

七、用一张矩阵验证 消息顺序、重复与丢失 的关键结论

矩阵按正文顺序列出“消息顺序、重复与丢失”的结论。一次实验只选择一行,只改变这一行对应的条件;不要把多行合并成一个无法归因的大实验。

正文章节 已解释的结论 本轮唯一变量 必须保存的证据
先定义“成功”发生在哪一层 必须分别定义生产确认、Broker 持久化、消费处理和业务副作用四个检查点。 只改变与“先定义“成功”发生在哪一层”相关的条件 执行计划、慢日志、锁等待、Offset、复制延迟、指标和数据校验
顺序为何会被打乱 Kafka 只保证单 Partition 日志顺序。 只改变与“顺序为何会被打乱”相关的条件 执行计划、慢日志、锁等待、Offset、复制延迟、指标和数据校验
重复与丢失如何产生 Offset 先提交、业务后失败,会丢失处理机会。 只改变与“重复与丢失如何产生”相关的条件 执行计划、慢日志、锁等待、Offset、复制延迟、指标和数据校验
验证与排障 故障注入要覆盖 Producer 超时重试、Broker 重启、消费者业务提交后崩溃、Offset 提交失败、重复和乱序事件。 只改变与“验证与排障”相关的条件 执行计划、慢日志、锁等待、Offset、复制延迟、指标和数据校验
Producer send 成功只表示消息达到配置要求的确认级别 Producer send 成功只表示消息达到配置要求的确认级别, 只改变与“Producer send 成功只表示消息达到配置要求的确认级别”相关的条件 执行计划、慢日志、锁等待、Offset、复制延迟、指标和数据校验
不表示消费者已经完成业务 不表示消费者已经完成业务。 只改变与“不表示消费者已经完成业务”相关的条件 执行计划、慢日志、锁等待、Offset、复制延迟、指标和数据校验

7.1 记录本次实际实验

下面的记录用于“消息顺序、重复与丢失”当前这一次实验,不是第二套知识目录。先从矩阵选择一个章节,再填写实际值;没有填写的字段表示尚未验证。

topic: "消息顺序、重复与丢失"
selected_chapter: required
claim_from_article: required
baseline_version: required
changed_condition: exactly_one
execution: "执行正常读写与故障场景,记录查询计划、锁、复制或消费状态"
evidence: "执行计划、慢日志、锁等待、Offset、复制延迟、指标和数据校验"
pass_when: "一致性与性能满足正文约束,故障恢复后没有丢失或重复副作用"
stop_when: "计划退化、死锁、热点击穿、消息重复丢失或恢复后数据不一致"
observed_result: required
first_deviation: null_or_evidence
recovery_replay: required_after_failure

7.2 边界实验必须证明能够停止和恢复

成功路径只能证明“消息顺序、重复与丢失”在当前样本上工作,不能证明它可以进入生产。边界实验需要主动制造:计划退化、死锁、热点击穿、消息重复丢失或恢复后数据不一致,并观察系统是否在产生不可逆副作用前停止。

场景 只改变什么 应保存什么 通过标准
正常路径 使用已知有效输入 执行计划、慢日志、锁等待、Offset、复制延迟、指标和数据校验 一致性与性能满足正文约束,故障恢复后没有丢失或重复副作用
边界路径 把一个输入推进到约束临界值 临界值前后的输出与指标 不静默降级,不把部分结果冒充成功
明确失败 注入:计划退化、死锁、热点击穿、消息重复丢失或恢复后数据不一致 原始错误、首个异常阶段和最终状态 失败被正确分类且没有扩大副作用
恢复重放 执行:从数据入口、存储状态、复制消费链路和恢复步骤定位根因 原失败样本的复测证据 原样本恢复,正常样本没有回归

恢复动作不是简单重启。对于“消息顺序、重复与丢失”,第一步是:从数据入口、存储状态、复制消费链路和恢复步骤定位根因。完成后使用原始失败样本复测;只验证一个新样本成功,不能证明触发条件已经消失。

“消息顺序、重复与丢失”边界实验结束后,应把正常、临界、失败和恢复四类记录放在同一个运行批次中。这样才能区分“候选方案真的修复问题”和“环境变化让问题暂时没有出现”。

八、消息顺序、重复与丢失 的结果解释

解释“消息顺序、重复与丢失”实验时先看首个偏差,而不是最后一条错误。最后的异常通常只是上游状态错误的结果;从末端反推容易误把症状当根因。

观察结果 可以支持的判断 下一步
主链路没有达到预期 计划退化、死锁、热点击穿、消息重复丢失或恢复后数据不一致 先执行:从数据入口、存储状态、复制消费链路和恢复步骤定位根因
异常链路无法恢复 计划退化、死锁、热点击穿、消息重复丢失或恢复后数据不一致 先执行:从数据入口、存储状态、复制消费链路和恢复步骤定位根因
新样本成功但原样本仍失败 修复没有覆盖原始触发条件 固定原失败输入,恢复基线后重新比较
指标改善但证据无法回链 数据、版本或中间状态没有固定 暂停发布,补齐可追溯记录后重跑

“消息顺序、重复与丢失”只有同时满足“一致性与性能满足正文约束,故障恢复后没有丢失或重复副作用”,并且没有出现“计划退化、死锁、热点击穿、消息重复丢失或恢复后数据不一致”,才可以认为主链路通过。这里的“通过”只对当前固定版本、样本和环境有效,不能外推到尚未测试的容量、权限或数据分布。

如果“消息顺序、重复与丢失”候选方案与基线差异很小,先检查证据分辨率是否足够;如果差异很大,先排除数据泄漏、环境漂移和版本不一致。两种情况都不能只看一个汇总均值,需要回到逐样本输出和中间状态。

“消息顺序、重复与丢失”故障定位完成后,记录“现象、首个偏差、根因、改动、原样本复测”五项。缺少原样本复测时,只能标记为待观察,不能标记为已解决。

九、消息顺序、重复与丢失 的发布判断

发布判断需要把“消息顺序、重复与丢失”的质量、失败边界和恢复能力放在同一份记录中。以下任一条件缺失,都应停止扩量,而不是用“基本正常”替代证据。

  • “消息顺序、重复与丢失”的基线与候选只存在一个计划内变量。
  • “消息顺序、重复与丢失”的输入、代码、依赖、配置和数据版本可以追溯。
  • “消息顺序、重复与丢失”的正常、临界、失败和恢复样本使用同一套断言。
  • “消息顺序、重复与丢失”的原始输出、中间状态和失败现场已经保留。
  • “消息顺序、重复与丢失”的日志、Trace、截图和测试数据已经脱敏。
  • “消息顺序、重复与丢失”的停止条件、负责人和回滚入口已经演练。
  • “消息顺序、重复与丢失”尚未覆盖的输入、权限、容量和外部依赖已经登记。

最终记录至少包含基线版本、唯一变量、原始证据、首个偏差、恢复复测和发布责任人。没有参与本次修改的人如果不能据此重放“消息顺序、重复与丢失”的判断,就不能发布。

十、总结

  • 先定义“成功”发生在哪一层:必须分别定义生产确认、Broker 持久化、消费处理和业务副作用四个检查点。
  • 顺序为何会被打乱:Kafka 只保证单 Partition 日志顺序。
  • 重复与丢失如何产生:Offset 先提交、业务后失败,会丢失处理机会。
  • 验证与排障:故障注入要覆盖 Producer 超时重试、Broker 重启、消费者业务提交后崩溃、Offset 提交失败、重复和乱序事件。

学完自测

选择所有正确答案;提交后逐项核对判断依据。

1在“消息顺序、重复与丢失”中,需要同时满足“先建立全局:消息顺序、重复与丢失 是什么?”与“核心对象之间怎样衔接”。给定正文约束“下表不另造概念,只把作者正文已经解释的章节按依赖顺序连起来。”,哪些判断保持了原有处理机制?多选
2“消息顺序、重复与丢失”出现偏差:“在“消息顺序、重复与丢失 / 再看失败:问题最早会出现在哪一步?”中,即使不满足“计划退化、死锁、热点击穿、消息重复丢失或恢复后数据不一致”,结果与副作用仍会保持不变。”已成为实际行为。围绕“再看失败:问题最早会出现在哪一步?”与“先定义“成功”发生在哪一层”,哪些判断能定位被改变的职责或边界?多选
3评审“消息顺序、重复与丢失”方案时,验收条件包含“消息使用不同 Key、增加 Partition、Producer 并发重试或消费者把同一 Partition 的记录无序提交到线程池,”。关于“顺序为何会被打乱”与“重复与丢失如何产生”的哪些决策符合正文机制?多选