Kafka
简介
Kafka 是一个分布式流处理平台,可以用来高效地收集、传输、存储、处理实时数据流。它最初由 LinkedIn 开发,用于解决海量日志处理的问题,后来开源给 Apache,成为现在广泛应用的数据中枢系统。
Kafka 能做消息队列、日志采集、事件驱动、实时数据流处理、数据管道。
主要应用场景:系统解耦 + 异步通信、大数据日志收集、实时数据分析、监控与告警系统、数据库变更同步(CDC)
重要的原因:高吞吐、可持久化、分布式架构、支持实时与离线、容错性强、生态成熟、可回溯消息。
Kafka 是现代数据架构的 核心中间件,让系统能更灵活、可扩展、实时响应世界的变化。
架构
Kafka 是一个分布式的发布订阅消息系统,核心由 Producer、Broker、Consumer、Topic、Partition、Consumer Group 组成,通过分区 + 副本 + 顺序写日志实现高吞吐与高可用。
Producer 负责将消息发送到 Kafka。
- 按 Topic + Partition 发送
- 可以指定 key,决定消息落到哪个 partition(保证局部有序)
- 支持同步 / 异步发送
- 支持 ack 机制(0/1/all)
Producer 不直接写 Broker 集群逻辑,只负责发送
Kafka 集群由多个 Broker 组成,每个 Broker 是一个独立 Kafka 服务实例。
- 每个 Broker 存储部分 Topic 的 Partition
- 负责读写请求处理
- 负责副本同步(Replica Sync)
Kafka 是去中心化 Broker 集群,没有单点 master(元数据除外)
Topic 是消息分类逻辑概念。
- 一个 Topic 可以有多个 Partition
- Topic 本身不存数据,数据存在 Partition 中
Partition 是 Kafka 高吞吐的关键。
- 每个 Partition 是一个追加写的有序日志文件
- Partition 内消息严格有序
- Topic 通过分区实现水平扩展
- 不同 Partition 可分布在不同 Broker 上
Kafka 的“高吞吐 + 扩展性”本质来自 Partition 并行读写
每个 Partition 有多个副本:
- Leader Replica(负责读写)
- Follower Replica(负责同步)
- ISR(In-Sync Replica)集合:同步中的副本
写入流程:
- Producer 写入 Leader
- Leader 同步给 Follower
- ISR 全部确认后(或部分)返回 ack
Kafka 通过 ISR 机制实现高可用与数据可靠性
Consumer(消费者)
- 从 Broker 拉取消息(pull 模型)
- 自己维护 offset(消费进度)
- 可以并行消费多个 partition
Consumer Group(消费者组)Kafka 并行消费核心机制:
- 一个 Partition 在同一个 Consumer Group 内,同一时刻只能被一个 Consumer 消费
- 不同 Consumer Group 内的 Consumer 可以并行消费同一个 Partition 中的数据。
- Consumer Group 内的 Consumer 数量 > Partition 会有空闲
Kafka 的扩展消费能力 = Partition 数量
写入流程:Producer → Broker Leader → Follower 同步 → ISR ack → 返回成功
消费流程:Consumer Group → 分配 Partition → pull 拉取数据 → commit offset
Kafka 高性能的核心原因:
- 顺序写磁盘(append log)
- 零拷贝(Zero Copy)
- 批量发送(batch)
- 分区并行读写
- PageCache 利用内存缓存
- 拉取模型(pull)减少服务端压力
Kafka 是一个基于 Partition + Replica + Consumer Group 的分布式消息系统,通过顺序日志写入 + 分区并行 + ISR副本机制实现高吞吐、高可用和可扩展。
数据倾斜问题
Kafka 本身并不是“专门解决数据倾斜的系统”,但它通过 Partition 机制 + 分区策略 + 扩展手段 来缓解或避免数据倾斜问题。
在 Kafka 中,数据倾斜本质是:某些 Partition 数据量或流量远大于其他 Partition,导致负载不均衡。
Kafka 的解决思路主要有以下几种:
通过 Partition 实现水平拆分(基础手段)
Kafka 将 Topic 拆分为多个 Partition,通过多 Partition 分摊数据写入和消费压力。
- Producer 写入时根据 key 进行 hash(默认策略)
- 不同 key 分布到不同 Partition
- Consumer Group 按 Partition 并行消费
本质上是通过“分片”减少单点压力
通过“自定义分区策略”避免热点 key(关键点)
- 如果使用默认 hash,可能导致某些 key 访问量极高,导致某个 Partition 成为热点(倾斜)
- Kafka 支持自定义 Partitioner 按业务维度拆分 key、加盐(salt)打散热点 key、按用户ID / 订单ID等均匀分布
- 典型手段热点 key + 随机前缀 / hash 二次打散
提升 Partition 数量(扩展能力)
- 当出现负载不均时,可以:增加 Partition 数量、提高并行度、分散写入压力
- Partition 增加后旧数据不会自动重分布,只能缓解新数据倾斜
Consumer 侧扩展消费能力(间接缓解)
- 通过增加 Consumer 数量,提升消费并行度,避免单 Consumer 处理过多 Partition
- 但本质仍受 Partition 数量限制
业务层“拆 Topic”解决极端倾斜(高级手段)
- 在极端情况下,单 Topic 已无法均衡,可以按业务维度拆分 Topic
Kafka 通过 Partition 分片机制实现天然的水平扩展能力,但对于数据倾斜问题,主要依赖合理的分区策略(如自定义 Partitioner、热点 key 打散)以及增加 Partition 和 Consumer 来提升整体并行度,从而缓解负载不均问题。Kafka 的数据倾斜本质不是系统自动解决的,而是需要业务在 key 设计和分区策略上进行控制。
在 k8s 中增加消费者实例
每个 Pod(在 K8s 中运行的消费者实例)都会作为一个独立的消费者加入到 Kafka 消费者组中。
如果只是将消费者代码在 K8s 中复制多个实例(例如通过增加副本数),每个实例都会作为一个独立的消费者加入到 Kafka 消费者组中,Kafka 会根据消费者组的机制将消息分配给各个消费者实例。这样可以通过增加实例来增加消费者的并发处理能力。
并不需要修改代码来增加消费者组中的消费者实例个数。只要 Kubernetes 中的 Pod 数量增加,Kafka 消费者组的成员数就会增加,Kafka 会自动重新分配分区给新的消费者实例。
消费者组的大小不能超过 Kafka 主题的分区数。如果消费者实例多于分区数,那么有些消费者将没有消息可消费。
保证消息可靠性
Kafka 通过 多副本机制、ACK 确认机制、ISR 机制、持久化存储以及消费端 offset 机制 来保证消息的可靠性。
多副本机制(Replication)
- Kafka 中每个 Partition 都有 1 个 Leader,多个 Follower,写入流程时,Producer 只写 Leader,Leader 负责同步给 Follower,Follower 从 Leader 拉取数据进行同步,即使某个 Broker 宕机,也可以从副本恢复数据
ISR 机制(In-Sync Replicas)
- Kafka 维护一个 ISR 集合,表示与 Leader 保持同步的副本,只有 ISR 中的副本才被认为是“可靠副本”,当写入数据时,Leader 收到消息,等 ISR 中副本同步完成(或满足条件),才认为写入成功,保证“不会丢已确认的数据”
ACK 机制(Producer 可靠性核心)
Producer 可以配置 ack:
- acks = 0:不等待确认(可能丢数据)
- acks = 1:Leader 写入成功即返回
- acks = all(-1):ISR 全部确认才返回(最可靠)
生产环境一般使用:acks = all + min.insync.replicas = 2
持久化机制(Log Append)
- Kafka 使用顺序写磁盘日志(append-only log):数据写入磁盘文件,利用 OS PageCache 提高性能,数据不会被立即覆盖,即使宕机,也可以通过日志恢复数据
消费端 Offset 机制(避免重复/丢失)
- Kafka 通过 offset 来记录消费进度:Consumer 自己维护 offset,可以提交到 Kafka(__consumer_offsets),支持手动提交 / 自动提交,消费语义:at-most-once(可能丢)、at-least-once(可能重复,最常用)、exactly-once(Kafka + 事务支持)。
事务机制(Exactly Once)
- Kafka 支持事务:Producer 发送多条消息要么全部成功,要么全部失败,避免“半成功数据”,配合 idempotent producer(幂等生产者),transactional API
Kafka 通过多副本机制保证高可用,通过 ISR 和 ACK 机制保证写入可靠性,通过持久化日志保证数据不丢失,并通过消费端 offset 和事务机制实现消费可靠性与语义一致性。Kafka 的可靠性本质是通过“副本复制 + 写入确认 + 可恢复日志 + 消费位点管理”四层机制共同实现的。
Exactly Once
Kafka 的 Exactly Once 语义(EOS)主要通过 幂等生产者 + 事务机制 + 消费偏移量原子提交 + 下游协调机制 来实现,保证消息在生产、存储和消费过程中不会重复也不会丢失。
幂等生产者(Idempotent Producer)
Kafka 通过为每个 Producer 分配一个 PID(Producer ID),并为每条消息增加 序列号(Sequence Number) 来实现幂等性。
- Broker 会记录每个 PID + Partition 的最大序列号
- 如果重复发送同一条消息(比如重试),Broker 会自动去重
- 从而避免“生产端重复写入”
事务机制(Transaction)
Kafka 引入事务机制来保证“多条消息的原子性写入”。Producer 可以将多条消息放入一个事务中:
- beginTransaction()
- send messages
- commitTransaction() / abortTransaction()
Kafka 会保证:要么全部写入成功,要么全部不写入,解决多条消息一致性 问题
事务 + Offset 绑定(核心关键)
Kafka EOS 的关键点在于:将“消息写入”和“消费 offset 提交”绑定在同一个事务中
流程如下:
- Consumer 拉取消息
- 处理业务逻辑
- Producer 在事务中写结果消息
- 同时提交 consumer offset(offset 作为事务的一部分)
- commitTransaction()
这样可以保证:
- 消息处理成功 + offset 提交 = 原子操作
- 避免“重复消费 / 消费丢失”
Read Process Write 模型(典型 EOS 模式)
Kafka EOS 常用于:consume → process → produce,例如:从 topic A 消费,处理后写入 topic B,同时提交 offset,Kafka 通过事务保证 A 的消费位置,B 的写入结果是同一个事务的一部分
Broker 端支持(事务协调器) Kafka 内部有:Transaction Coordinator(事务协调器)、__transaction_state 主题,负责:事务状态管理、两阶段提交(类似 2PC)、保证事务一致性。
Kafka 的 Exactly Once 语义是通过幂等生产者保证“不重复写入”,通过事务机制保证“多消息原子性”,并将消费 offset 与生产结果绑定在同一事务中,从而实现端到端的精确一次语义。Kafka 的 EOS 本质不是“绝对不重复”,而是在分布式环境下通过“幂等 + 事务 + offset 原子提交”实现逻辑上的 Exactly Once。
offset 提交问题
Kafka 默认使用自动提交 offset,基本上每隔 5 秒会自动提交一次,也可以手动提交,手动提交分为:手动同步提交和手动异步提交。
Kafka 中既可以由 Consumer 提交 offset,也可以由 Producer 提交 offset,两者适用的场景不同。Consumer 提交适用于普通消息消费,不参与事务,而 Produce 提交用于事务。
高水位(HW)和 LEO
Kafka 面试中真正高频的是 High Watermark(HW,高水位)和 LEO(Log End Offset,日志末端偏移量)。
"Low Watermark(LW,低水位)"并不是 Kafka 官方对消息同步机制中的核心概念,很多资料把它和 HW 对称地讲,但在 Kafka 官方实现中几乎不会这样表述。因此,面试时不要重点讲 “LW”,否则容易被追问。
Kafka 为了保证副本数据一致性和消费者不会读取到未同步的数据,引入了 LEO(Log End Offset) 和 HW(High Watermark,高水位) 两个重要概念。
LEO(Log End Offset) 表示一个 Partition 当前日志的末尾位置,也就是下一条消息将要写入的位置。Leader 和每个 Follower 都维护自己的 LEO,因此不同副本的 LEO 可能不同。
HW(High Watermark) 表示已经被 ISR(In-Sync Replicas)中所有副本成功同步的最大 Offset。只有 Offset 小于 HW 的消息,才被认为是可靠的数据。
消费者读取消息时,并不是读取到 Leader 的 LEO,而是最多只能读取到 HW。这是因为 LEO 后面的消息可能还没有同步到所有 ISR 副本,如果此时 Leader 宕机,这些消息可能会丢失。如果消费者提前读取了这些消息,就会出现数据不一致的问题。
Kafka 会根据 ISR 中各个副本的同步进度计算 HW。通常情况下,HW 等于 ISR 中所有副本 LEO 的最小值。这样可以保证 HW 之前的数据已经在所有同步副本中存在,即使 Leader 发生切换,也不会影响消费者已经读取到的数据。
Kafka 通过 LEO 记录副本当前写入位置,通过 HW 记录所有 ISR 副本都已同步的数据位置,消费者只能读取 HW 之前的消息,从而保证即使 Leader 宕机发生切换,也不会读取到可能丢失的数据,确保消息的一致性和可靠性。
Rebalance
Rebalance 是 Kafka Consumer Group 的一种协调机制,当 Consumer Group 中的成员发生变化,或者 Topic 的 Partition 数量发生变化时,Kafka 会重新分配 Partition 与 Consumer 的对应关系,以保证消费负载均衡。
Rebalance 由 Group Coordinator 负责协调,目的是确保每个 Partition 在同一个 Consumer Group 内,同一时刻只会被一个 Consumer 消费,同时尽可能均衡地分配消费任务。
Kafka 在以下几种情况下会触发 Rebalance:
- Consumer 加入 Consumer Group。
- Consumer 正常退出或异常宕机。
- Consumer 长时间没有发送心跳,被 Coordinator 判定为失效。
- Topic 新增或减少 Partition。
- Consumer Group 的订阅关系发生变化。
Rebalance 的执行过程主要包括以下几个步骤:
第一步:停止消费
Group Coordinator 通知 Consumer Group 进入 Rebalance 状态,此时所有 Consumer 会暂停消费。
第二步:重新加入 Group
所有 Consumer 向 Group Coordinator 发送 JoinGroup 请求。Coordinator 会选择其中一个 Consumer 作为 Group Leader。
第三步:生成分配方案
Group Leader 根据分区分配策略(如 Range、RoundRobin、Sticky 等)计算 Partition 的分配结果,并发送给 Coordinator。
第四步:同步分配结果
Coordinator 将分配结果发送给所有 Consumer,各 Consumer 根据新的分配结果开始消费对应的 Partition。
Rebalance 虽然保证了负载均衡,但也会带来一些问题:
- Rebalance 期间所有 Consumer 都会暂停消费,影响系统吞吐量。
- 如果 Consumer 数量频繁变化,会导致频繁 Rebalance,影响系统稳定性。
- Partition 重新分配后,本地缓存可能失效,需要重新建立状态。
因此,在生产环境中,应尽量减少不必要的 Rebalance。常见优化方式包括:
- 合理设置
session.timeout.ms和heartbeat.interval.ms,避免因心跳超时导致误判 Consumer 下线。 - 增大
max.poll.interval.ms,避免业务处理时间过长触发 Rebalance。 - 使用 Cooperative Sticky Assignor,采用渐进式 Rebalance,减少所有 Consumer 同时停止消费的情况。
- 保持 Consumer 实例数量稳定,避免频繁扩缩容。
Rebalance 是 Kafka Consumer Group 的负载均衡机制。当 Consumer 成员或 Partition 数量发生变化时,由 Group Coordinator 协调重新分配 Partition。它保证了一个 Partition 在同一 Consumer Group 内只会被一个 Consumer 消费,但 Rebalance 期间会暂停消费,因此生产环境需要尽量减少频繁 Rebalance,提高系统稳定性。
早期 Kafka 使用的是 Eager Rebalance,一旦触发 Rebalance,所有 Consumer 都需要先释放自己持有的 Partition,再统一重新分配,消费会全部暂停。后来引入 Cooperative Rebalance(增量 Rebalance) 后,只回收需要迁移的 Partition,其余 Partition 可以继续消费,大大减少了停顿时间,提高了 Consumer Group 的稳定性。
** Partition 的分配方案不是由 Group Coordinator 自己计算的。在经典的 Consumer Group 协议中,Coordinator 负责协调流程,而由被选出的 Group Leader Consumer 根据分配策略(Assignor)计算分配方案,然后 Coordinator 将方案同步给其他 Consumer。
零拷贝
零拷贝(Zero Copy)是一种高性能 I/O 技术,它并不是完全没有数据拷贝,而是减少 CPU 参与的数据复制以及用户态和内核态之间的数据传输,从而降低 CPU 开销,提高数据传输效率。
传统的数据发送过程中,数据需要先从磁盘读取到内核缓冲区,再复制到用户空间,之后又复制回内核的 Socket 缓冲区,最后发送到网卡,整个过程需要多次数据拷贝和用户态、内核态切换。
Kafka 使用操作系统提供的 sendfile 等零拷贝技术,数据从磁盘读取到内核后,可以直接发送到 Socket,再由网卡发送出去,不需要经过用户空间,因此减少了 CPU 拷贝和上下文切换。
因此,Kafka 能够以较低的 CPU 消耗实现高吞吐量,这也是 Kafka 高性能的重要原因之一。
Kafka Broker 的主要工作是把磁盘中的消息发送给 Consumer,本身几乎不处理消息内容,因此非常适合使用零拷贝技术,直接将数据从磁盘传输到网络,减少 CPU 消耗,提高消息传输效率。
Kafka vs RocketMQ
Kafka 和 RocketMQ 都是分布式消息队列,但两者的设计目标有所不同。
Kafka 更偏向于高吞吐、大数据和日志流处理场景,而 RocketMQ 更偏向于互联网业务消息场景,更注重消息可靠性和丰富的业务特性。
主要有以下几个区别:
第一,设计目标不同。
Kafka 最初是为日志采集和流式计算设计的,追求极致的吞吐量;RocketMQ 最初由阿里开发,主要服务于电商业务,更注重消息可靠性、事务消息和业务解耦。
第二,性能方面。
Kafka 采用顺序写磁盘、零拷贝、Page Cache 等技术,吞吐量通常更高,适合大规模日志和实时数据处理;RocketMQ 吞吐量虽然略低,但仍能满足绝大多数业务场景。
第三,消息可靠性。
Kafka 主要通过副本机制、ISR、ACK 和事务保证消息可靠性;RocketMQ 除了提供高可靠存储外,还原生支持事务消息、延迟消息、定时消息等业务能力,在业务消息场景更有优势。
第四,消息顺序。
Kafka 可以保证同一个 Partition 内消息有序;RocketMQ 可以支持普通顺序消息和严格顺序消息,顺序消息支持更加完善。
第五,消费模型。
Kafka 采用 Pull 模式,由 Consumer 主动拉取消息;RocketMQ 底层也是 Pull,但客户端封装成 Push 方式,对业务开发更加友好。
第六,应用场景。
Kafka 更适合日志采集、埋点分析、大数据平台、实时计算等场景;RocketMQ 更适合订单、支付、库存、交易等对可靠性要求较高的业务系统。
Kafka 的优势是高吞吐、高性能,适合大数据和流式计算;RocketMQ 的优势是业务功能丰富、可靠性高,更适合互联网核心业务消息场景。实际项目中,应根据业务特点选择合适的消息队列。
同一个分区只能由消费者组中的一个消费者消费的原因
Kafka 规定,在同一个 Consumer Group 内,一个 Partition 同一时刻只能分配给一个 Consumer 消费,主要是为了保证消费顺序、消费一致性和 Offset 管理的简单性。
首先,保证消息顺序。Kafka 只能保证 Partition 内消息有序,如果多个 Consumer 同时消费同一个 Partition,就无法保证消息按照 Offset 的顺序处理,消息可能出现乱序。
其次,避免重复消费和漏消费。一个 Consumer Group 共用一份 Offset,如果多个 Consumer 同时消费同一个 Partition,就会同时更新 Offset,容易导致 Offset 冲突,从而出现重复消费或漏消费的问题。
最后,简化消费协调机制。Kafka 通过 Rebalance 将 Partition 分配给不同 Consumer,每个 Partition 只对应一个 Consumer,这样 Consumer 之间无需竞争同一个 Partition,也不需要复杂的分布式锁,大大降低了协调成本。
因此,Kafka 选择了 “一个 Partition 对应一个 Consumer” 的消费模型,用顺序性换取更简单、更高效、更可靠的消费机制。
批次
Kafka 中的批次(Batch)是指 Producer 将多条消息先缓存起来,再一次性发送给 Broker,而不是每产生一条消息就立即发送。
采用批量发送主要有两个目的:减少网络请求次数和提高吞吐量。如果每条消息都单独发送,会频繁进行网络通信,性能较低;而将多条消息组成一个 Batch 后一次发送,可以有效降低网络开销,提高消息发送效率。
Kafka 的 Batch 是以 Partition 为单位进行维护的。Producer 会为每个 Partition 维护一个 Batch,同一个 Partition 的消息会先写入对应的 Batch,当 Batch 满了、等待时间到了,或者调用 flush() 时,再统一发送给 Broker。
Batch 的大小主要由两个参数控制:
- batch.size:单个 Batch 的最大容量,默认 16KB。
- linger.ms:等待组装 Batch 的最大时间。如果在规定时间内 Batch 没有装满,也会立即发送。
合理增大 batch.size 和 linger.ms 可以提高吞吐量,但也会增加消息发送的延迟,因此需要根据业务场景进行权衡。
消息堆积问题
消息堆积是指 Producer 的生产速度持续大于 Consumer 的消费速度,导致消息不断积压在 Broker 中。
Kafka 本身不会自动消除消息堆积,通常需要从提升消费能力、提升写入能力和优化业务处理三个方面进行解决。
首先,提升消费能力。如果 Consumer 的处理能力不足,可以增加 Consumer 实例数量,提高消费并行度。但需要注意,消费并行度最终受 Partition 数量限制,因此如果 Consumer 数量已经达到 Partition 数量,就需要增加 Partition。
其次,增加 Partition 数量。增加 Partition 后,可以部署更多 Consumer,实现更高的并行消费能力。不过新增 Partition 只能提升后续消息的并行处理能力,不会自动重新分布已有消息。
再次,优化 Consumer 的业务逻辑。如果消费过程中存在数据库慢查询、远程调用耗时等问题,应优化业务代码,采用异步处理、批量处理等方式,提高单个 Consumer 的消费速度。
最后,适当提高 Broker 的写入和读取性能。例如增加 Broker 节点、优化磁盘和网络配置,可以提升整个 Kafka 集群的吞吐能力,但如果瓶颈在 Consumer,这种优化效果有限。
Kafka 解决消息堆积的核心思路是:先定位瓶颈,如果是消费能力不足,就增加 Consumer 和 Partition;如果是业务处理慢,就优化消费逻辑;如果是集群性能不足,再扩容 Broker,提高整体吞吐能力。
高吞吐量的原因
Kafka 之所以具有很高的吞吐量,主要是因为它在存储、网络和并发模型等方面都进行了优化。
首先,顺序写磁盘。Kafka 采用追加写(Append Only)的方式将消息顺序写入磁盘,避免了随机 I/O,即使数据存储在磁盘中,也能够获得接近内存的写入性能。
其次,零拷贝技术。Kafka 使用操作系统提供的 sendfile 等零拷贝机制,将数据直接从内核缓冲区发送到网卡,减少了 CPU 数据拷贝和用户态、内核态切换,提高了数据传输效率。
第三,批量发送和批量拉取。Producer 会将多条消息组成 Batch 后再发送,Consumer 也会批量拉取消息,减少了网络请求次数,提高了网络利用率。
第四,Partition 并行机制。一个 Topic 可以划分为多个 Partition,不同 Partition 可以分布在不同 Broker 上,实现消息的并行写入和并行消费,大幅提升系统吞吐量。
最后,充分利用操作系统 Page Cache。Kafka 大部分读写都发生在操作系统的 Page Cache 中,减少了磁盘 I/O,使热点数据能够直接从内存读取,进一步提升了性能。
最核心的是顺序写磁盘和Partition 并行机制。顺序写保证了单机的高写入性能,Partition 并行保证了整个集群能够横向扩展,这两点是 Kafka 高吞吐的基础。
AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。
更多推荐



所有评论(0)