同一订单的“已支付”先在页面出现,“已创建”随后才到。团队确认两条消息 Key 都是订单号,于是把问题归因给 Kafka。可同 Key 只约束分区映射,不替应用生成顺序,也不替下游串行提交。
Kafka 能保证分区日志中的 offset 顺序;业务最终可见顺序还要同时守住生产意图、分区映射、发送重试和消费并发四道边界。
先问清楚“哪个顺序”
| 顺序 | 证据 | 常见破坏者 |
|---|---|---|
| 业务发生顺序 | 业务版本、事件时间、状态机 | 多线程、跨服务时钟与提交竞争 |
| Producer 调用顺序 | 客户端日志、trace | 多 Producer 并发调用 |
| 分区日志顺序 | partition + offset | 分区映射变化、非幂等重试 |
| 处理完成顺序 | 消费任务开始/结束时间 | 线程池、异步 I/O、批量提交 |
| 用户可见顺序 | 目标库版本与更新时间 | 下游覆盖写、缓存、索引延迟 |
“同 Key”只使默认分区选择在分区数不变时保持稳定。它不能证明两个服务对同一 Key 的调用先后,也不能阻止 Consumer 把连续 offset 扔进线程池后逆序完成。
Producer 端的两个典型窗口
多实例并发
实例 A 先生成version=10,网络稍慢;实例 B 后生成version=11,先抵达 Leader。Kafka 只按抵达并追加的顺序给 offset,不按业务时间重新排序。
非幂等重试
若关闭幂等并允许多个 in-flight 请求,前一批失败重试、后一批先成功,日志可能出现反转。官方 Producer 配置明确提示这一风险;启用幂等时,max.in.flight.requests.per.connection<=5的允许范围内会保序。Producer Configs
这不是说幂等能创造业务顺序:它只能保护同一 Producer 到同一分区的协议顺序。
分区数变化会让同一个 Key 改道
典型分区映射可抽象为:
partition = hash(serializedKey) mod partitionCountTopic 从 6 个分区扩到 12 个后,同一 Key 的余数可能改变;历史事件留在旧分区,新事件进入新分区。两个分区之间不存在总序,Consumer 并行读取时就会出现跨分区可见反转。
扩分区不会自动重分布已有数据,也不能再缩回去。内部状态 Topic 更不应手工扩分区。Basic Operations
Consumer 端最常见:拉取有序,完成无序
offset 100: CREATE → 远程调用 800 ms offset 101: PAID → 远程调用 20 ms若两个任务并发执行,PAID先落库并不违反 Kafka 分区顺序。更危险的是批次末尾统一提交 offset:101 完成并推动提交后,100 若失败,重启时还可能直接越过未完成业务。
解决方式不是“所有 Consumer 单线程化”,而是按 Key 分片串行、不同 Key 并行;或者目标表用单调业务版本做条件更新:
UPDATEorder_stateSETstatus=:status,version=:versionWHEREorder_id=:idANDversion<:version;状态机还应拒绝非法跃迁,避免“最后写入者获胜”把旧状态覆盖新状态。
只读取证:用 partition、offset、version 闭合链路
1. 核对 Topic 分区数和 Consumer 分配
bin/kafka-topics.sh --bootstrap-server broker:9092\--describe--topicorder-events bin/kafka-consumer-groups.sh --bootstrap-server broker:9092\--describe--grouporder-projection--members--verbose两条命令都只读。重点不是 Lag,而是同一业务 Key 是否跨分区、哪一个实例持有相关分区。
2. 给每条事件保留六元组
event_id, business_key, business_version, producer_instance, partition, offset如果 offset 顺序与业务版本相反,问题在 Kafka 追加之前;如果 offset 正确而完成时间相反,问题在消费执行;如果两条记录跨分区,继续查扩分区或自定义 Partitioner。
3. 不要用时间戳代替版本
跨主机时间可能漂移,事件时间也可能来自不同阶段。业务版本应由单一事实源递增,或使用能表达因果的状态机版本。
修复选择
| 约束 | 合适方案 | 代价 |
|---|---|---|
| 单 Key 必须严格顺序 | 稳定分区 + Key 内串行 | 热 Key 限制吞吐 |
| 只要求最终状态正确 | 版本条件更新 + 幂等 | 需定义冲突规则 |
| 已经扩分区 | 新 Topic 重分区并双读校验 | 迁移复杂、需切换窗口 |
| 多服务共同生产 | 统一事件出口或序列号服务 | 增加协调点 |
修改分区数属于不可逆操作。执行前要抽样计算 Key 新旧映射,确认顺序 SLA,准备新 Topic/双写迁移和回退路径;不能把“扩完再观察”当方案。
验证标准
构造多个 Key,每个 Key 携带严格递增版本,同时制造 Producer 重试和 Consumer 处理时延差异。技术上验证每个 Key 的落库版本单调递增、旧版本被拒绝且 offset 提交不越过未完成记录;业务上验证订单状态不回退、重复通知不发生。
源码与 Java:同时观察业务版本、Partition 和 Offset
KafkaProducer.doSend取得元数据并选分区,再进入RecordAccumulator.append。Key 路由实现在BuiltInPartitioner,分区数变化会改变哈希映射,但不会搬迁旧日志。
以下示例按 Kafka 4.3.1 API 静态审阅,未在本环境执行扩分区对照实验。
importjava.util.Properties;importorg.apache.kafka.clients.producer.*;importorg.apache.kafka.common.serialization.StringSerializer;publicclassKeyOrderProbe{publicstaticvoidmain(String[]args)throwsException{Propertiesp=newProperties();p.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,"localhost:9092");p.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,StringSerializer.class);p.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,StringSerializer.class);p.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG,"true");try(KafkaProducer<String,String>producer=newKafkaProducer<>(p)){for(intversion=1;version<=3;version++){RecordMetadatam=producer.send(newProducerRecord<>("order-events","order-42","version="+version)).get();System.out.printf("version=%d partition=%d offset=%d%n",version,m.partition(),m.offset());}}}}扩分区前后分别运行并比较 partition。映射是key bytes + partition count → BuiltInPartitioner → 分区批次 → offset。它证明路由和日志顺序,不能证明异步下游按同序完成。
结论
同 Key 是局部路由条件,不是端到端顺序协议。只有稳定分区、Producer 保序、Key 内处理顺序、业务版本和目标端防回退同时成立,用户看到的顺序才真正可控。