Kafka 消息明明都设了相同的 Key 保证分区有序,为什么消费端拿到的顺序还是乱的? Kafka 如何保证消息有序标准答案是用同一个 Key 发消息让同一用户的订单落在同一个分区里。比如先创建、再支付、最后发货Kafka 分区内有序顺序自然就保住了。发消息的时候用 user_id 做 Key同一个用户的消息一定落在同一个分区。理论上是这样但实际用户uid_872451的订单发货消息可能比支付消息更早消费了。排查 Broker 端的日志会发现消息在分区 3 里的 offset 是 12001创建、12002支付、12003发货顺序完全没问题。但消费端的业务日志里12003 的处理时间比 12002 还早了 400 毫秒。Broker 存储有序生产者也确实用了同一个 Key可中间环节把顺序搞乱了。Kafka 的分区有序怎么理解千万别把分区有序理解成只要消息在同一个分区里消费端一定按顺序拿到并按顺序处理。这个理解只对了一半。Kafka 保证的是分区内消息的存储顺序同一个分区里 offset 小的消息一定比 offset 大的消息先写入。消费者在调用poll()拉取消息时返回的那批消息确实是按 offset 顺序排列的。但拉取有序不等于处理有序从 Broker 到最终业务处理完成中间有好几个环节都可能把顺序搞乱。Producer 端重试把消息顺序搞反了很多时候大家压根不知道消息还没到消费端顺序就已经乱了。Kafka Producer 发送消息时要知道底层可不是发一条等一条的它有一个参数叫max.in.flight.requests.per.connection默认值是 5意思是一个连接上最多允许 5 个请求同时在飞不用等前一个请求收到 Broker 确认再发下一个。问题就出在这里假设 Producer 往分区 3 连续发了两条消息请求 1 → 消息 A用户下单 请求 2 → 消息 B用户支付正常情况下 Broker 会按收到的顺序写入A 的 offset 比 B 小没毛病。但如果请求 1 因为网络抖动失败了Producer 开启了重试请求 1 会重新发送这时候请求 2 可能已经成功写入 Broker 了。时序大概是这样的Broker 里存的顺序变成了 B 在前、A 在后消费者拿到的时候先处理支付再处理下单业务直接炸。这个问题的根因是当max.in.flight.requests.per.connection 1时多个请求并发在飞某个请求重试成功的时间点可能晚于后续请求。Kafka 在 0.11 版本之后引入了幂等 Producer可以解决这个问题后面方案部分会讲。消费端多线程并行处理打乱顺序调用poll()拉到一批消息假设拉到了 offset 100 到 199 这 100 条消息它们在ConsumerRecords里的顺序是严格按 offset 排列的。但为了提高消费吞吐量我们通常会把这批消息丢到线程池里并行处理。代码通常长这样KafkaListener(topics order_topic, groupId order-service) public void onMessage(ConsumerRecordString, String record) { // 拿到消息后丢给线程池异步处理 executorService.submit(() - { orderService.process(record); }); }或者更常见的直接在 Spring Kafka 里配了concurrency大于 1。但就算concurrency 1只有一个消费线程如果业务代码里自己起了线程池去处理顺序照样乱。看看具体会发生什么支付 20ms 就处理完了发货 30ms 处理完了创建要 80ms。最终的处理顺序变成了支付 → 发货 → 创建。完全反了。问题的本质很简单**poll()返回的消息有序不代表多线程并行处理的完成时间有序。** 每条消息的处理耗时不同先开始处理的不一定先完成。分区扩容让同一个 Key 跑到不同分区Kafka 用 Key 决定消息落到哪个分区默认的分区策略是对 Key 做 hash 然后取模// Kafka 默认分区器 DefaultPartitioner 的核心逻辑 int partition Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions;这里的numPartitions是分区总数。假设原来 Topic 有 6 个分区user_id uid_872451算出来的 hash 值对 6 取模落在分区 3。某天运维觉得消费速度不够把分区数从 6 扩到了 10。扩容之后同样的uid_872451对 10 取模可能就落到分区 7 了。从这一刻开始这个用户的新消息全部写到分区 7但之前的历史消息还在分区 3 里。如果分区 3 和分区 7 分别由不同的消费者实例处理两个消费者之间没有任何协调机制顺序就彻底没法保证了。举个具体的例子分区扩容本质上改变了 Key 到分区的映射关系。Kafka 不会把历史消息从旧分区迁移到新分区所以扩容前后同一个 Key 的消息会分散在不同分区里。Kafka 分区策略内部到底怎么算的那 Kafka 能不能搞一个一致性 hash让扩容的时候大部分 Key 还映射到原来的分区答案是不行。Kafka 的默认分区器用的就是简单的murmur2 hash 取模不是一致性 hash。这是官方有意为之的设计简单、快、可预测。代价就是分区数变了之后几乎所有 Key 的映射都会变。翻一下 Kafka 源码里DefaultPartitioner的实现public class DefaultPartitioner implements Partitioner { public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) { ListPartitionInfo partitions cluster.partitionsForTopic(topic); int numPartitions partitions.size(); if (keyBytes null) { // 没有 Key用 Sticky 策略轮询 return stickyPartitionCache.partition(topic, cluster); } // 有 Keyhash 取模 return Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions; } }所以如果你的业务对消息顺序有强依赖分区数一旦确定就别随便改。要改的话得做数据迁移确保旧分区的消息全部消费完再切流。怎么解决Producer 端开启幂等性Kafka 0.11 引入了幂等 Producer。开启之后Broker 会为每个 Producer 的每个分区维护一个递增的序列号。即使请求重试了Broker 也会根据序列号把消息放到正确的位置重复的消息会被自动去重。配置很简单加一行就行# Producer 配置 enable.idempotencetrue开启幂等之后Kafka 会自动把max.in.flight.requests.per.connection限制在 5 以内其实默认就是 5并且在 Broker 端通过序列号保证同一个 Producer 发往同一个分区的消息严格有序。不过要注意幂等性只保证单 Producer 实例 单分区的顺序。如果你有多个 Producer 实例同时往同一个分区写同一个 Key 的消息幂等性管不了。这种场景需要业务层自己保证只有一个 Producer 负责特定 Key 的写入。消费端单线程顺序处理最简单粗暴的方案别用多线程处理同一个分区的消息。KafkaListener(topics order_topic, groupId order-service, concurrency 1) public void onMessage(ConsumerRecordString, String record) { // 直接在消费线程里同步处理不丢线程池 orderService.process(record); }concurrency 1意味着只有一个消费线程拉到的消息按 offset 逐条处理。这是保证顺序的最可靠方式但吞吐量会比较低。如果只是部分 Key 需要保序而大部分消息不关心顺序可以更精细一点按 Key 做分组同一个 Key 的消息串行不同 Key 的消息并行。// 按 Key 分组的有序消费方案 publicclass OrderedConsumer { privatestaticfinalint MAX_WORKERS 16; // 预创建固定数量的单线程执行器同一个 Key 始终路由到同一个 privatefinal ExecutorService[] workers new ExecutorService[MAX_WORKERS]; public OrderedConsumer() { for (int i 0; i MAX_WORKERS; i) { workers[i] Executors.newSingleThreadExecutor( r - new Thread(r, order-worker- i) ); } } KafkaListener(topics order_topic, groupId order-service) public void onMessage(ConsumerRecordString, String record) { String key record.key(); // 同一个 Key 的 hashCode 取模后始终落到同一个 worker int slot Math.abs(key.hashCode()) % MAX_WORKERS; workers[slot].submit(() - orderService.process(record)); } }这个方案的思路是用 Key 的 hash 值把消息路由到固定数量的单线程执行器上。同一个 Key 永远在同一个线程里执行天然有序不同 Key 可以并行不浪费吞吐。不过要注意 Key 的分布均匀性如果某个 Key 的消息量特别大对应的那个线程就会成为瓶颈。扩容只加消费者不加分区针对分区扩容导致 Key 映射变化的问题前期规划好分区数后续不要改。Kafka 消费的并行度上限是分区数所以一开始就要根据业务的峰值吞吐量算好分区数。比如你预估未来三年峰值需要 12 个消费者并行处理那就直接建 12 个分区。如果消费速度真的不够了先看是不是单条消息处理太慢优化业务逻辑比加分区有用得多消费者实例数还没到分区数的上限先加消费者实在要扩分区必须做好切流方案等旧分区的消息全部消费完再让新消息按新的分区数走# Spring Kafka 消费者配置示例 spring: kafka: consumer: group-id:order-service # 关闭自动提交手动控制 offset enable-auto-commit:false max-poll-records:100 listener: # 消费者并发数不超过分区数 concurrency:6 ack-mode:manual说在最后Kafka 的分区有序只管到 Broker 存储这一层从 Producer 发送到 Consumer 处理完成有三个链路都可能乱序Producer 的并发重试、Consumer 的多线程处理、分区扩容导致 Key 映射漂移。从存储有序到处理有序之间的差距得靠自己的代码来补。