主题
Kafka 核心概念:Topic、Partition 与 Consumer Group
本文是消息队列系统学习系列的 L2 核心篇。前置:RabbitMQ 入门:Exchange、路由与可靠投递。 学完可以配合面试题食用:05-kafka-consumer-group-rebalance、07-kafka-partition-assignment-strategies
为什么需要 Kafka 这一套抽象
RabbitMQ 用 Exchange 和 Binding 做灵活路由,适合业务系统内部的消息流转。但当日志、埋点、流式计算这类数据量上亿/天的场景出现时,RabbitMQ 的单机吞吐瓶颈就暴露了——它把消息都堆在内存里的队列,磁盘 IO 也是随机写入。
Kafka 做了两个根本不同的设计选择:日志(append-only) 代替队列,分区(Partition) 代替单兵作战。数据按顺序追加到磁盘,读的时候走 page cache + 零拷贝,吞吐量从万级跳到百万级。代价是路由能力弱了——没有 Exchange 级别的灵活匹配,投递全靠 topic 名硬匹配。
Topic → Partition → Replica → Segment:四层物理结构
Kafka 的数据组织是四层嵌套的。从顶往下看:
mermaid
graph TD
Topic["Topic<br/>(逻辑分类)"] --> P0["Partition 0"]
Topic --> P1["Partition 1"]
Topic --> P2["Partition 2"]
P0 --> Leader["Leader Replica"]
P0 --> F1["Follower Replica 1"]
P0 --> F2["Follower Replica 2"]
Leader --> Seg["Segment 文件序列<br/>(.log + .index + .timeindex)"]
Seg --> Msg["消息体<br/>(offset + key + value + 元数据)"]Topic 是逻辑容器,一个 Topic 对应一类数据流(比如"订单事件")。Partition 是物理存储单元,每个 Partition 独立追加、独立管理 offset。Partition 的副本(Replica)分布在不同的 broker 上,Leader 负责读写,Follower 异步拉取同步。每个 Partition 的日志文件再切分成若干 Segment,默认每 1GB 或每 7 天滚动一个新 Segment。
这样一个 Topic 如果有 6 个 Partition、3 副本,集群里就有 18 个 Replica,读写在 6 个 Leader 上并行。这是 Kafka 水平扩展的基石。
Consumer Group 语义:组内竞争,组间广播
Consumer Group 是 Kafka 最核心的消费抽象。同一个 Group 内的消费者竞争消费 Partition 里的消息——每条消息只会被组内一个消费者处理。不同 Group 则独立消费同一条消息,像是广播。
mermaid
graph LR
subgraph Broker
P0["Topic-X / Partition-0"]
P1["Topic-X / Partition-1"]
P2["Topic-X / Partition-2"]
end
subgraph Group-A["Group A (订单处理)"]
C1["Consumer A1"]
C2["Consumer A2"]
end
subgraph Group-B["Group B (数据仓库)"]
C3["Consumer B1"]
end
P0 --> C1
P1 --> C2
P2 --> C2
P0 --> C3
P1 --> C3
P2 --> C3关键约束:一个 Partition 只能被同一个 Group 内的一个消费者消费。所以如果你的消费者数多于分区数,多出来的消费者永远闲置。
Rebalance 是 Consumer Group 的再分配机制。当组成员增减、分区数变化时,Kafka 触发一轮 Rebalance——所有消费者 stop 消费,协调器(Group Coordinator)重新分配分区,消费者再按新分配开始拉取。这期间消费中断,短则几百毫秒,长则几十秒(取决于堆大小和 GC)。Rebalance 频繁是生产环境最常见的坑之一,后面面试题 05 会专门拆解。
Offset 管理:从哪里读,为什么重复
每个 Partition 的消息按写入顺序排号,这个号就是 offset(逻辑位移,从 0 开始递增)。消费者需要知道"我已经读到哪了",这个位置信息也存 Kafka 内部——__consumer_offsets 这个内部 Topic。
提交 offset 的方式有三种:
- 自动提交(enable.auto.commit=true):默认每 5 秒自动提交当前 poll 返回的最大 offset。简单但容易重复——如果消费者在两次提交之间挂了,重启后会从上一次提交的 offset 继续读,这 5 秒内的消息就被重复消费了。
- 手动同步提交(commitSync):当前批次处理完手动提交,阻塞直到 broker 确认。最安全,但吞吐会受影响。
- 手动异步提交(commitAsync):不阻塞,回调处理失败。吞吐高,但失败时可能提交了更大的 offset,下一次重启丢消息。
实际生产中常见的做法:手动异步提交 + 周期性同步提交兜底。
Partition 与并行度的关系
Partition 数就是 Kafka 最大并行度。一个分区同一时刻只能被一个消费者读,所以想让消费速度翻倍,就加倍分区数(和消费者数)。但反过来,分区数多了也有代价:
- 更多文件句柄(每个 Partition 对应一组 Segment 文件)
- 更多 Rebalance 开销(协调器要管理更多分区)
- 更高的端到端延迟(Leader 要维护更多 Follower 的同步)
经验值:分区数 = broker 数 × 2~4,或者按峰值 TPS / 单分区吞吐量反推。单分区写吞吐约 10 MB/s(取决于消息大小和硬件),单分区读吞吐更高。
Broker 与 Zookeeper / KRaft
早期 Kafka 依赖 Zookeeper 存储元数据(集群信息、Leader 选举等)。ZooKeeper 写性能差,当集群规模大时元数据变更(如新增 Topic 引发大量分区 Leader 选举)会成为瓶颈。
Kafka 2.8 引入 KRaft(Kafka Raft),把元数据管理移到 Kafka 自身,去掉 ZK 依赖。KRaft 用 Raft 共识协议在 Controller 节点间同步元数据,元数据变更就变成 Kafka 内部 Topic 的 append-only 操作,一致性更好。生产环境 3.0+ 建议直接上 KRaft 模式,省掉 ZK 的运维成本。细节见面试题 21。
动手实操
命令行实验:起一个单机 Kafka 并收发消息
bash
# 下载 Kafka(假设已安装,版本 3.x)
cd ~/kafka_2.13-3.7.0
# 启动 KRaft 模式(单节点)
# 1. 格式化日志目录(首次启动)
KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)"
bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/kraft/server.properties
# 2. 启动 broker
bin/kafka-server-start.sh config/kraft/server.properties &
# 3. 创建 Topic(3 个分区,1 副本)
bin/kafka-topics.sh --create --topic test-topic \
--bootstrap-server localhost:9092 \
--partitions 3 --replication-factor 1
# 4. 描述 Topic 看分区信息
bin/kafka-topics.sh --describe --topic test-topic \
--bootstrap-server localhost:9092
# 输出示例:
# Topic: test-topic PartitionCount: 3 ReplicationFactor: 1
# Topic: test-topic Partition: 0 Leader: 0 ...
# Topic: test-topic Partition: 1 Leader: 0 ...
# Topic: test-topic Partition: 2 Leader: 0 ...
# 5. 开一个控制台消费者
bin/kafka-console-consumer.sh --topic test-topic \
--bootstrap-server localhost:9092 --group my-group
# 6. 另一个终端开生产者,打字发消息
bin/kafka-console-producer.sh --topic test-topic \
--bootstrap-server localhost:9092
# >hello kafka
# >partition 2 message
# Ctrl+C 退出Java 生产消费 20 行版
java
import org.apache.kafka.clients.producer.*;
import org.apache.kafka.clients.consumer.*;
import java.time.Duration;
import java.util.Properties;
import java.util.List;
// 生产者
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("acks", "all"); // 等待所有副本确认
Producer<String, String> producer = new KafkaProducer<>(props);
producer.send(new ProducerRecord<>("test-topic", "my-key", "hello kafka"));
producer.close();
// 消费者(手动提交示例)
Properties cprops = new Properties();
cprops.put("bootstrap.servers", "localhost:9092");
cprops.put("group.id", "java-consumer-group");
cprops.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
cprops.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
cprops.put("enable.auto.commit", "false"); // 手动控制
Consumer<String, String> consumer = new KafkaConsumer<>(cprops);
consumer.subscribe(List.of("test-topic"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
for (ConsumerRecord<String, String> record : records) {
System.out.printf("offset=%d, key=%s, value=%s%n",
record.offset(), record.key(), record.value());
}
consumer.commitSync(); // 处理完一批提交一次
}这段代码跑起来就能看到消费者在三个分区之间分配消息。开两个消费者实例(同一个 group.id),观察分区分配怎么变化。
常见误区与小结
- 分区数越多越好:不是。每个分区有文件句柄和内存开销,分区数超过 1000 后 Rebalance 超时和 Leader 选举延迟会显著上升。
- 消费者数超过分区数也能加速:不能,多出来的消费者永远闲置。加机器前先检查分区数。
- 自动提交是偷懒好选择:重复消费概率高。生产环境至少用手动异步提交 + 定时同步兜底。
- 不同 Group 的消费者可以共享 Partition:可以,每个 Group 独立维护 offset,互不干扰。
- Kafka 保证分区内有序,不保证跨分区有序:如果需要全局有序,只能用一个分区(牺牲吞吐)。业务上通常用关键 key 路由到同一个分区就够了。
小结:Topic/Partition/Consumer Group 是 Kafka 的三大基石。Partition 决定了吞吐上限和并行度,Consumer Group 决定了消费隔离和广播范围。理解了这三层,再看 Kafka 的可靠投递和高性能设计就顺了。
下一篇进入生产者与消费者内部原理:拦截器、分区器、累加器、Sender 线程和 Rebalance 全流程。
参考
参考:Kafka 官方文档 (https://kafka.apache.org/documentation/)、《Kafka 权威指南》第二部分