「恰好一次」是句谎话
一句话总结
“恰好一次”不是传输层替你实现的,而是由至少一次投递 + 接收方去重共同实现,后半部分必须由我们负责。
为什么使用消息队列
同步集成成立的前提,是对方此刻处于可用状态;异步集成取消了这个前提。
- 时间解耦——对方现在不可用,也可稍后处理
- 吸收负载——即使每秒涌入一万条,也可先在队列积压,再以每秒一千条处理
- 故障隔离——接收方故障不会蔓延到发送方
- 重处理——失败消息可以重新投入
代价是复杂度与重复。“无法立即得到响应”意味着必须另行设计结果通知渠道;“会重试”则意味着同一条消息可能被处理两次。
投递保证的三个级别
| 级别 | 含义 | 现实用途 |
|---|---|---|
| at-most-once | 最多一次,可能丢失 | 日志、指标等允许丢失的数据 |
| at-least-once | 至少一次,可能重复 | 大多数实际消息系统 |
| exactly-once | 恰好一次 | 仅在特定条件下成立 |
必须理解一点。
“恰好一次” = “至少一次投递” + “接收方去重”
也就是说,exactly-once 不是传输层施展魔法得到的,而是接收方过滤重复消息,使最终结果看起来只发生一次。即使消息代理宣传“支持 exactly-once”,也只是在特定条件下成立,例如使用同一集群与事务 API;一旦写入外部系统,这项保证就会被打破。
因此,消费者必须始终假设消息会重复。 这不是例外,而是默认前提。
顺序保证的范围
另一个常见误解是:
顺序只在分区(或队列)内部得到保证,并不保证整个主题的全局顺序。
因此,如果业务要求“同一订单号的消息必须按顺序处理”,就必须使用订单号作为分区键。这样,同一订单的消息会进入同一分区,顺序也得以保持。
如果不了解这一点而采用轮询分发,就可能先处理取消消息,再处理下单消息。这种情况每天只出现一两次,因而极难排查。
消费者设计的标准形式
1. 메시지 수신
2. 파싱 및 형식 검증 → 실패: 즉시 error/DLQ (재시도해도 똑같다)
3. 멱등 확인 (이미 처리했나) → 이미 처리: 아무것도 안 하고 ack
4. 업무 처리 (DB 트랜잭션)
5. 처리 이력 기록 ← 4와 같은 트랜잭션 안에서
6. ack (큐에서 제거)
关键是第 5 步必须与第 4 步位于同一事务中。如果分开执行,就可能出现“业务已经处理,但处理记录没写下来”的状态;此时重试就会造成重复处理。
而且,第 6 步绝不能早于第 4 步。 如果先确认消息、随后处理过程中崩溃,消息就永久消失,这正是系统退化为 at-most-once 的位置。
解析失败不要重试
把第 2 步单独列出的原因很明确:JSON 损坏或缺少必填字段的消息,重试 100 次也会失败 100 次。 如果它进入普通重试逻辑,就会在队首不断失败,阻塞后续正常消息。这类消息称为毒药消息(poison message)。
因此,应把错误分成两类。
- 永久错误(格式错误、必填值缺失、代码不存在)→ 立即送入 DLQ
- 临时错误(数据库连接失败、对方系统返回 5xx、超时)→ 重试
如果不做区分,队列就会阻塞,而阻塞的队列很快会变成整体业务停摆。
队列积压是最重要的指标
如果只能用一个数字观察异步集成的健康状态,应选择消费者延迟(consumer lag),也就是“积压消息数”或“最旧未处理消息的年龄”。
- 平时接近 0 → 正常
- 持续增长 → 消费者跟不上(性能问题或消费者宕机)
- 突然暴增 → 生产方流量激增或消费者故障
- 长时间保持固定值且不下降 → 可能被毒药消息阻塞
告警不应只按绝对值,而应结合趋势与持续时间。“lag 超过 1000”不如“lag 连续 10 分钟增长”有效,因为在存在批量流入的系统中,瞬时 lag 很大可能是正常现象。
文件队列也是队列
许多 SI 现场没有 Kafka、RabbitMQ 等消息代理,此时可以使用基于目录的队列,这种方式比想象中可靠。
/data/if/inbox/ ← 도착
/data/if/processing/ ← 처리 중 (원자적 mv 로 이동 = 잠금)
/data/if/done/ ← 성공
/data/if/error/ ← 실패 (DLQ 역할)
关键在于,mv 在同一文件系统内是原子操作。只有成功完成 inbox → processing 移动的进程才拥有该消息,即使启动多个消费者也不会重复处理。
有两点需要注意。
- 跨文件系统的
mv实际是复制后删除,并非原子操作。 必须在同一挂载点内移动。 - 消费者可能在发送方写完文件之前就将其取走。 因此,应先用临时名称写入,完成后再改名,或采用同时创建**完成标志文件(
.ok)**的约定。
提前设计重处理流程
故障发生后一定会需要重处理,应提前准备以下信息。
- 什么失败了——error/ 中的消息及失败原因
- 为什么失败——按原因分类(格式、业务、系统)
- 能否修复——修改数据后重新投入,还是请求源系统重发
- 重处理是否安全——是否具备幂等性
- 由谁批准——涉及资金的接口可能需要审批
这套流程应在上线前形成文档并进行演练。等到故障当天再设计,往往就要通宵。
在实际项目中
引入消息队列的项目,真正出问题的通常不是队列本身,而是消费者一侧的错误假设。
最常见的是重复。网络短暂中断后恢复时,消息代理会重新发送尚未收到确认的消息。这不是故障,而是按规范运行;如果消费者没有预料到,同一订单就会写入两次。而且这类事故通常要到月末结算金额对不上时才被发现——距离事故发生可能已经三周。
第二是顺序。相信整个队列都能维持顺序,一旦增加分区或消费者就会失效。顺序保证通常只在单个分区内部成立,因此必须选择键,让同一订单的消息进入同一分区。
第三是积压。消息队列能够吸收负载,所以即使消费者很慢,发送方也感觉不到异常。 如果不监控积压(lag),系统可能延迟数小时,直到有人询问“今天的数据为什么看不到”才会被发现。