目录
关于其他消息队列的文章: 带你轻松学习RocketMQ-CSDN博客https://blog.csdn.net/2401_88959292/article/details/149296505?spm=1001.2014.3001.5501 带你轻松学习RabbitMQ-CSDN博客https://blog.csdn.net/2401_88959292/article/details/149294023?spm=1001.2014.3001.5501
一、核心组件/概念
1.主题(Topic)
是消息的逻辑分组,也是存储消息元数据的容器。
2.分区(Partition)
是Topic的并行分片,每个分区都是有序的队列,不同分区的消息完全并行处理。
3.分区键(Key)
是将消息写入指定分区的标识键。
4.副本(Replica)
分为Leader和Follower,Leader可读可写,但Follower仅同步无读写权限。
5.消息代理(Broker)
负责消息存储以及消费者的请求处理,同时也负责的协调工作。
二、消息流程


流程说明:
1.生产者发送消息
生产者连接Broker集群,构建消息并发送。
消息在发送前需要指定以下三个属性:
- Topic
- Key(可选)
- 消息实体
生产这会先将消息存到本地缓冲区,进行批量发送,以减少网络IO次数。
Broker接收到消息后根据分区将消息发送到对应的分区。
分局策略如下:
- 哈希策略(默认):若指定Partition Key,用Key哈希值%分区数计算分区ID(保证相同Key的消息进入同一分区,保证分区内顺序)。
- 轮询策略:未指定Partition Key时,按顺序将消息分配到各分区(负载均衡)。
- 自定义策略:用户可实现
Partitioner接口,定义自定义分区逻辑(如按地域分配)。
2.Leader处理消息
分区的Leader接收到消息后,会将其以追加模式写入磁盘的文件中(追加模式保证顺序写)进行持久化操作。
随后Leader会将消息写入ISR(与Leader保持同步的Follower集合)中的Follower,Follower收到消息后同样会写入Log进行持久化,并向Leader返回确认。
ISR会将长期未同步或者同步进度缓慢的Follower移出集合,待恢复后再重新加入。 当Leader宕机后会从选举拥有最新offset的Follower作为新的Leader。 副本同步机制是Kafka高可用的核心。
3.Leader返回ACK
Leader根据不同的策略会在对应的时机返回ACK给生产者:
- acks=0:Leader不等待任何确认,直接返回成功(吞吐量最高,但可能丢失消息,如Leader未刷盘就故障)。
- acks=1:Leader确认消息写入本地日志后返回成功(默认配置,可靠性中等,若Leader故障但Follower未同步,可能丢失消息)。
- acks=all(或-1):Leader确认消息写入本地日志,且大多数Follower(ISR中的多数)同步完成后返回成功(可靠性最高,适合金融等核心业务)。
4.消费者拉取消息
消费者组先订阅Topic,然后Broker会根据分区分配策略将Topic的分区分配给组内的消费者:
- Range策略(默认):按分区ID顺序分配,每个消费者分配连续的分区(如3个分区,2个消费者:消费者1分配分区0、1,消费者2分配分区2)。
- Round-Robin策略:按消费者顺序轮询分配分区(如4个分区,2个消费者:消费者1分配分区0、2,消费者2分配分区1、3)。
- Sticky策略(粘性策略):尽量保持现有分配,仅在必要时调整(减少分区迁移,提升稳定性)。
但每个分区只会同时被一个消费者消费,以此确保消息不会被重复消费。
消费者主动向分区的Leader拉取数据,拉取时消费者需要指定分区Id以及offset。
此环节使用了零拷贝技术,详细可见以下博客中的4.消费者拉取消息: 带你轻松学习RocketMQ-CSDN博客
5.消费者返回offset
消费者处理完业务逻辑后,会根据策略选择不同的时机提交offset给Broker:
- 自动提交:消费者会定期(默认5000ms)自动提交当前offset(可能会丢失消息)。
- 手动提交:
- 消费者处理成功后手动提交:
- 同步提交:等待Broker确认后返回(可靠性高,但阻塞线程)。
- 异步提交:不等待确认,通过回调函数处理结果(吞吐量高,但可能丢失消息)。
6.Broker维护offset
Broker直接将消费者提交的offset存储到内部的__consumer_offsets主题,确保持久化和高可用性。
三、优劣分析
(一)优点
1.极致的高吞吐
Kafka的高是其最显著的优势,主要依赖以下机制:
- 顺序写优化:每个分区的日志文件采用Append-Only(追加)模式,磁盘IO为顺序写。
- 零拷贝技术:通过
sendfile(网络传输)和mmap(磁盘存储)减少数据在“内核空间”与“用户空间”的拷贝(详见之前博客),降低CPU占用,提升IO效率。 - 分区并行:Topic的多个分区分布在不同Broker节点,生产者可并行发送消息,消费者组可并行消费(每个分区由组内一个消费者处理),实现负载均衡。
- 批量发送/拉取:生产者将消息缓存到本地缓冲区,批量发送;消费者批量拉取消息,减少网络请求次数。
2.的高扩展性
Kafka的集群设计采用无状态Broker+集中式元数据管理(ZooKeeper,但新版本使用KRaft更佳),支持水平扩展:
- Broker扩展:新增Broker节点后,可通过
kafka-reassign-partitions工具将现有Topic的分区迁移到新Broker,提升集群吞吐量。 - 分区扩展:Topic的分区数可动态增加,每个分区独立存储和处理消息,支持业务流量的线性增长。
- 消费者组扩展:新增消费者到消费者组后,Broker会重新分配分区,保证消费负载均衡。
3.高可用性机制
Kafka通过副本集和Leader选举机制保证服务的高可用性:
- 副本同步:每个分区有多个副本(默认1个Leader+2个Follower),Leader处理读写请求,Follower同步Leader的日志数据。ISR动态维护,确保副本的一致性。
- Leader选举:若Leader节点宕机,Kafka会从ISR中选举新的Leader,整个过程秒级完成,不影响服务。
- 数据可靠性:生产者可配置
acks=all(或-1),要求Leader确认消息写入本地日志且多数Follower同步完成后返回成功,保证数据不丢失(适合金融等核心业务)。
4. 消息持久化与可回溯
Kafka的消息存储采用磁盘日志文件,支持长期保留和可回溯消费:
- 持久化策略:消息默认保留7天,超过保留时间或日志大小后,旧消息会被删除。
- 可回溯消费:消费者可通过修改Offset,从过去的某个位置重新消费消息。
- Offset管理:消费者的Offset存储在Kafka内部的
__consumer_offsets主题中,即使消费者宕机,也能从上次提交的Offset继续消费。
(二)缺点
1. 全局顺序消息的瓶颈
Kafka只能保证分区内的消息顺序(因为每个分区是顺序写的),若需要全局顺序,必须将所有消息发送到同一个分区。此时,单分区的处理能力会成为性能瓶颈,无法应对高并发场景。
2. 事务支持较弱
Kafka的事务是生产者端的事务,支持原子性发送多个消息到多个Topic/分区,但相对于的事务消息,其事务处理更复杂:
- 配置复杂:生产者需要开启事务,消费者需要设置
isolation.level。 - 性能开销:事务消息需要协调多个Broker(如提交事务时需要同步所有参与的分区),吞吐量会下降。
- 功能有限:不支持“半消息+事务回查”的机制,需要生产者自行处理事务失败的情况。
所以实际上并不会使用Kafka进行事务性强的业务相关的操作。
3.无法应对场景
Kafka的默认配置是高吞吐量优先,对于延迟敏感的场景,需要调整参数来降低延迟,但会牺牲部分吞吐量:
- 批量发送的延迟:生产者的
linger.ms(默认0)配置表示等待多久来积累批量,若设置为10,会增加10ms的延迟,但提升批量效果(吞吐量增加);若设置为0,延迟降低,但批量效果差(吞吐量下降)。 - 拉取的延迟:消费者的
fetch.min.bytes(默认1B)配置表示每次拉取的最小字节数,若设置为1024,会等待足够的消息积累后再拉取,增加延迟,但减少网络请求次数(吞吐量增加);若设置为1,延迟降低,但网络请求次数增加(吞吐量下降)。
可以看出高吞吐和低延迟其实是难以兼容的,这也是Kafka优先高吞吐权衡之下的选择。
4.存在漏消费或者重复消费问题
Kafka为了确保消息可靠性一般不会启用配置acks=all(或-1),因为这样反而会降低吞吐量,违背了使用Kafka作为MQ的初衷,所以一般会采用以下两种模式:
- 至少一次:消费者先拉取消息,再提交offset。这是最为常见的方案,以容忍重复消费换取消息一定被消费的保障。
- 最多一次:消费者先提交offset,再拉取消息。这是以容忍漏消费换取高吞吐量。
5. 对小消息的处理效率较低
Kafka的批量发送和顺序写适合大消息或批量消息,但对于小消息,其处理效率较低:
- 批量效果差:小消息的批量大小(
batch.size)需要设置较小,否则等待时间会增加延迟,导致批量效果不明显。 - 元数据开销大:每个消息的元数据约占10-20B,对于小消息,元数据的比例较高,增加了存储和网络开销。
- 网络传输次数多:小消息的批量大小小,导致网络请求次数增加,降低了吞吐量。
四、最佳实践
** 一个Topic负责一个业务场景,一个消费者组仅订阅一个Topic**
因为使用Kafka的场景大多不在意消息的高可靠性而是追求其高并发,所以尽量不使用Key以此节省筛选时间。
注意!必须避免以下情况:
- 一个消费者组订阅多个Topic,这会导致业务耦合且负载均衡失效,增加维护成本。
- 消费者组中的消费者数量大于Topic的分区数量,这样会导致多余的消费者空闲,造成资源浪费。
五、总结
Kafka作为分布式流处理平台,其核心特性决定了其适合大规模数据流式传输、实时分析、日志收集等场景,同时在低延迟、强事务、小规模应用等场景下存在局限性。
一、适合的场景
- **大规模数据流式传输场景:**如电商用户行为数据采集、物联网传感器数据传输、社交媒体信息流。 Kafka的高吞吐量和分区并行特性,能支撑TB级/天的大规模数据传输。
- **实时流处理场景:**实时推荐系统、实时监控报警。Kafka与实时计算框架深度整合,支持流-流 joins、窗口计算等实时操作。
- 日志与事件收集场景:服务器日志收集、应用事件跟踪、运维监控数据。Kafka的持久化存储和可回溯消费特性,适合作为“日志中心”。
- 数据备份与同步场景:数据库变更同步、跨地域数据复制。Kafka的Kafka Connect工具支持与关系型数据库、NoSQL数据库的双向同步,实现数据的增量备份(仅同步变更数据)。
二、不适合的场景
- 极度低延迟场景:高频交易、实时视频流、工业控制。Kafka的端到端延迟为毫秒级,无法满足微秒级延迟需求。
- 小规模低并发场景:创业公司内部通知系统、小流量应用。 Kafka的架构复杂,资源占用高**。**
- 强事务要求的场景:金融转账交易、电商订单支付。Kafka的事务支持仅能保证原子性发送多个消息到多个Topic/分区,但不支持半消息+事务回查。
- ** 需要严格全局顺序的场景**:订单的严格流程、物流状态的更新。 Kafka仅能保证分区内的消息顺序(通过Key哈希到同一分区),无法保证全局顺序。若需要全局顺序,必须将所有消息发送到单分区(此时吞吐量会下降到单分区的处理能力)。
- 简单点对点通信场景: 如内部通知系统、简单消息传递等。Kafka的模型比RabbitMQ的
Direct交换机更复杂,简单场景下用RabbitMQ更便捷。
码文不易,留个赞再走吧
原文链接: 带你轻松学习Kafka 作者: Yilena
如果这篇文章对你有帮助,欢迎分享给更多人!
部分信息可能已经过时










