Skip to content

Apache Kafka 分区与分发

本页结论:Kafka 的「路由」是 key 到分区的哈希映射——它同时决定局部顺序与负载分布;消费组在分区之上做分配,组内瓜分、组间广播。

适用场景

  • 同业务键有序的事件流(key=orderId → 同分区 FIFO)。
  • 多消费者水平扩展:分区数是组内消费者的并行度上限。
  • 多个下游独立消费:每组一份完整数据。

核心模型:key → Partition → Consumer Group

  • 同一 key 的消息总是进同一分区 → 分区内顺序写入(ordering-replay 实验断言 samePartitionOnProduce=1)。
  • 无 key 的消息按 Sticky 策略分布到各分区(批量内轮转、批间尽量粘连),目的是摊匀负载并保留批内局部性。
  • 跨分区没有全局顺序——「Kafka 保证全局顺序」是禁止表述(规格 §7.2)。

动手验证:

bash
npm run lab -- kafka ordering-replay
verifiedkafka / ordering-replaybroker 4.3.1 · java-kafka-clients-4.3.1
镜像apache/kafka:4.3.1@sha256:77e3df9054047a88b520d0cc46e16696d3b22022e1d580aeccd2632df6532837
捕获时间2026-08-19T09:29:57.825Z
耗时 / 退出码23875 ms / exit 0
断言
produced6
samePartitionOnProduce1
samePartitionOnConsume1
observedOrder[ 1, 2, 3, 4, 5, 6 ]
replayed6
replayFromOffset00
replayUniqueMessageIds6
归一化日志
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 offset=0 seq=1 durationMs=<ms> status=produced
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-2 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 offset=1 seq=2 durationMs=<ms> status=produced
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-3 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 offset=2 seq=3 durationMs=<ms> status=produced
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-4 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-4 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 offset=3 seq=4 durationMs=<ms> status=produced
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-5 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-5 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 offset=4 seq=5 durationMs=<ms> status=produced
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-6 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-6 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 offset=5 seq=6 durationMs=<ms> status=produced
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay destination=orders.ordering produced=6 status=done
[assert] produced=6 PASS
[assert] samePartitionOnProduce=1 PASS
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay destination=orders.ordering consumerGroup=orders-ordering-g1 consumer=g1 partitions=0,1,2 status=assigned
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 offset=0 seq=1 consumerGroup=orders-ordering-g1 consumer=g1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 seq=1 consumerGroup=orders-ordering-g1 attempt=1 status=business_committed
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-2 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 offset=1 seq=2 consumerGroup=orders-ordering-g1 consumer=g1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-2 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 seq=2 consumerGroup=orders-ordering-g1 attempt=1 status=duplicate_skipped
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-3 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 offset=2 seq=3 consumerGroup=orders-ordering-g1 consumer=g1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-3 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 seq=3 consumerGroup=orders-ordering-g1 attempt=1 status=duplicate_skipped
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-4 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-4 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 offset=3 seq=4 consumerGroup=orders-ordering-g1 consumer=g1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-4 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-4 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 seq=4 consumerGroup=orders-ordering-g1 attempt=1 status=duplicate_skipped
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-5 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-5 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 offset=4 seq=5 consumerGroup=orders-ordering-g1 consumer=g1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-5 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-5 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 seq=5 consumerGroup=orders-ordering-g1 attempt=1 status=duplicate_skipped
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-6 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-6 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 offset=5 seq=6 consumerGroup=orders-ordering-g1 consumer=g1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-6 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-6 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 seq=6 consumerGroup=orders-ordering-g1 attempt=1 status=duplicate_skipped
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay destination=orders.ordering consumer=g1 received=6 status=done
[assert] samePartitionOnConsume=1 PASS
[assert] observedOrder=1,2,3,4,5,6 PASS
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay destination=orders.ordering consumerGroup=orders-ordering-g2 consumer=g2 partitions=0,1,2 status=assigned
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 offset=0 seq=1 consumerGroup=orders-ordering-g2 consumer=g2 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 seq=1 consumerGroup=orders-ordering-g2 attempt=1 status=business_committed
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-2 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 offset=1 seq=2 consumerGroup=orders-ordering-g2 consumer=g2 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-2 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 seq=2 consumerGroup=orders-ordering-g2 attempt=1 status=duplicate_skipped
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-3 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 offset=2 seq=3 consumerGroup=orders-ordering-g2 consumer=g2 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-3 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 seq=3 consumerGroup=orders-ordering-g2 attempt=1 status=duplicate_skipped
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-4 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-4 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 offset=3 seq=4 consumerGroup=orders-ordering-g2 consumer=g2 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-4 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-4 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 seq=4 consumerGroup=orders-ordering-g2 attempt=1 status=duplicate_skipped
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-5 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-5 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 offset=4 seq=5 consumerGroup=orders-ordering-g2 consumer=g2 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-5 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-5 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 seq=5 consumerGroup=orders-ordering-g2 attempt=1 status=duplicate_skipped
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-6 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-6 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 offset=5 seq=6 consumerGroup=orders-ordering-g2 consumer=g2 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay messageId=mid-6 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-6 correlationId=order-1001 destination=orders.ordering partitionOrQueue=0 seq=6 consumerGroup=orders-ordering-g2 attempt=1 status=duplicate_skipped
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=ordering-replay destination=orders.ordering consumer=g2 received=6 status=done
[assert] replayed=6 PASS
[assert] replayFromOffset0=0 PASS
[assert] replayUniqueMessageIds=6 PASS
npm run lab -- kafka ordering-replay

