一、Kafka是什么
一句话:Kafka是一个分布式的、基于日志(log)的消息系统。它的核心思想很简单——把消息像写日志一样"追加"到磁盘文件里,消费者按顺序"读"这个文件,仅此而已。这个简单模型是它高吞吐的根本原因。
二、核心架构组件
在讲数据流之前,先认识几个名词:
- Producer(生产者):发消息的一方,比如你的业务代码
- Broker:Kafka的服务器节点,多个Broker组成一个Cluster(集群)
- Topic(主题):消息的分类,比如"订单创建"topic
- Partition(分区):每个Topic会被切成多个分区,分区是Kafka并行处理和存储的最小单位
- Replica(副本):每个分区有多个副本分布在不同Broker上,其中一个是Leader(负责读写),其余是Follower(只同步数据,不对外服务)
- Consumer(消费者) / Consumer Group(消费组):消费消息的一方,同一个组内的消费者共同"瓜分"一个Topic的所有分区
Controller:集群里选出的一个特殊Broker,负责管理分区Leader选举等元数据工作(以前靠ZooKeeper协调,新版本Kafka用内置的KRaft协议取代了ZooKeeper)## 三、数据从Producer到Broker再到Consumer的完整过程
Producer发送消息:业务代码调用
send(),消息先进入客户端本地的一个内存缓冲区(RecordAccumulator),并不是发一条就立刻走网络- 选择分区:如果你指定了key,Kafka默认对key做hash来决定发到哪个分区(保证同key的消息永远进同一个分区,从而保证顺序);没指定key则轮询或用粘性分区策略
- 批量发送:客户端把同一分区的多条消息打包成一个batch,达到大小或时间阈值后一起发给对应分区的Leader Broker(这是Kafka高吞吐的关键之一——减少网络往返)
- Leader写入日志:Leader收到消息后,以顺序追加的方式写入本地磁盘的log文件,并分配一个递增的offset
- 副本同步:Follower副本主动从Leader拉取新消息进行同步,当所有处于ISR(In-Sync Replicas,同步副本集合)中的副本都写完后,才算这条消息"提交"(commit)
- 返回ack:根据Producer的
acks配置(0/1/all)决定什么时候返回成功,这一步直接决定了可靠性和性能的取舍 - Consumer拉取:消费者是主动poll(拉模式,不是broker推给你),按上次记录的offset位置继续往后读
- 提交offset:消费完成后,消费者把"我读到哪了"这个offset提交到一个特殊的内部Topic
__consumer_offsets里,下次重启就从这个位置接着消费
四、日志和数据结构细节
Kafka的性能秘密基本都藏在这套存储结构里:
- Partition = 一个只追加的文件序列:每个分区在磁盘上对应一个目录,消息按offset严格递增地写入
- Log Segment(日志段):一个分区的数据不会塞进一个文件,而是切成多个segment,每个segment有三类文件:
.log:真正存消息内容的文件.index:offset到物理文件位置的稀疏索引,用于快速定位.timeindex:时间戳到offset的索引,支持"按时间查消息"
- 消息本身的结构包含:key、value、时间戳、offset、headers、CRC校验、压缩标记等
- 顺序写 + OS Page Cache + 零拷贝(sendfile):这三者结合是Kafka比传统消息队列快得多的根本原因——磁盘顺序写接近内存速度,读取时利用系统缓存和零拷贝减少用户态/内核态数据拷贝
数据清理两种策略:按时间/大小滚动删除旧segment(
retention),或者Log Compaction(只保留每个key最新的一条,常用于存"最新状态"场景,比如Kafka自己用它存offset)## 五、实际使用中要注意的坑分区数不是越多越好:分区数决定了并行度上限,但太多分区会增加Controller管理开销、增加端到端延迟、增加文件句柄消耗。一般根据目标吞吐量和消费者数量来定,而不是越多越好
- 分区键设计要避免热点:如果大量消息用同一个key(比如都用"CN"作为地区key),会导致某个分区数据量远超其他分区,造成热点和消费延迟不均
- acks参数是可靠性与性能的权衡:
acks=0最快但可能丢数据;acks=1只要leader写成功就返回,leader挂了还没同步给follower的数据会丢;acks=all配合min.insync.replicas才能真正保证不丢消息,但延迟更高 - 顺序性只在分区内保证,跨分区不保证顺序。如果业务要求严格顺序,必须把相关消息发到同一个分区
- 消费者组的Rebalance是常见坑:新消费者加入、消费者挂掉、
session.timeout.ms太短都会触发rebalance,期间会短暂停止消费。处理消息耗时过长导致max.poll.interval.ms超时是最常见的诱因,要么加快处理速度,要么调大这个参数,要么用异步处理+手动提交offset - 重复消费是常态,不是异常:Kafka是"至少一次"(at-least-once)语义为主,消费者重启、rebalance都可能导致重复消费,业务逻辑必须做幂等处理(比如用消息里的唯一ID做去重)
- 要监控消费延迟(Lag):consumer的offset落后producer offset太多说明消费能力跟不上,这是生产环境最需要盯的指标之一
- 磁盘容量规划:Kafka数据是先落盘再清理的,retention设置太长+topic量大很容易把磁盘写爆
- 序列化格式建议提前定好:生产环境不建议直接传JSON字符串,用Avro/Protobuf配合Schema Registry能省很多后期兼容性的麻烦
六、延时队列怎么实现
关键点:Kafka本身并不原生支持"延时消息",它只是一个按offset顺序读写的日志,没有"消息到期才可见"这种机制。要在Kafka上做延时队列,常见思路是:
- 多级延时Topic + 轮询转发:仿照RocketMQ思路,建几个不同粒度的延时topic(比如
delay-5s、delay-30s、delay-1m),生产者把消息连同"目标可见时间"一起写入对应topic,由一个专门的转发服务不断轮询这些topic,发现到期的消息就转发到真正的业务topic;没到期的可以选择sleep等待或者重新写回原topic(类似时间轮的思路) - 借助外部存储做调度器:把消息连同到期时间存入Redis的ZSET(score=到期时间戳)或数据库表,用一个独立的定时任务/时间轮扫描到期的记录,到期后再发到Kafka真正的topic。这是目前生产环境用得最多的方式,因为Kafka本身不适合频繁"回查"
- 本质上都是"外挂"方案:因为Kafka缺少消息级别的可见性控制(不像队列型系统可以"暂扣"某条消息),延时能力必须靠上层组件补齐,精度和维护成本都比原生支持延时的MQ要高
七、和RocketMQ、RabbitMQ的区别
| 维度 | Kafka | RocketMQ | RabbitMQ |
|---|---|---|---|
| 核心模型 | 分布式日志,partition顺序追加 | commitlog + 多队列,借鉴Kafka又做了增强 | AMQP协议,交换机(Exchange)+队列路由 |
| 消费模式 | 纯拉模式(pull) | 推拉结合(长轮询) | 推模式为主 |
| 延时消息 | 不原生支持,需外部方案 | 原生支持,早期版本固定18个延时级别(1s~2h),新版本支持任意精度延时 | 通过TTL+死信队列实现,或装delayed-message-exchange插件支持任意延时 |
| 事务消息 | 支持(较复杂,主要面向流处理场景) | 原生支持,半消息+回查机制,是其强项之一 | 支持,但性能开销较大 |
| 顺序消息 | 分区内有序 | 支持全局顺序或分区顺序 | 单队列内基本有序,但集群/多消费者场景较弱 |
| 吞吐量 | 最高,适合海量日志/流数据 | 高,略低于Kafka | 相对最低,适合中小规模、低延迟场景 |
| 路由能力 | 简单(topic+partition) | 中等(tag过滤) | 最强,支持direct/topic/fanout/headers多种灵活路由 |
| 典型场景 | 日志采集、流处理、大数据管道 | 电商订单、金融交易(强调可靠性和延时/事务消息) | 微服务间的低延迟RPC、任务分发 |
简单总结一下选型直觉:要吞吐量和流处理选Kafka;要原生延时消息、事务消息、订单类强可靠性场景选RocketMQ;要灵活路由、复杂场景解耦、开发效率优先选RabbitMQ。