【Kafka 】从入门到生产:一篇搞懂消息队列的核心原理
一、消息队列是干啥的?为什么选 Kafka?
1.1 消息队列的本质
消息队列(Message Queue)本质上是个中转站。你的系统A要发数据给系统B,直接调接口行不行?行,但问题会越来越多:
- A发得太快,B处理不过来,数据直接冲垮B
- B挂了,A发的数据全丢了
- 后来又要加个子系统C也要这份数据,A的代码得改
消息队列就是来解决这些破事的。A只管往队列里丢,B和C只管按需取,两边解耦,互不耽误。
核心模型就三个角色:
- Producer(生产者):发消息的
- Consumer(消费者):接消息的
- Broker(服务端):存消息的,Kafka集群里的每台机器都是一个Broker
1.2 两种消费模式
| 模式 | 特点 | 代表 |
|---|---|---|
| 点对点(P2P) | 一条消息只被一个消费者消费,消费完删除 | RabbitMQ的Queue |
| 发布订阅(Pub/Sub) | 一条消息可被多个消费者组同时消费,消息不删 | Kafka就是这种模式 |
Kafka用发布订阅模式,所以消息会被保留一段时间(默认7天),不同消费者组可以独立消费同一份数据。这个设计特别适合大数据场景——同一份日志,既要给实时流处理(Flink/Spark),又要给离线批处理(Hive/HDFS),各取所需互不干扰。
1.3 为什么Kafka这么猛?
Kafka最初是LinkedIn那帮人为了处理海量日志搞出来的,后来捐给了Apache。它火起来有几个硬实力:
- 吞吐量变态高:单机每秒几十万条消息是基本功
- 数据不丢:多副本 + 持久化,配置对了数据就是安全的
- 水平扩展:加机器就能扩容,几乎不用停服务
- 生态成熟:大数据圈子里,Kafka几乎是消息队列的事实标准
实际生产链路基本是:业务系统 → Flume/Logstash → Kafka → Flink/Spark Streaming → 存储/告警
二、Kafka 架构全貌(Kafka 4.0 已彻底告别 Zookeeper)
这是重点中的重点,很多老教程还在讲Zookeeper,但时代变了。
2.1 KRaft模式:Kafka自己管自己
先上一张整体架构图:
老版本的Kafka强依赖Zookeeper做几件事:Broker注册、选举Controller、存Topic配置、消费者Offset。问题是维护两套分布式系统(Kafka + ZK)太痛苦了,ZK本身就是性能瓶颈和故障源。
所以Kafka社区搞了个KRaft模式(Kafka Raft),从2.8开始预览,3.3生产可用,到4.0正式版已经完全移除Zookeeper支持。现在新集群直接上KRaft就行。
KRaft的核心变化:
| 对比项 | 老版本(依赖ZK) | 新版本(KRaft) |
|---|---|---|
| 元数据存储 | 外部Zookeeper | 内置__cluster_metadata主题 |
| 控制器选举 | ZK选举Controller | Raft共识协议自选举 |
| 故障恢复 | 分钟级(从ZK拉全量元数据) | 秒级(本地快照+增量日志) |
| 运维复杂度 | 高,两套系统 | 低,只维护Kafka |
| 分区上限 | 万级(ZK瓶颈) | 百万级 |
换个方式理解:Kafka现在自己管自己了,不需要外部"管家",架构更干净,跑得更快,运维更省心。
踩坑点:如果你还在用Kafka 2.x甚至更老的版本,升级路径是先升到2.8+,再迁到3.x的KRaft。直接从1.x跳到4.0会出事!!!
2.2 核心组件一览
- Topic(主题):消息的分类,比如"order-events"、“user-logs”
- Partition(分区):Topic的数据被切分成多个分区,分区是Kafka并行处理的基本单位
- Replica(副本):每个分区可以有多个副本(Leader + Follower),保证数据安全
- Consumer Group(消费者组):一组消费者共同消费一个Topic,组内每条消息只被消费一次
- Offset(偏移量):每条消息在分区里的位置编号,消费者用这个标记自己读到哪了
三、分区与副本:Kafka高吞吐的根基
3.1 为什么要有分区?
想象一个Topic有100万条消息/秒的写入量,单台机器根本扛不住。分区就是把一个Topic的数据拆成多份,散到不同Broker上,每台机器只处理一部分,压力均摊。
同时,多个消费者可以并行消费不同分区,消费速度也线性提升。
3.2 分区策略:消息该进哪个分区?
生产者发消息时,Kafka有四种方式决定消息进哪个分区:
- 指定分区号:你明确说"这条进0号分区",直接进
- 有Key但没指定分区:
key的hash值 % 分区总数,同一个Key总会落到同一个分区。这个特性很有用,比如你想保证同一个用户ID的消息有序,就把用户ID当Key - 既没Key也没指定分区:轮询(Round Robin),挨个分区分摊
- 自定义分区器:实现Partitioner接口,按业务逻辑自己决定
踩坑点:如果用了Key做分区,后来Topic扩容(增加分区数),
hash(key) % 新分区数的结果会变,导致同一个Key的消息分散到不同分区,有序性被破坏。所以一开始就规划好分区数,尽量别扩容。真扩了的话,只能建新Topic迁移数据。
3.3 副本机制:数据怎么保证不丢?
每个分区可以配多个副本(replication.factor),比如3副本意味着数据存3份。副本分两种角色:
- Leader副本:唯一负责读写请求的,生产者和消费者都只跟Leader打交道
- Follower副本:从Leader拉取数据做备份,Leader挂了的时候顶替上去
ISR(In-Sync Replicas) 是副本机制的核心概念。ISR列表里放的是跟Leader保持同步的副本集合。一个副本要留在ISR里,必须满足:
- 与Leader保持心跳(
replica.lag.time.max.ms,默认30秒)
注意,老版本还有个replica.lag.max.messages参数(消息数差值限制),但Kafka 2.5+已经移除了这个参数,只保留时间判断。因为消息数阈值在生产环境中很难调——流量高峰时正常副本也可能被踢出ISR,造成不必要的抖动。
Leader挂了怎么办? 从ISR里选一个新的当Leader。如果ISR全空了(比如所有副本都挂了),Kafka会等(unclean.leader.election.enable=false,默认值就是false),宁可不可用也不丢数据。如果你把这个参数开成true,允许选非ISR副本当Leader,数据可能会丢,慎重。
踩坑点:线上曾经出现过某个分区ISR只剩Leader,Follower因为网络抖动被踢出,然后Leader所在机器挂了,整个分区不可用。查了半天发现是Follower机器磁盘满了导致同步卡住。所以监控每个Broker的磁盘水位是必须的,90%就要告警。
四、消息是怎么存到磁盘上的?
4.1 Segment + Index:高效存储的秘密
很多人以为Kafka快是因为纯内存,错了。Kafka的消息是持久化到磁盘的,但它把磁盘用出了内存的感觉。
每个分区在磁盘上是一个目录,里面被切成一段一段的Segment文件(默认1GB一个):
.log文件:存实际的消息数据,只追加不写改(Append Only).index文件:稀疏索引,存(relative offset → physical position)的映射.timeindex文件:按时间戳的索引(用于按时间查找)
查找offset=500的消息时:
- 用文件名二分定位到哪个Segment(比如
00000000000000000000.log存的是offset 0~368767) - 进对应的.index文件二分查找,找到不大于500的最大索引项(比如offset 368 → position 23415)
- 从.log文件的23415字节位置开始,顺序扫描找到offset=500的消息
稀疏索引的好处是索引文件很小,不会占太多内存。顺序扫描几百条消息对磁盘来说轻而易举。
4.2 Kafka为什么这么快?三个关键技术
从图里也能看到,Kafka高性能依赖操作系统层面的三个底层技术:
① 顺序写磁盘
机械磁盘最怕随机读写(磁头来回寻道),但顺序写磁盘的速度接近内存随机写。Kafka只追加写日志,纯顺序IO,直接把磁盘性能拉满。顺序写能做到600MB/s,随机写可能只有100KB/s,差了6000倍。
② PageCache 页缓存
Kafka并不自己做缓存,而是依赖操作系统的PageCache。消息先写到PageCache(内存),后台异步刷盘。消费者读消息时,如果数据还在PageCache里,直接从内存读,根本不走磁盘。
这带来一个巨大的好处:Kafka进程重启了,PageCache还在(因为是内核维护的),性能不受影响。
调参建议:Kafka是吃内存的大户,但不是吃JVM堆内存——它主要依赖操作系统的PageCache。所以JVM堆内存不需要配太大,**一般610G足够覆盖大部分场景**(几百个分区、几十GB数据量级)。如果集群规模特别大(几千个分区),可以适当加到1216G,但超过这个数通常意味着GC压力变大,收益反而下降。剩下的物理内存全部留给OS做PageCache。我见过有人给Kafka配了64G堆内存,结果GC一直卡,吞吐量反而上不去。
③ 零拷贝 Zero-Copy
传统方式把文件从磁盘发到网络,要经过4次数据拷贝(磁盘→内核缓冲区→用户缓冲区→内核Socket缓冲区→网卡),还要4次上下文切换。
Kafka用sendfile()系统调用,直接从PageCache发送到网卡,数据根本不进用户空间,CPU完全不参与拷贝,从4次变成1次。这就是为什么Kafka能做到磁盘读取速度几乎等于网卡带宽。
注意:开SSL/TLS加密或者开启Exactly-Once(事务消息)会打破零拷贝,因为数据必须进用户空间做加解密和校验。所以高吞吐场景下要权衡。
五、数据一致性:HW与LEO
5.1 什么是HW和LEO?
- LEO(Log End Offset):每个副本最后一条消息的offset+1。Leader和各个Follower的LEO可能不一样
- HW(High Watermark)高水位线:所有ISR副本中最小的LEO。消费者只能读到HW之前的消息
5.2 为什么要有HW?
看上面的图就明白了:
- Leader(LEO=4)已经写到了message 3
- Follower1(LEO=3)同步到了message 2
- Follower2(LEO=2)同步滞后,只到message 1
此时HW=2,意味着消费者只能读到message 0和message 1。message 2虽然Leader已经有了,但还没被所有ISR副本同步完,不算"已提交"。
如果这时候Leader挂了,Follower1被选为新Leader,Follower2同步到新Leader的数据。但消费者之前读到的message 0~1是安全的,因为所有副本都有。HW的本质就是保证:消费者读到的数据,一定不会因为Leader切换而丢失。
六、生产者:怎么发消息才靠谱?
6.1 批量发送机制
Kafka生产者默认就是批量发送的,不能关闭。流程是:
- 你调用
send()发一条消息 - 消息先进
RecordAccumulator(一个内存缓冲区,默认32MB) - 被攒成一个个Batch(默认16KB)
- 独立的Sender线程把Batch发送到Broker
触发发送的条件(满足任一):
- 一个Batch满了(16KB)
linger.ms超时(默认5ms,“等一等,看有没有更多消息一起发”)- 缓冲区满了
- 你手动调了
flush()
调参建议:高吞吐场景可以调大
linger.ms到10~50ms,牺牲一点延迟换更高的吞吐量。但如果你的业务对延迟敏感(比如实时风控),保持默认值5ms或更小。
6.2 ACK机制:三种安全级别
生产者发送消息后,Broker要不要给确认?Kafka给了三个选项:
| acks值 | 行为 | 吞吐 | 安全性 |
|---|---|---|---|
| 0 | 发完就忘,不等确认 | 最高 | 可能丢数据,不推荐 |
| 1 | 等Leader确认就返回 | 中高 | Leader挂了会丢 |
| all/-1 | 等ISR中所有副本都确认 | 最低 | 最安全,不丢数据 |
生产环境建议:
- 日志采集等对丢失不敏感的场景:
acks=1 - 订单、支付等核心场景:
acks=all(等所有副本确认) +min.insync.replicas=2(至少2个副本写入成功才算)
踩坑点:只配
acks=all不够,如果ISR只剩Leader一个,all实际上只等Leader确认(跟acks=1没区别)。必须配合min.insync.replicas(最小同步副本数)一起用。比如3个副本配min.insync.replicas=2,意思是至少2个副本写入成功,这次写入才算成功。
6.3 幂等性与事务消息
Kafka 0.11引入了幂等性生产者。开启方式很简单:
props.put("enable.idempotence", "true"); // 默认就是true(Kafka 3.x+)
开启后,Kafka会给每个生产者分配一个PID(Producer ID),每条消息带一个单调递增的Sequence Number。Broker会自动去重,如果收到重复的消息(比如网络超时生产者重试了),直接丢弃。
但幂等性只解决单分区单会话的问题。如果你需要跨分区的事务(比如"转账扣款和记账要么都成功要么都失败"),要用事务API:
producer.initTransactions();
try {
producer.beginTransaction();
producer.send(record1);
producer.send(record2);
producer.commitTransaction(); // 一起提交
} catch (Exception e) {
producer.abortTransaction(); // 一起回滚
}
注意:事务消息有性能开销(约10~20%),会打破零拷贝。非必要别开。
七、消费者与消费者组
7.1 消费者组的核心规则
一个分区,同一时间只能被一个消费者组里的一个消费者消费。这是Kafka消费者模型的铁律。
所以消费者个数和分区数的关系:
- 消费者数 = 分区数:完美一对一,最高并行度
- 消费者数 < 分区数:部分消费者处理多个分区
- 消费者数 > 分区数:多余的消费者闲着没事干,浪费
最佳实践:消费者数 = 分区数,这样能达到最大并行度。
7.2 Offset管理:消费到哪了?
消费者需要记录自己读到哪了,这个标记就是Offset。Offset存在哪?
- Kafka 0.8以前:存在Zookeeper里(时代的眼泪)
- Kafka 0.8以后:存在Kafka内部主题
__consumer_offsets里(50个分区,自动创建)
Offset的提交方式:
① 自动提交(默认)
props.put("enable.auto.commit", "true");
props.put("auto.commit.interval.ms", "5000"); // 每5秒自动提交一次
简单但有问题:先消费了消息,再提交Offset。如果消费完消息还没提交就挂了,下次重启会重复消费。
② 手动同步提交
props.put("enable.auto.commit", "false");
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
process(records); // 处理业务逻辑
consumer.commitSync(); // 处理完再提交
}
这样保证了"业务处理成功,Offset才提交",但性能比自动提交差一些(每次poll后阻塞提交)。
③ 手动异步提交
consumer.commitAsync((offsets, exception) -> {
if (exception != null) {
// 提交失败,记录日志
}
});
不阻塞主线程,性能好,但可能丢Offset(提交失败没重试)。
踩坑点:很多新手处理消息时抛了异常,导致Offset没提交,重启后一直重复消费同一条消息,进入死循环。建议捕获异常,把失败消息写入死信队列(DLQ),然后正常提交Offset,别让一条坏消息卡死整个消费者。
7.3 指定消费位置
消费者启动时,如果找不到已提交的Offset(比如新消费者组第一次消费),可以通过auto.offset.reset决定从哪开始:
earliest:从头开始消费(历史数据都要)latest:从最新的位置开始(只消费启动后的新数据,默认值)none:找不到Offset就抛异常
踩坑点:测试环境经常换消费者组名,每次
auto.offset.reset=latest导致看不到历史数据,查半天以为是生产者没发消息。新消费者组想看全量数据,一定要配earliest。
八、数据不丢失的完整保障方案
"怎么保证Kafka消息不丢?"这是面试高频题。完整答案要从三个维度说:
8.1 生产者端不丢
acks=all(等所有副本确认) +min.insync.replicas >= 2(至少2个副本写成功才算成功)- 开启幂等性(
enable.idempotence=true,自动去重) - 配置重试:
retries=Integer.MAX_VALUE(无限重试),delivery.timeout.ms=120000(2分钟超时) - 生产者缓冲区满了时的行为:
max.block.ms(超时抛异常,避免静默丢数据)
8.2 Broker端不丢
replication.factor >= 3(至少3副本)min.insync.replicas=2(至少2个副本写成功,生产者才收到成功响应)unclean.leader.election.enable=false(ISR为空时宁可暂停服务,也不选不同步的副本当Leader)- 磁盘做RAID或使用云盘多副本,防止单盘损坏
8.3 消费者端不丢
- 关闭自动提交,手动提交Offset
- 业务处理成功后,再提交Offset
- 处理失败的消息进死信队列,不要卡住主流程
8.4 精准一次性消费(Exactly-Once)怎么实现?
前面说的At-Least-Once是"至少消费一次,可能重复"。但有些场景真的不能重复——比如转账扣款,重复扣一次用户就要骂娘了。实现Exactly-Once需要生产端 + 消费端双管齐下:
① 生产端:幂等性 + 事务
// 开启幂等性(Kafka 3.x默认开启)
props.put("enable.idempotence", "true");
props.put("acks", "all");
props.put("retries", Integer.MAX_VALUE);
// 如果需要跨分区事务(比如多条消息要么都成功要么都失败)
props.put("transactional.id", "my-transactional-id"); // 必须唯一!
producer.initTransactions();
try {
producer.beginTransaction();
producer.send(new ProducerRecord<>("topicA", "order_001", "扣款100元"));
producer.send(new ProducerRecord<>("topicB", "log_001", "扣款日志"));
producer.commitTransaction(); // 一起提交
} catch (Exception e) {
producer.abortTransaction(); // 一起回滚
}
幂等性解决的是单分区内的重复发送问题(网络超时导致生产者重试)。原理是Broker给每个生产者分配PID,每条消息带递增的Sequence Number,Broker自动丢弃重复。
事务解决的是跨分区的原子性——要么都写入成功,要么都不写。
注意:
transactional.id每个生产者实例必须唯一,重启后不能变。如果两个实例用同一个ID,Kafka会踢掉先那个,导致"Producer Fenced"异常。
② 消费端:业务层幂等去重
生产者端只管住了"发送到Kafka不重复",但消费者从Kafka拉出来处理后,如果业务逻辑执行了但Offset没提交(或者提交了但处理结果没生效),还是会出问题。
消费端Exactly-Once的标准做法是:把业务处理和Offset提交绑定成一个原子操作。最常见的方案:
// 方案A:数据库事务(业务写入 + Offset存入同一张表)
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
// 开启数据库事务
dbConnection.setAutoCommit(false);
try {
for (ConsumerRecord<String, String> record : records) {
// 1. 业务处理:比如写入订单表(order_id是唯一键)
insertOrder(record.value());
// 2. 同时记录消费位点
saveOffset(record.topic(), record.partition(), record.offset());
}
dbConnection.commit(); // 业务和Offset一起提交
} catch (Exception e) {
dbConnection.rollback(); // 一起回滚
}
}
// 方案B:Redis原子操作(适合简单去重场景)
for (ConsumerRecord<String, String> record : records) {
String orderId = extractOrderId(record.value());
// SETNX = SET if Not eXists,只有不存在才设置成功
boolean success = redis.setnx("kafka:consumed:" + orderId, "1", 24, TimeUnit.HOURS);
if (success) {
process(record); // 第一次消费,处理
} else {
// 重复消息,跳过
}
}
consumer.commitSync();
③ Exactly-Once vs At-Least-Once 怎么选?
| 场景 | 推荐方案 | 原因 |
|---|---|---|
| 日志采集、监控数据 | At-Least-Once | 重复几条不影响 |
| 订单、支付、转账 | Exactly-Once | 重复就是事故 |
| 普通业务消息 | At-Least-Once + 业务幂等 | 性价比最高 |
实践经验:大部分公司不会用Kafka的事务API(太重了),而是选At-Least-Once + 数据库唯一键/Redis去重。比如订单表天然有order_id唯一索引,重复消费时插入会报主键冲突,捕获异常继续就行。这样简单又可靠。
九、常用命令速查(新版命令,基于–bootstrap-server)
老教程里一堆--zookeeper的命令已经过时了,Kafka 3.x/4.x的推荐命令:
# ========== Topic操作 ==========
# 创建Topic(3分区,2副本)
kafka-topics.sh --bootstrap-server node1:9092 --create --topic my-topic \
--partitions 3 --replication-factor 2
# 查看Topic列表
kafka-topics.sh --bootstrap-server node1:9092 --list
# 查看Topic详情
kafka-topics.sh --bootstrap-server node1:9092 --describe --topic my-topic
# 增加分区(只能增不能减!)
kafka-topics.sh --bootstrap-server node1:9092 --alter --topic my-topic --partitions 6
# 删除Topic
kafka-topics.sh --bootstrap-server node1:9092 --delete --topic my-topic
# ========== 生产/消费消息 ==========
# 发送消息
kafka-console-producer.sh --bootstrap-server node1:9092 --topic my-topic
# 消费消息(从最新开始)
kafka-console-consumer.sh --bootstrap-server node1:9092 --topic my-topic
# 消费消息(从头开始)
kafka-console-consumer.sh --bootstrap-server node1:9092 --topic my-topic --from-beginning
# 查看消费者组列表
kafka-consumer-groups.sh --bootstrap-server node1:9092 --list
# 查看消费者组详情(消费进度、Lag)
kafka-consumer-groups.sh --bootstrap-server node1:9092 --describe --group my-group
踩坑点:
--zookeeper参数在Kafka 4.0里已经彻底移除了,如果你看到还这么写的教程,直接划走,那是古董。另外分区只能增加不能减少,建Topic时务必规划好分区数。
十、生产环境 checklist
最后整理一份上线前必查清单:
| 检查项 | 推荐配置 |
|---|---|
| 副本数 | replication.factor=3(最低) |
| 最小同步副本 | min.insync.replicas=2 |
| ACK级别 | acks=all(核心业务) |
| 分区数 | 预估吞吐量,宁可多勿少(只能增不能减) |
| 消息保留 | retention.ms=7天 或 retention.bytes(按磁盘容量) |
| 消费者Offset提交 | 手动提交 + 失败进死信队列 |
| JVM堆内存 | 610G(一般场景),1216G(大规模集群),其余给OS PageCache |
| 磁盘 | SSD推荐,机械盘也行但要独立磁盘(避免和OS/日志抢IO) |
| 监控 | Consumer Lag、Broker磁盘水位、ISR列表变化 |
总结
这篇梳理了Kafka从架构到原理再到实战的核心知识点。几个最关键的概念再强调一下:
- KRaft替代Zookeeper是Kafka 4.0最大的变化,新集群直接上KRaft
- 分区是并行处理的基础,分区数一开始就规划好
- ISR + HW机制保证了数据一致性和不丢失
- 顺序写 + PageCache + 零拷贝是Kafka高性能的三大基石
- acks=all + min.insync.replicas是生产环境数据安全的底线配置
如果对你有帮助,欢迎点赞收藏。有问题评论区见~
更多推荐




所有评论(0)