消费组分配与再均衡

  • 组内:每个分区至多分配给组内一个消费者;消费者数 > 分区数时,多出来的消费者空转(这就是「消费者数量上限 = 分区数」的关系)。
  • 成员变化(加入/退出/崩溃、订阅 Topic 分区数变化)触发 再均衡(rebalance):期间消费暂停,分区重新分配;已提交 offset 决定新属主从哪继续。
  • 本仓库 consumer-group 实验用两个消费者瓜分 3 分区,随后独立组 b 再次全量接收:
bash
npm run lab -- kafka consumer-group
verifiedkafka / consumer-groupbroker 4.3.1 · java-kafka-clients-4.3.1
镜像apache/kafka:4.3.1@sha256:77e3df9054047a88b520d0cc46e16696d3b22022e1d580aeccd2632df6532837
捕获时间2026-08-19T09:29:21.932Z
耗时 / 退出码135508 ms / exit 0
断言
produced3
groupAReceived3
groupAUnique3
a1Assignedtrue
a2Assignedtrue
a1ExitCode0
a2ExitCode0
groupBReceived3
groupBBusinessRows3
groupALag0
归一化日志
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.group partitionOrQueue=0 offset=0 seq=1 durationMs=<ms> status=produced
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1002 traceId=trace-2 correlationId=order-1002 destination=orders.group partitionOrQueue=0 offset=1 seq=2 durationMs=<ms> status=produced
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1003 traceId=trace-3 correlationId=order-1003 destination=orders.group partitionOrQueue=1 offset=0 seq=3 durationMs=<ms> status=produced
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group destination=orders.group produced=3 status=done
[assert] produced=3 PASS
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group destination=orders.group consumerGroup=orders-group-a consumer=a-1 partitions=0,1,2 status=assigned
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.group partitionOrQueue=0 offset=0 seq=1 consumerGroup=orders-group-a consumer=a-1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.group partitionOrQueue=0 seq=1 consumerGroup=orders-group-a attempt=1 status=business_committed
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1002 traceId=trace-2 correlationId=order-1002 destination=orders.group partitionOrQueue=0 offset=1 seq=2 consumerGroup=orders-group-a consumer=a-1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1002 traceId=trace-2 correlationId=order-1002 destination=orders.group partitionOrQueue=0 seq=2 consumerGroup=orders-group-a attempt=1 status=business_committed
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1003 traceId=trace-3 correlationId=order-1003 destination=orders.group partitionOrQueue=1 offset=0 seq=3 consumerGroup=orders-group-a consumer=a-1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1003 traceId=trace-3 correlationId=order-1003 destination=orders.group partitionOrQueue=1 seq=3 consumerGroup=orders-group-a attempt=1 status=business_committed
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group destination=orders.group consumerGroup=orders-group-a consumer=a-1 partitions=2 status=assigned
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group destination=orders.group consumer=a-1 received=3 status=done
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group destination=orders.group consumerGroup=orders-group-a consumer=a-2 partitions=0,1 status=assigned
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group destination=orders.group consumerGroup=orders-group-a consumer=a-2 partitions=0,1,2 status=assigned
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group destination=orders.group consumer=a-2 received=0 status=done
[assert] groupAReceived=3 PASS
[assert] groupAUnique=3 PASS
[assert] a1Assigned=true PASS
[assert] a2Assigned=true PASS
[assert] a1ExitCode=0 PASS
[assert] a2ExitCode=0 PASS
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group destination=orders.group consumerGroup=orders-group-b consumer=b-1 partitions=0,1,2 status=assigned
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.group partitionOrQueue=0 offset=0 seq=1 consumerGroup=orders-group-b consumer=b-1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.group partitionOrQueue=0 seq=1 consumerGroup=orders-group-b attempt=1 status=business_committed
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1002 traceId=trace-2 correlationId=order-1002 destination=orders.group partitionOrQueue=0 offset=1 seq=2 consumerGroup=orders-group-b consumer=b-1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1002 traceId=trace-2 correlationId=order-1002 destination=orders.group partitionOrQueue=0 seq=2 consumerGroup=orders-group-b attempt=1 status=business_committed
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1003 traceId=trace-3 correlationId=order-1003 destination=orders.group partitionOrQueue=1 offset=0 seq=3 consumerGroup=orders-group-b consumer=b-1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1003 traceId=trace-3 correlationId=order-1003 destination=orders.group partitionOrQueue=1 seq=3 consumerGroup=orders-group-b attempt=1 status=business_committed
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group destination=orders.group consumer=b-1 received=3 status=done
[inspect] timestamp=<ts> level=INFO service=order-service product=kafka lab=consumer-group business_rows=3 processed_rows=3 status=snapshot
[assert] groupBReceived=3 PASS
[assert] groupBBusinessRows=3 PASS
[assert] groupALag=0 PASS
npm run lab -- kafka consumer-group

分区数怎么选

  • 分区数 ≈ 目标消费并行度;扩容消费者只能到分区数为止,扩分区不可逆(且会打乱既有 key 分布)。
  • 分区过多会增加元数据、文件句柄与端到端延迟(Producer 按分区分批)。
  • 顺序需求决定 key 的选择(orderId 而非 userId 还是二者组合),热点 key 会造成单分区倾斜。

常见误区

  • 「key 相同就一定相邻消费」——只保证同分区顺序;消费端多线程拆分处理同一分区会破坏顺序。
  • 「增加消费者总能提速」——超过分区数后新消费者拿不到分区。
  • 「再均衡是无感的」——cooperative 协议减少了停顿,但 rebalance 仍意味着短暂的消费中断与位点重读风险(提交太早会重复)。
  • 「无 key 消息完全随机分布」——Sticky 分区器有批内粘连语义,不是纯随机轮转。

官方资料

以统一实验验证消息系统语义边界