第 12 章 · Offset 与消息语义可视化

At Most Once / At Least Once / Exactly Once · 自动提交陷阱 · seek · 业务幂等

动画对比:消息流转 → 处理 → 提交 → 宕机

分别在 At Most Once / At Least Once / Exactly Once 三种语义下,模拟同一段消息流,观察「丢」「重」结果。点击「下一步」逐步推进,可以在任意时刻按「💥 模拟宕机」看后果。

At Most Once(先提交后处理)

已 poll 已处理
就绪

At Least Once(先处理后提交)

已 poll 已处理
就绪

Exactly Once(处理 + 提交原子)

已 poll 已处理
就绪
10
原始消息数
0
AMO 真实处理
0
ALO 真实处理(含重)
0
EOS 真实处理

自动提交:处理时长 vs commit interval

调节「单条处理时长」和「auto.commit.interval.ms」,观察自动提交在错误时机触发,导致「丢」或「重」的窗口大小。

运行结果

0
真处理条数
0
已提交 offset
0
丢失条数
0
重复条数(重启后)
尚未运行

seek 任意 offset:把消费者拨到任意位置

这是一个 Partition 的 30 条消息(offset 0~29),点击任意 offset 模拟 seek(tp, offset),观察消费者光标跳转。

当前消费者光标 未消费
点击任意方块跳转
方法含义典型场景
seek(tp, offset)跳到指定 offset排查特定问题消息
seekToBeginning(tps)跳到分区最早全量重放、回溯计算
seekToEnd(tps)跳到分区最新只关心最新数据
offsetsForTimes(tps,ts)按时间戳找 offset线上排障定位「14:00 的异常」

业务幂等:去重表 + 唯一约束

模拟同一条订单消息在 At Least Once 下被消费 3 次。第一次插入成功,第二/三次因 UNIQUE 约束IntegrityError,被业务侧识别为重复并跳过。

消息流

consumer_dedup 表

order_idprocessed_at消息序号
空表