跳至正文
来两杯美式
返回

Kafka(下):消费者、消息有序性、命令与集群搭建

By 来两杯美式
发布于更新于

Kafka 系列收尾篇。前两篇讲完「是什么」和「为什么快」,本篇进入怎么用、怎么运维——Consumer Group 怎么分组、消息怎么保证有序、为什么会重复消费、本地怎么起集群、设计 MQ 该考虑什么。

一、Consumer Group 与 groupid

1.1 groupid 的作用

Kafka 的每个 topic,可以被不同 groupid 的消费者单独消费——每个 group 独立记录自己的 offset,互不影响。

Topic: order-events (3 partitions)
   ├── Partition 0
   ├── Partition 1
   └── Partition 2

Group A (独立消费全量)        Group B (独立消费全量)
   ├── Consumer A1              ├── Consumer B1
   ├── Consumer A2              └── Consumer B2
   └── Consumer A3

1.2 消费者数量与 partition 数量的关系

通常情况下,一个 groupid 里的消费者数量,最好等同于一个 topic 下的 partition 数量,以保证最好的性能。

Consumer 数 vs Partition 数行为
consumer = partition每个 consumer 拿一个 partition,最佳
consumer < partition部分 consumer 消费多个 partition
consumer > partition多余的 consumer 闲置(资源浪费)

关键认知一个 partition 在一个 group 内只能被一个 consumer 消费。这是 Kafka 实现消费有序性的基础,也是为什么 partition 数直接决定消费端并行上限。

1.3 Rebalance(消费者重平衡)

当 consumer 加入或退出 group 时,会触发 rebalance

Rebalance 是 Kafka 消费端最常见的「短暂停顿」来源——rebalance 期间整个 group 暂停消费。生产环境调优的重点是减少不必要的 rebalance(如 session timeout 配得太短、poll 间隔过大都会触发)。

二、消息有序性

2.1 单 partition 是有序的

Kafka 的单 partition 是有序的——消息按写入顺序追加,消费者按 offset 顺序消费。

2.2 全局有序的代价

如果真的把 topic 的 partition 设置为单个,就失去了 Kafka 并行的优势

单 partition:全局有序 ✅,但吞吐量退化为单线 ❌
多 partition:吞吐量恢复 ✅,但跨 partition 无序 ❌

2.3 业务级有序:key + offset

实战解法:用 kafka key + offset 实现业务有序性:

Producer 发送时指定 key:
   msg(key=order_123, action=created)  ┐
   msg(key=order_123, action=paid)     ├─ 同一 partition
   msg(key=order_123, action=shipped)  ┘   按顺序消费 → 订单 123 的状态有序

   msg(key=order_456, action=created)  ─► 另一个 partition

生产推荐做法:用业务 ID 做 key,让 Kafka 自动按 hash 分配 partition——同一业务实体的所有事件都进同一 partition,自然有序。比单 partition 全局有序高效几个数量级。

三、为什么会重复消费

本质上讲,Kafka 的消息重复消费是 offset 控制出了问题

三大原因:

3.1 人为原因:Consumer 使用不当

→ 下次启动,会从上次提交的 offset 重新消费,于是这条消息被消费 2 次。

3.2 程序原因

→ consumer 被 group 踢出,rebalance 后 partition 转给别人,对方从未提交 offset 的位置开始消费

3.3 消费者 Rebalance 时 offset 未控制好

3.4 怎么解决

