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
- 不同 groupid:各自完整消费一遍(互相独立,类似广播)
- 相同 groupid:组内分摊消费(每条消息只被组内一个 consumer 消费一次)
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:
- 新 consumer 加入 → partition 重新分配
- consumer 宕机 → 它的 partition 转交给组内其他 consumer
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 实现业务有序性:
- key 是一个有特定业务含义的字段(如订单 ID、用户 ID)
- 具有相同 key 的消息会被路由到同一个 partition
- 在该 partition 内,按 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
→ 下次启动,会从上次提交的 offset 重新消费,于是这条消息被消费 2 次。
3.2 程序原因
- 程序消费消息时,依赖了外部资源不能申请到(如 DB 连接打满)
- 死锁、锁等待导致 session 超时
- offset 不能及时上报
→ consumer 被 group 踢出,rebalance 后 partition 转给别人,对方从未提交 offset 的位置开始消费。
3.3 消费者 Rebalance 时 offset 未控制好
- 新的消费者加入消费者组
- 有消费者离组(宕机 / 手动关闭)
- 这时如果 offset 还没提交,新接管方会从老 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 字段解读
- 第一行:所有分区的摘要
- 后面每一行:一个分区的详细信息(这里只有 1 个 partition,所以只有一行)
| 字段 | 含义 |
|---|---|
| leader | 该节点负责该分区的所有读写(每个分区的 leader 是随机选择的) |
| replicas | 备份节点列表(无论该节点是否是 leader、是否还活着,都列出来) |
| isr | In-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 系列(如果还没迁的话),或者开始写新内容。看心情。