2.7 Java 帝国之消息队列
从同步调用下游到由代理保存消息,沿生产、持久化、消费、落库和确认边界诊断重复投递与幂等。
学习目标
- 能沿生产消息、代理持久化、分派消费、业务落库和确认/重试追踪一条消息
- 能解释消息队列如何把生产者与消费者解耦,并用业务键、确认和重试处理重复投递
- 能在正常、边界和故障场景中回答:消费者处理成功却在确认前断线时,怎样避免重复扣款或重复发货
2.7 Java 帝国之消息队列
本页依据刘欣《码农翻身》(2018 年第 1 版)及出版社公开书目信息,独立重构 2.7 Java 帝国之消息队列。正文、代码、图示、实验和练习都是本课程重新设计的教学材料,不复制原书正文、插图、练习答案或代码。
先想象一家餐馆:下单的人把单子交给前台就可以离开,前台把单子放进可追踪的篮子,厨房按自己的速度取单。没有这个中间篮子,下单的人就必须站在厨房门口等结果,厨房一慢,所有人都被拖住。
这一章解决的就是“提交动作是否必须同步等到下游完成”。中间篮子不能凭空保证成功:它要保存消息、记录谁拿走、判断是否已经处理,并在失败时让消息再次出现;否则系统只会把等待换成丢单或重复处理。
三个会让消息系统失真的陷阱
九个目录节点到消息证据
2.7 Java 帝国之消息队列
本章的总问题是:一次业务提交怎样从调用者手中的同步等待,变成可追踪的异步状态。用 messageId、businessKey 和 attempt 固定身份,任何重试都必须能够解释“同一业务是否已经生效”。
张家村的历史
↡把一个同步调用拆成生产、保存、消费和结果查询,使调用方不必一直占住下游线程的演进背景。是理解问题压力的入口:当所有人都直接找同一个下游,慢调用会把等待传播给整条链。历史叙事在本章只承担动机,真正的验收证据是代理中是否有可查询消息,以及业务状态是否能被重放。
拆分
↡把一次端到端动作切成发送、保存、投递、处理和确认等边界,并为每个边界规定状态与失败处理。不是把类名切成更多类,而是让每个节点都有单独的输入、输出和故障证据。生产者只负责提交可验证的消息;消费者只负责处理契约允许的消息,不能把网络等待伪装成业务完成。
新问题
↡同步等待被移走后出现的丢失、重复、乱序、积压和最终一致性问题。会在异步边界之后出现:发送方先返回了,但业务结果尚未发生;消费者处理完却没来得及确认;代理重启后又送来同一消息。每个问题都要由消息状态、业务状态和重放记录共同回答,不能只看 HTTP 成功码。
消息队列
↡暂存消息并按协议把它交给消费者的中间服务,通常负责保存、投递、确认和失败后的再次投递。是一个有状态的缓冲边界,不是“更快的函数调用”。它至少要说明消息保存多久、谁可以读取、确认何时生效、失败怎样回退,以及积压如何被发现。
互不兼容的MQ
↡不同消息代理在确认、顺序、可见性、分区、重试和持久化语义上不相同,不能只替换客户端名称。提醒我们:接口相似不代表行为相同。迁移前先把交付保证写成测试矩阵,再验证代理客户端的确认时机、重试次数、死信策略和消费并发。
消息队列接口设计
↡把消息身份、载荷、版本、确认、重试和错误结果写成生产者与消费者共同遵守的边界。要回答五件事:消息谁产生、谁消费、业务键是什么、成功何时确认、失败如何重试。推荐的最小信封包含
messageId、businessKey、schemaVersion、occurredAt
和脱敏载荷;接口越明确,适配不同代理越安全。
配置和代码的分离
↡把队列地址、主题、消费组、并发和重试策略放在可审计配置中,让业务处理代码只表达领域规则。让同一消费者可以在不同环境使用不同代理参数,但不应把危险默认值藏起来。配置必须经过版本控制和启动校验,代码仍要显式处理确认、幂等和不可重试错误。
再次抽象
↡在代理客户端之上提炼稳定的业务消息合同与状态机,把厂商差异留在适配层。不是再包一层万能客户端,而是把业务真正关心的“已接收、处理中、已生效、待重试、不可重试”暴露出来。适配层可以翻译确认 API,却不能抹掉顺序、投递保证和失败语义的差别。
最小消息合同
record OrderMessage(String messageId, String businessKey,
int schemaVersion, String orderId, int attempt) {}
ConsumeResult consume(OrderMessage message) {
if (dedupe.alreadyApplied(message.businessKey())) return ACK;
try {
orders.applyOnce(message.businessKey(), message.orderId());
dedupe.markApplied(message.businessKey());
return ACK;
} catch (RetryableFailure failure) {
return RETRY;
}
}这段草图把幂等键和确认结果放在同一条可审计路径中。真实实现还要让 applyOnce 与幂等记录使用同一事务或等价的原子机制;如果业务写入已经提交、确认却丢失,下一次投递必须返回同一个业务结果,而不是再次扣款。
五步复核一条消息链
1. 生产消息并固定业务身份
先写出 messageId、businessKey、版本、载荷摘要和生产时间。只改变一个输入,例如把订单重复提交一次;预期是两个传输消息仍映射到一个业务效果,而不是用消息 ID 假装去重。
Lab
业务键与确认时机实验
只改变投递或断线位置,观察业务副作用是否重复,以及重放后能否安全确认。
消息保存、一次落库,随后确认
m-104 → saved → delivered → applied(order-7) → ACK
判定
accept:一次业务效果,积压回到基线
当前样本:正常投递;保存消息 ID、业务键、尝试次数、落库状态和确认结果。
正常、边界与故障证据矩阵
| 样本 | 只改变的变量 | 预期判定 | 必存证据 |
|---|---|---|---|
| 正常 | 代理可用、一次投递、业务写入成功 | 一次落库后确认,积压回到基线 | 消息 ID、业务键、确认、业务状态 |
| 边界 | 重复提交或消费并发增加 | 幂等命中,不重复产生副作用 | 尝试次数、唯一键、处理记录 |
| 故障 | 落库后、确认前消费者断线 | 重投后返回同一业务结果 | 首次提交、断线点、重放结果 |
故障诊断:先定位状态断在哪个边界
- 生产侧:查业务键、消息 ID、版本和代理接收结果;没有代理确认时不能把请求标成已完成。
- 代理侧:查主题/队列、分区、消费组、保存期限和积压;确认丢失通常意味着消息会再次可见。
- 消费侧:查实例、租约、尝试次数和顺序键;并发重复处理要回到唯一约束或幂等表验证。
- 业务侧:查业务事务、幂等记录和外部副作用;提交成功而确认失败必须能安全重放。
如果消息既没有代理确认也没有生产日志,先判定为发送不确定;如果代理有消息但业务没有状态,查分派与消费错误;如果业务状态已生效但消息仍重试,查确认时机和幂等记录。每次只改变一个故障点,并保留恢复后的基线。
术语表
名词解释
本章出现的专业名词,用大白话再讲一遍。
- 张家村的历史
用故事说明同步调用为什么会把等待传播给下游,但验收仍要回到可查询状态。
- 拆分
把发送、保存、投递、处理和确认分成有输入输出的边界。
- 新问题
异步化后必须面对的丢失、重复、积压、乱序和最终一致性风险。
- 消息队列
暂存消息并按协议投递给消费者的中间服务。
- 互不兼容的MQ
在确认、顺序、重试和持久化行为上存在差异的消息代理集合。
- 消息队列接口设计
对消息身份、版本、确认和失败结果作出的共同约定。
- 配置和代码的分离
把代理连接和运行策略放入可审计配置,把业务规则留在代码中。
- 再次抽象
在代理客户端之上保留稳定的业务合同,同时隔离厂商差异。
练习
练习
问题 1: 消费者已经写入订单,随后在确认前断线。下一次投递怎样设计,才能不重复发货?
问题 2: 代理 A 切换到代理 B 后,消息偶尔乱序且重试次数不同。你会先比较哪些合同?
问题 3: 修改本页“业务键与确认时机实验”的故障场景,使它再显示一次“业务提交后断线”轨迹,并说明重置后应恢复哪些状态。
本页小结
- 一条消息链必须区分生产、保存、投递、落库和确认。
- 消息 ID 追踪传输,业务键约束业务副作用。
- 确认丢失会带来重投,幂等记录让重放安全。
- 不同代理的顺序、重试和持久化语义必须显式适配。
读完后的自测问题是:消费者处理成功却在确认前断线时,你能否用消息轨迹、业务状态和幂等记录证明重放不会重复扣款或重复发货?