方案适用场景
关闭自动提交,改手动提交生产环境标配
业务幂等即使重复消费也无副作用(最稳
唯一 ID + 去重表不能改业务代码时
Transactions(exactly-once)Kafka 2.5+,强一致场景

铁律MQ 的「至少一次」语义是事实标准——不要幻想「不会重复」。业务必须做幂等,把重复消费当成常态来设计。

四、本地启动 Kafka

4.1 单机模式

# 进入 kafka 安装目录,先启动 zookeeper
bin/zookeeper-server-start.sh config/zookeeper.properties

# 启动 kafka
bin/kafka-server-start.sh config/server.properties

ZooKeeper 是 Kafka 早期版本的元数据存储 + 选主组件。生产环境 ZK 集群建议 ≥ 3 节点(奇数,避免脑裂)。Kafka 2.8+ 引入 KRaft 模式,可以不用 ZK

4.2 简单 CLI 命令

# 创建一个名字为 "test" 的 topic
bin/kafka-topics.sh --create --zookeeper localhost:2181 \
  --replication-factor 1 --partitions 1 --topic test

# 查看 topic 列表
bin/kafka-topics.sh --list --zookeeper localhost:2181

# 查看 topic 描述
bin/kafka-topics.sh --describe --zookeeper localhost:2181 --topic test

# 向 topic 发送消息(启动交互式 producer)
bin/kafka-console-producer.sh --broker-list localhost:9092 --topic test

# 消费消息(启动交互式 consumer)
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
  --topic test [--from-beginning]

--from-beginning 让 consumer 从最早的 offset 开始消费(默认从最新的 offset)。调试时常用

4.3 新版 CLI 参数差异(2026 视角补充)

Kafka 3.x 起,部分命令的 --zookeeper 参数被废弃,改为 --bootstrap-server

# 旧(2.x)
bin/kafka-topics.sh --create --zookeeper localhost:2181 ...

# 新(3.x+)
bin/kafka-topics.sh --create --bootstrap-server localhost:9092 \
  --partitions 3 --replication-factor 1 --topic test

五、搭建集群

5.1 准备多个配置文件

cp config/server.properties config/server-1.properties
cp config/server.properties config/server-2.properties

编辑这两个新文件,设置以下属性:

# config/server-1.properties
broker.id=1
listeners=PLAINTEXT://:9093
log.dir=/tmp/kafka-logs-1

# config/server-2.properties
broker.id=2
listeners=PLAINTEXT://:9094
log.dir=/tmp/kafka-logs-2

分别启动三个 broker:

bin/kafka-server-start.sh config/server.properties       # broker.id=0
bin/kafka-server-start.sh config/server-1.properties     # broker.id=1
bin/kafka-server-start.sh config/server-2.properties     # broker.id=2

5.2 创建有副本的 topic

bin/kafka-topics.sh --create --zookeeper localhost:2181 \
  --replication-factor 3 --partitions 1 --topic replicated-topic

5.3 查看新 topic 的描述

bin/kafka-topics.sh --describe --zookeeper localhost:2181 --topic replicated-topic

输出:

Topic: replicated-topic   PartitionCount: 1   ReplicationFactor: 3   Configs:
Topic: replicated-topic   Partition: 0   Leader: 1   Replicas: 1,2,0   Isr: 1,2,0

5.4 字段解读

字段含义
leader该节点负责该分区的所有读写(每个分区的 leader 是随机选择的)
replicas备份节点列表(无论该节点是否是 leader、是否还活着,都列出来)
isrIn-Sync Replicas——「同步备份」节点列表,活着的 + 正在同步 leader 的
Leader: 1   Replicas: 1,2,0   Isr: 1,2,0
            │         │       │
            │         │       └─ 当前健康在同步的副本
            │         └─ 配置的所有副本(含可能掉线的)
            └─ 当前分区的主节点

ISR 是 Kafka 高可用的核心——acks=all 实际等的是「所有 ISR 副本确认」。某个 replica 长时间不同步会被踢出 ISR,leader 挂掉时只有 ISR 内的副本有资格当选。

六、设计一个 MQ 该考虑什么

最后一节——如果你自己设计一个 MQ,会考虑哪些?这是面试常见开放题,也是检验是否真的理解了 MQ 的试金石。

维度关键决策
1. 定位是流处理(Kafka)还是业务消息(RabbitMQ)?定位决定取舍
2. 产品设计目标受众 + 特性选择——RocketMQ 选了死信 + 事务消息,Kafka 选了吞吐量
3. 吞吐量目标充分考虑前两点后预估——前两者都会对性能有损耗
4. 集群模式主从 / Control / 集群——可用性 + 易用性的保障
5. 客户端模式 & 易用性SDK 多语言、API 简洁度
6. 数据安全和加密是否支持 TLS、SASL、ACL

三大主流 MQ 的取舍很有意思:

  • Kafka:定位是流处理平台,舍弃了延迟队列、事务消息等 MQ 标配,换来了吞吐量
  • RocketMQ:定位是企业级 MQ,提供了死信队列 + 分布式事务消息,更符合当下企业诉求
  • RabbitMQ:定位是经典消息队列,强在路由灵活性 + 协议完备性,但 Erlang 限制了它的扩展性

没有银弹——选 MQ 的本质是选一组最适合业务的取舍

七、整个 Kafka 系列回顾

「一文拿下 Kafka」3 篇覆盖了 Kafka 的核心:

主题重点
特性、场景、消息传递保障流处理定位、acks 三档语义
吞吐量为什么这么大日志、partition、批量、零拷贝
消费者、有序性、命令、集群groupid、rebalance、CLI、ISR

整个「MQ 系列」到这一篇全部收尾——RabbitMQ 4 篇 + Kafka 3 篇 = 7 篇。从「为什么用 MQ」到「怎么选型」,从「RabbitMQ 怎么集成」到「Kafka 为什么快」,再到「设计 MQ 该考虑什么」,闭环完成。

现代视角补一句(2026):Kafka 的核心知识到 2026 年几乎没变——partition、offset、ISR、Consumer Group 这些概念从 2014 年到今天一脉相承。变的只是:KRaft 去 ZK、Tiered Storage、Transactions 增强、客户端 SDK 性能优化。这套笔记的内核,再放 10 年依然适用。


至此,从 JVM(5 篇 + 锁篇)→ GC(5 篇)→ MQ(Rabbit 4 + Kafka 3 = 7 篇),2017~2018 年的基础笔记全部迁移完成。下一站:可能是 Redis、MySQL、DDD 系列(如果还没迁的话),或者开始写新内容。看心情。


分享这篇文章:
通过邮件分享这篇文章✓ 链接已复制
所属专题
Kafka
第 3 / 3 篇
查看系列全部文章
  1. 01.Kafka(上):特性、应用场景与消息传递保障
  2. 02.Kafka(中):吞吐量为什么这么大——日志、零拷贝、批量
  3. 03.Kafka(下):消费者、消息有序性、命令与集群搭建

上一篇
Dubbo 入门:什么是 RPC 与服务治理?
下一篇
Kafka(中):吞吐量为什么这么大——日志、零拷贝、批量