Skip to content

顺序、消费组与回放(Kafka)

本页结论:Kafka 的顺序只在分区内成立——同 key 进同分区、组内瓜分分区、新消费组可从 earliest 全量回放;本页用三个实验分别验证,并给出与 RabbitMQ 的语义对照。

适用场景

  • 需要「同 orderId 事件按序处理」的业务流。
  • 验证消费组并行度与分区的关系。
  • 验证回放能力(新消费组/位点重置从头读)。

拓扑

实验一:同 key 有序(ordering-replay)

bash
npm run lab -- kafka ordering-replay

步骤:Producer 用同一 key(order-1001)发送 6 条带序号消息 → 断言 6 条落在同一分区samePartitionOnProduce=1)→ 消费组 g1 按序接收(observedOrder=[1..6])→ 新消费组 g2 以 auto.offset.reset=earliest 从 offset 0 全量回放(replayed=6replayFromOffset0=0)。

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

结论与边界:

  • 同 key → 同分区 → 分区内写入有序;消费端单线程读分区即保序。
  • 回放不改变日志本身:g1 的消费不影响 g2,offset 是组级位点。
  • 跨分区无顺序可言;消费端多线程处理同一分区会打乱顺序(见 分区与分发)。

实验二:消费组瓜分分区(consumer-group)

bash
npm run lab -- kafka consumer-group

步骤:先发 3 条到 3 分区 Topic → 组 A 两个消费者并行消费(空闲超时退出后合并统计:每条消息恰被组内一个消费者收到一次,两个消费者都观察到分区分配)→ 组 B 独立消费,同样收到全量 3 条。

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

结论与边界:

  • 组内瓜分:每个分区至多一个消费者;消费者数超过分区数时多余者空转。
  • 组间广播:位点按组独立,组 B 的接收与组 A 互不影响。
  • 分配的具体归属由协议决定(本实验断言「两者都有分配」而非具体归属,因为分配结果不保证可复现)。

与 RabbitMQ / Pulsar 的顺序语义对照

维度KafkaRabbitMQPulsar
顺序单位Partition(key 决定归属)单个 Queue(binding 决定归属)分区(key 路由)+ 订阅类型共同决定
同键有序怎么做key=orderId → 同分区routing key=orderId → 专属队列 + 单消费者key=orderId → 分区 + Key_Shared 订阅(同 key 粘连同一消费者)
并行与顺序的冲突分区数 = 并行上限,全局顺序需单分区队列数类似;单队列多消费者会竞争乱序分区数类似;Shared 订阅换并行但放弃顺序
失败后顺序无 requeue;失败消息通常转发 DLQ Topic,原分区继续NACK+requeue 会把消息送回队列,顺序可能抖动negativeAck/ack 超时触发重投,同分区内重试消息与后续消息的顺序会被打乱
回放原生:位点重置/新组 earliest不适用(ACK 即删;Streams 除外)原生:reset-cursor 到 earliest/时间戳(redelivery-replay 实验验证)
多订阅多消费组各自位点多队列各自绑定同一 topic 可并存多个订阅(Exclusive/Shared/Failover/Key_Shared),各自独立游标

统一结论(对应 顺序语义):三家都只提供局部顺序;「全局顺序」需要牺牲并行度,且都要在消费端保持单线程处理同一顺序单位。Pulsar 的额外维度是订阅类型:Shared 换吞吐但无跨消费者顺序,Key_Shared 才能同时兼得同键有序与并行(见 subscriptions 实验)。

断言汇总

实验关键断言期望
ordering-replaysamePartitionOnProduce / observedOrder / replayed / replayFromOffset01 / [1..6] / 6 / 0
consumer-groupgroupAReceived / groupAUnique / a1Assigned+a2Assigned / groupBReceived / groupALag3 / 3 / 均有 / 3 / 0

官方资料与版本说明

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