学习目标
学完本章你应该能够:
- 用自己的话讲清 Kafka 是什么,以及它和 RabbitMQ、ActiveMQ 等消息队列的取舍逻辑。
- 画得出 Kafka 的整体架构图(Producer / Broker / Consumer / Topic / Partition / ZooKeeper 的关系)。
- 讲清多副本机制里 leader/follower 与 ISR 的协同,以及"消息不丢失"到底靠哪些参数兜底。
- 解释 partition 内有序、pull 拉模型、零拷贝存储这些设计为什么能让 Kafka 扛住高吞吐。
- 面试时把"Kafka 如何保证不丢、不乱序、节点怎么判活"讲成一个完整的故事。
前置知识:分布式系统基础概念(集群、主从复制)、消息队列"生产者/消费者"基本模型。
本章你会动手做的事:
- 在纸上画一张 Kafka 集群拓扑图,标出 topic、partition、replica 落在哪台 broker。
- 用一句话总结 acks=-1 与 min.insync.replicas 的配合关系。
- 对比"推"和"拉"两种消费模型,各写一条优缺点。
一、Kafka 是什么?主要应用场景有哪些
1.1 用生活类比先建立直觉
类比:Kafka 就像一个物流分拣中心——商家(生产者)把包裹源源不断丢进传送带,分拣中心按目的地(topic/partition)把包裹码放整齐,快递员(消费者)按自己的节奏来取件。分拣中心不关心你什么时候来取,它只负责"先收好、摆整齐、不丢件"。
对应到工程里,Kafka 官方给自己的定位是高吞吐的分布式流平台(distributed streaming platform):它既能当消息队列做异步解耦,又能当流处理管道做实时计算的数据源,还能当存储系统持久化事件日志。
1.2 它到底解决什么问题
传统消息队列也在做"异步、解耦、削峰",但 Kafka 把"吞吐量和持久化“做到了极致:
- 高吞吐:单机就能扛住每秒几十万条消息的写入(得益于顺序写磁盘、批量发送、零拷贝)。
- 可持久化:消息写入磁盘并按 retention 策略保留,消费完也不删(默认可回放)。
- 分布式:天生为集群设计,水平扩容几乎线性。
1.3 主要应用场景
| 场景 | 说明 |
|---|---|
| 日志 / 埋点采集 | 把业务系统、服务器产生的海量日志统一汇聚到 Kafka,下游再落 ES / HDFS |
| 系统解耦 | 订单系统只管往 Kafka 发"下单事件”,库存、积分、推荐各消费自己的副本,互不阻塞 |
| 流量削峰 | 秒杀时瞬时流量先堆进 Kafka,下游按自身处理能力匀速消费,保护数据库 |
| 流处理 / 事件溯源 | 配合 Flink、Kafka Streams 做实时聚合、监控告警 |
考点总结:Kafka 不只是"消息队列",更准确的定位是分布式流平台;面试被问"用在哪"时,优先答日志采集、解耦、削峰、流处理四件套。
二、和其他消息队列相比,Kafka 的优势在哪里
2.1 先看清对比对象
常见的消息中间件大致分两类:
- RabbitMQ / ActiveMQ:主打低延迟、复杂路由、事务消息,单队列吞吐相对低。
- RocketMQ:阿里开源,兼顾高吞吐与事务/定时消息,偏业务场景。
- Kafka:主打超高吞吐、持久化、水平扩展,路由模型简单(靠 topic/partition)。
2.2 Kafka 的核心优势
flowchart LR
A[Kafka 优势] --> B[吞吐量极高
顺序写+批量+零拷贝]
A --> C[消息持久化
落盘可回放]
A --> D[水平扩展
加 broker 即可扩容]
A --> E[拉模式消费
consumer 自控节奏]- 吞吐高:顺序写磁盘 + 批量发送 + 零拷贝(sendfile)三重加持,单机轻松十万级~百万级 QPS。
- 持久化:消息落到磁盘 log 文件,按 retention 保留,消费位移独立管理(可重复消费、可回放)。
- 水平扩展:topic 拆成多个 partition 分散到不同 broker,加机器就能加吞吐,几乎线性。
- 拉模式(pull):consumer 主动来拉,broker 不被慢消费者拖垮,还能做批量拉取。
考点总结:Kafka 的差异化优势 = 高吞吐 + 持久化 + 水平扩展 + 拉模式;劣势是"实时性不如 RabbitMQ、路由没有那么花哨",选型要看场景。
三、核心概念:Producer / Consumer / Broker / Topic / Partition
3.1 类比:邮局系统
- Producer(生产者) = 寄信人,负责把消息投进 Kafka。
- Consumer(消费者) = 收信人,从 Kafka 读取消息处理。
- Broker(代理节点) = 邮局分局,Kafka 集群里的一台服务器,负责存消息。
- Topic(主题) = 信箱类别,比如"账单信箱"“广告信箱”,是消息的逻辑分类。
- Partition(分区) = 信箱里的多个抽屉,一个 topic 可以拆成多个 partition 并行存、并行读。
3.2 架构关系图
flowchart LR
P[Producer] -->|发送| T1[topic-A]
T1 --> PA1[partition-0]
T1 --> PA2[partition-1]
PA1 --> B1[(Broker 1)]
PA2 --> B2[(Broker 2)]
B1 --> C1[Consumer Group]
B2 --> C1要点:
- 一个 topic 是逻辑概念,下面由若干 partition 组成;partition 才是真正的存储与并行单元。
- 同一 topic 的不同 partition 可以分布在不同 broker 上,实现负载分散。
- 一条消息在 partition 内会被分配一个递增的 offset(位移),offset 唯一标识它在分区里的位置。
考点总结:topic 是逻辑分类,partition 是物理并行单元;offset 只在"单个 partition 内"有意义,跨 partition 不保证全局顺序。
四、多副本机制:Replication、Leader/Follower 与 ISR
4.1 为什么要副本
单台 broker 挂了,上面的 partition 数据就没了。Kafka 给每个 partition 配了多个副本(replica),数量由 replication.factor(副本因子)决定,比如设为 3 就表示一份数据存 3 台机器。
4.2 Leader / Follower 模型
flowchart TB
subgraph PF[partition-0 的三个副本]
L[Leader 副本
负责读写]
F1[Follower 副本1]
F2[Follower 副本2]
end
P[Producer] -->|只发 Leader| L
L -->|同步| F1
L -->|同步| F2
C[Consumer] -->|只从 Leader 读| L- 每个 partition 的多个副本里,有且仅有一个 Leader 副本,负责所有读写。
- 其余副本是 Follower,只做一件事:从 Leader 异步拉取数据保持同步(不对外服务读)。
- 这叫"主写主读"——和 MySQL 主从"主写从读"不同,Kafka 的 Follower 不参与读,简化了一致性。
4.3 ISR(In-Sync Replicas,同步副本集合)
不是所有 Follower 都算"靠谱的"。Kafka 维护一个 ISR 列表,里面是"和 Leader 差距不超过 replica.lag.time.max.ms 阈值"的副本(含 Leader 自己)。
- 只有 ISR 里的副本才有资格在 Leader 挂掉时被选为新 Leader。
- 某个 Follower 同步太慢 / 心跳断了,会被踢出 ISR;追上后又会重新加入。
考点总结:副本因子决定冗余度;Leader 统管读写,Follower 只同步;ISR 是"可靠的同步副本小圈子",新 Leader 只能从 ISR 里选。
五、多分区与多副本的好处
5.1 多分区(Partition)的好处
- 并行度:partition 是并行单元,consumer group 里一个 partition 同一时刻只被一个 consumer 消费,partition 越多并发消费能力越强。
- 水平扩展:partition 可分散到不同 broker,突破单机容量/IO 上限。
- 负载均衡:写入按 key 哈希或轮询分散到各 partition,避免热点集中。
5.2 多副本(Replica)的好处
- 高可用:单 broker 宕机,该 partition 的 Leader 自动切到 ISR 里的 Follower,服务不中断。
- 数据可靠性:配合
acks=-1,消息要写入足够多的副本才算成功,降低丢失概率。
flowchart LR
subgraph 好处
A[多分区
并行+扩展+均衡]
B[多副本
高可用+可靠性]
end考点总结:分区解决"吞吐和扩展",副本解决"可用和不丢";两者正交——分区数管横向扩展,副本因子管冗余度。
六、ZooKeeper 在 Kafka 中的角色(早期版本)
注意:这是针对 Kafka 早期版本(2.8 之前,依赖外部 ZK)的经典面试题。新版 Kafka(KRaft 模式)已内置元数据仲裁,逐步去 ZK 化。
6.1 ZK 到底管了什么
flowchart TB
ZK[(ZooKeeper)]
ZK --> R[Broker 注册
谁在线]
ZK --> L[Controller 选举
选集群大脑]
ZK --> O[老版本 offset 存储
consumer 位移]
B1[Broker 1] -->|心跳注册| ZK
B2[Broker 2] -->|心跳注册| ZK- Broker 注册与发现:每台 broker 启动时在 ZK 上注册临时节点,集群由此知道"谁活着"。
- Controller(控制器)选举:ZK 负责从 broker 中选举出一个 Controller,由它统筹 partition 的 Leader 选举、副本分配等集群级决策。
- 老版本 offset 存储:早期 consumer 的位移(offset)存在 ZK 的
/consumers节点下(后来改为存在 Kafka 内部的__consumer_offsetstopic,不再依赖 ZK)。
考点总结:早期 ZK 负责 broker 注册、Controller/Leader 选举、老版 offset 存储;新版走 KRaft 去 ZK,但面试仍常考 ZK 职责。
七、如何保证消息的消费顺序
7.1 Kafka 的顺序语义边界
Kafka 只保证 partition 内有序,不保证跨 partition 全局有序。原因是:消息在单个 partition 里是append-only 的顺序日志,offset 递增;但多个 partition 之间是互相独立的并行流。
flowchart LR
M1[消息A] --> M2[消息B] --> M3[消息C]
subgraph partition-0 内严格有序
M1; M2; M3
end
P0[partition-0 有序]
P1[partition-1 独立] 7.2 如何做到"想有序就有序"
- 单 partition:所有相关消息发到同一个 partition,自然有序(但牺牲并行度)。
- 按 key 路由:Producer 指定
key(如用户ID),Kafka 默认用hash(key) % 分区数把"同一 key 的消息"固定落到同一个 partition,于是"同一用户的事件"保持顺序。 - Consumer 顺序消费:单个 partition 同一时刻只被一个 consumer 实例消费,所以消费端天然不会乱序拉取。
考点总结:Kafka 的秩序边界是 partition——partition 内严格有序,跨 partition 不保证;想保序就用相同 key 路由到同一分区。
八、如何保证消息不丢失
8.1 三个环节都要兜底
“不丢消息"要从生产端、服务端、消费端三侧一起看:
flowchart LR
P[Producer] -->|acks=-1
min.insync.replicas| B[(Broker 多副本)]
B -->|offset 提交时机| C[Consumer]8.2 生产端:acks 与副本数
acks=-1(或all):Producer 要等Leader 且所有 ISR 副本都写入成功才认为发送成功,可靠性最高。min.insync.replicas(默认 1,建议 ≥2):要求至少有这么多 ISR 副本写入,acks=-1 才有意义;否则 ISR 只剩 Leader 时仍可能丢。- 重试:
retries设大,配合enable.idempotence=true开启幂等,避免网络抖动导致丢或重。
8.3 消费端:offset 提交时机
- 先处理再提交:务必等业务逻辑处理成功后再提交 offset。如果"先提交 offset 再处理”,处理中途崩溃就会丢消息(offset 已前进,崩了的那批不会再消费)。
- 反之"处理完才提交"最多导致重复消费(at-least-once 语义),比丢消息好处理。
考点总结:不丢 = acks=-1 + min.insync.replicas≥2 + 生产重试/幂等 + 消费端"处理完再提交 offset";Kafka 默认是 at-least-once,精确一次要靠幂等+事务。
九、Kafka 如何判断一个节点是否还活着
Kafka(早期依赖 ZK)判定一个 broker 节点"存活"需要同时满足两个条件:
- 与 ZooKeeper 保持心跳(session 未过期):broker 定期向 ZK 发送心跳,ZK 认为 session 还活着。
- 它是 Follower 时,必须能及时与 Leader 副本保持同步:即该节点(作为某 partition 的 Follower)落后于 Leader 的差距在允许范围内,没有被踢出 ISR。
注意:第 2 条本质是说——即使 broker 进程还在、ZK 心跳还在,如果它作为 Follower 拷贝数据太慢(落后超
replica.lag.time.max.ms),也会被认定"跟不上",从而被移出 ISR,丧失成为新 Leader 的资格。
考点总结:节点存活 = ①和 ZK 心跳未断 ②作为 Follower 能跟上 Leader 同步;两条缺一则视为"不够健康"。
十、Producer 是否直接将数据发送到 Broker 的 Leader
是的,Producer 只把消息发给对应 partition 的 Leader 副本,不会发给 Follower。
原因很直接:Kafka 是"主写主读"模型——只有 Leader 负责读写,Follower 只被动从 Leader 拉数据同步。如果 Producer 写到 Follower,还要再转发一次到 Leader,既多一跳延迟又让写入路径变复杂。所以客户端会先从 broker 拉取元数据(哪台 broker 是某 partition 的 Leader),然后直连 Leader 发送。
sequenceDiagram
participant P as Producer
participant B as Broker(元数据)
participant L as Leader副本
participant F as Follower副本
P->>B: 拉取 topic 分区 Leader 信息
B-->>P: 返回 Leader 所在 broker
P->>L: 直接发送消息到 Leader
L->>F: 异步同步给 Follower考点总结:Producer 直连 Leader 写,Follower 只同步不接收生产者的直接写入;客户端靠元数据缓存知道 Leader 在哪。
十一、Consumer 能否消费指定分区消息
可以。 Consumer 默认由消费者组(Consumer Group)的 rebalance 自动分配分区,但你也能绕过自动分配,手动指定要消费的分区。
方式是通过 consumer.Assign()(Java 客户端为 assign())显式传入 TopicPartition 列表,而不是用 subscribe() 订阅整个 topic。典型场景:
- 需要严格按 partition 顺序单线程处理某几个分区;
- 做重置/补数据时,只重放某个特定分区;
- 想脱离"组消费位移"机制,自己管理 offset。
⚠️ 新手必踩的坑:一旦用
assign()手动分配,该 consumer 就不参与组的 rebalance,group 的自动位移管理也基本失效,offset 得自己维护(比如提交到__consumer_offsets或外部存储)。
考点总结:能指定——用 assign() 显式绑定 TopicPartition;代价是失去组内自动均衡和自动位移管理,offset 要自己管。
十二、Kafka 高效文件存储设计特点
Kafka 能扛高吞吐,存储层的设计是关键。核心有四招:
flowchart LR
A[分段 Segment] --> B[顺序写磁盘]
B --> C[稀疏索引]
C --> D[零拷贝 sendfile]- Partition 分段(Segment):一个 partition 的日志物理上切成多个固定大小(默认 1GB)的 segment 文件,旧 segment 只读不写,便于清理和查找。
- 顺序写磁盘:消息只追加(append-only),避免随机写,磁盘顺序写性能接近内存。
- 稀疏索引:每个 segment 配
.index索引文件,只记录"每隔若干字节的消息 offset → 物理位置",不是每条都建索引,兼顾查找速度与空间。 - 零拷贝(Zero-Copy):consumer 读消息时用
sendfile系统调用,数据直接从磁盘页缓存经网卡发出,不经过用户态内存拷贝,大幅降低 CPU 和延迟。
考点总结:高效存储四件套 = 分段 + 顺序写 + 稀疏索引 + 零拷贝;本质是把"磁盘"用出了"顺序 + 不拷贝"的接近内存的速度。
十三、Partition 的数据如何保存到硬盘
Partition 在磁盘上以一组 log segment 文件形式存在,每个 segment 由三个配套文件组成:
flowchart TB
subgraph 一个 Partition 的磁盘文件
LOG[xxx.log
消息体,顺序append]
IDX[xxx.index
偏移量→物理位置 稀疏索引]
TIDX[xxx.timeindex
时间戳→偏移量 时间索引]
end.log文件:真正存消息内容,按批次顺序追加;达到log.segment.bytes(默认 1GB)就滚动新建下一个 segment。.index文件:偏移量索引,记录"某 offset 相对该 segment 基址的字节位置",用于快速定位消息。.timeindex文件:时间戳索引,支持"按时间找 offset"(如offsetForTimes)。
⚠️ 新手必踩的坑:不要把
.log当成"数据库表"去随机改——它是只追加的,删除/清理是按 segment 整体过期(retention)或按 compact(相同 key 只留最新)策略做的,不能原地修改单条。
考点总结:partition 落盘 = 多个 segment(.log + .index + .timeindex);.log 存消息、.index 按偏移稀疏索引、.timeindex 按时间索引起定位与清理作用。
十四、Producer 的分区路由策略
Producer 发消息时要决定"这条消息进哪个 partition",策略优先级如下:
flowchart LR
K{指定了 key?}
K -->|是| H[hash(key) % 分区数
相同 key 固定分区]
K -->|否| R[粘性分区/轮询
均匀分散]
C{自定义 Partitioner?}
C -->|是| U[走用户逻辑]- 显式指定 partition:调用时直接传 partition 编号,优先级最高,完全由你控制。
- 按 key 哈希(默认):没指定 partition 但指定了 key,则用
hash(key) % 分区数决定,保证相同 key 进同一分区(保序关键)。 - 无 key 轮询 / 粘性分区:没 key 时,新版本用"粘性分区(sticky partitioner)"——先粘在一个分区批量发,攒够一批再换,减少碎片化、提升批效率(老版本是轮询)。
- 自定义 Partitioner:实现
Partitioner接口,按业务自定义(如按地区、按用户等级)。
考点总结:路由优先级 = 指定分区 > key 哈希 > 轮询/粘性;key 哈希是"保序 + 均匀分布"的核心手段。
十五、Consumer 是推还是拉
Kafka 的 Consumer 是"拉(pull)“模式——consumer 主动向 broker 请求"给我下一批消息”,而不是 broker 主动推送。
flowchart LR
subgraph 推 PUSH
S1[Broker] -->|主动塞| C1[慢 Consumer 被压垮]
end
subgraph 拉 PULL
C2[Consumer] -->|按节奏请求| S2[Broker]
end为什么选 pull 而不是 push:
- 保护消费者:push 模式下,如果 consumer 处理慢,broker 还不停推,会把 consumer 内存/连接压垮;pull 让 consumer “能吃多少取多少”,按自己速率消费。
- 批量友好:consumer 可以一次拉一批(
fetch.min.bytes/max.poll.records),提升吞吐。 - 缺点:consumer 要不停轮询,空轮询有少量开销;Kafka 用"长轮询(不足量时 broker hold 住请求一小会儿)“缓解。
考点总结:Kafka 是 pull 拉模型,核心动机是"别把慢消费者压垮"并支持批量取;代价是需要轮询,靠长轮询补偿。
十六、消费状态跟踪:Offset 的提交方式
Consumer “读到哪了"靠 offset(位移) 记录,位移的持久化方式就是"提交(commit)"。
16.1 自动提交
enable.auto.commit=true(默认),配合auto.commit.interval.ms(默认 5s),consumer 周期性把"已拉取"的 offset 自动提交到__consumer_offsets。- 坑:自动提交的是"拉到的位置"而非"处理完的位置”,如果拉到后还没处理就崩了,会丢消息。
16.2 手动提交
enable.auto.commit=false,业务处理成功后调用commitSync()(同步、稳)或commitAsync()(异步、快)。- 这是保证"处理完再提交”(见第八章不丢消息)的关键手段。
flowchart LR
A[消费消息] --> B{处理成功?}
B -->|是| C[手动 commit offset]
B -->|否| D[不提交
下次重拉]考点总结:offset 跟踪两路——自动提交(简单但有丢消息风险)、手动提交(处理完再 commit,保不丢);位移实际存在 Kafka 内部 topic __consumer_offsets。
十七、Kafka Tiered Storage 在 Kubernetes 中配置 S3 凭证而不暴露 Secret
17.1 用生活类比先建立直觉
类比:Kafka 的分层存储(Tiered Storage)像"图书馆把旧书搬进市立档案馆(S3)"。平时热读的书(本地 SSD)留在馆内,查得少的旧书归档到档案馆,省下馆内空间。问题来了:去档案馆取书要"门禁卡"(S3 凭证)。如果把门禁卡密码直接写在墙上的告示(明文 ConfigMap)或随便给每个馆员配一张万能卡(长期 AK/SK 塞进 Secret 并挂载到所有 Pod),一旦墙被拍、卡被偷,档案馆就门户大开。
正确做法是:门禁卡由档案馆自动发给"有正当事由的馆员"(IAM Role 绑定到 ServiceAccount,Pod 通过 IRSA 自动拿到临时凭证),而且这张卡不落地、不写在墙上、不挂进容器文件系统——这就是"不暴露 Secret"的核心。
flowchart LR
Pod[Kafka Broker Pod] -->|绑定| SA[ServiceAccount
带 IRSA 注解]
SA -->|AssumeRole| STS[AWS STS 临时凭证
不落盘]
Pod -->|用临时凭证| S3[(S3 分层存储桶)]
Cfg[ConfigMap
只放非敏感配置] --> Pod这张图是"无 Secret 暴露"的凭证流:Pod 绑 ServiceAccount → 经 IRSA 换 STS 临时凭证 → 直连 S3;敏感凭证从不出现在 ConfigMap 或挂载卷里。
17.2 工程要点
步骤1:给 ServiceAccount 绑定 IAM Role(IRSA,EKS 为例)
凭证不进 Secret,而是让 Pod 的 ServiceAccount 关联一个 IAM Role,由云厂商在运行时注入临时凭证(环境变量 AWS_ROLE_ARN / AWS_WEB_IDENTITY_TOKEN_FILE)。
# kafka-sa-irsa.yaml
apiVersion: v1
kind: ServiceAccount
metadata:
name: kafka-broker
namespace: kafka
annotations:
# 步骤1:把 SA 绑定到拥有 S3 访问权的 IAM Role(由 Terraform/eksctl 预先建好)
eks.amazonaws.com/role-arn: "arn:aws:iam::123456789012:role/kafka-tiered-s3"
# 步骤2:只给分层存储所需的最小权限(s3:PutObject/GetObject/ListBucket 等)
eks.amazonaws.com/sts-regional-endpoints: "true"
对应的 IAM Role 信任策略只允许这个 SA 的 OIDC 身份
AssumeRole,且权限策略只含s3:GetObject、s3:PutObject、s3:ListBucket、s3:DeleteObject作用于指定的分层存储桶——最小权限,无长期密钥。
步骤2:用 ConfigMap 只放"非敏感"的分层存储配置
把不带密码的配置放进 ConfigMap(即使被读也不泄密),敏感部分交给 IRSA 自动注入:
# kafka-tiered-config.yaml(ConfigMap,可公开,无 Secret)
apiVersion: v1
kind: ConfigMap
metadata:
name: kafka-tiered-config
namespace: kafka
data:
server.properties: |
# 步骤1:开启分层存储
remote.storage.enable=true
# 步骤2:选用 S3 作为远程存储后端
remote.log.storage.system=S3
# 步骤3:桶名与区域(非敏感)
s3.bucket.name=kafka-tiered-archive
s3.region=us-east-1
# 步骤4:凭证交给"默认链 + IRSA"——不写 AK/SK!
s3.credentials.provider.class=com.amazonaws.auth.DefaultAWSCredentialsProviderChain
# 步骤5:本地保留时长,超时段才上云
local.log.segment.bytes=1073741824
log.retention.ms=86400000
remote.log.retention.ms=2592000000
步骤3:Broker Pod 绑定 SA、不挂任何 Secret 卷
# kafka-broker-deployment.yaml(节选)
apiVersion: apps/v1
kind: StatefulSet
metadata:
name: kafka
namespace: kafka
spec:
serviceName: kafka
replicas: 3
template:
spec:
serviceAccountName: kafka-broker # 步骤1:绑定带 IRSA 的 SA
containers:
- name: kafka
image: confluentinc/cp-kafka:7.6.0
env:
# 步骤2:把 ConfigMap 的配置挂为环境变量(不含任何密码)
- name: KAFKA_REMOTE_STORAGE_ENABLE
value: "true"
volumeMounts:
- name: tiered-config
mountPath: /etc/kafka/tiered
# 步骤3:注意——这里【没有任何】secret 卷挂载,
# S3 凭证由 IRSA 经 STS 自动注入到 AWS_WEB_IDENTITY_TOKEN_FILE
volumes:
- name: tiered-config
configMap:
name: kafka-tiered-config
⚠️ 新手必踩的坑: 千万不要为了"省事"把 AK/SK 写进 ConfigMap 或
server.properties明文——ConfigMap 在集群内默认所有能读 namespace 的人都能看,且会进 etcd 明文。更别把 Secret 以subPath挂成文件后还让容器把它cat出来打日志。分层存储的正确姿势是 IRSA/OIDC + 临时凭证,Secret 只应在"无法用 IRSA 的环境(如自建 K8s)“才用,且必须是加密 etcd + 定期轮转。
考点总结
Kafka Tiered Storage 把冷段存到 S3,但 S3 凭证绝不能进 ConfigMap 明文或随意挂 Secret。K8s 上最佳实践是用 IRSA(IAM Roles for Service Accounts):ServiceAccount 注解绑定 IAM Role,Pod 运行时由云厂商注入临时 STS 凭证(环境变量 + token 文件),从不落盘、不出现在配置里。server.properties 里只配 s3.bucket.name 等非敏感项,并用 DefaultAWSCredentialsProviderChain 自动走临时凭证链。自建集群无 IRSA 时才退用加密 Secret + 定期轮转。一句话:敏感凭证走"身份换临时令牌”,配置走"可公开 ConfigMap"。
十八、基于 Cruise Control 在分层存储场景下重新平衡副本的 JSON 参数
18.1 用生活类比先建立直觉
类比:Kafka 集群像一座仓库,分区副本是货架上的货箱。随着货物进出(分区写入、节点上下线),货箱会越摆越歪——有的货架(broker)挤爆、有的空着,甚至"冷热货箱"全堆在同一排(机架)。Cruise Control 就是"智能理货员",你给它一张需求清单(goals 参数),它算出一套搬运方案,把货箱重新摆匀。
在分层存储场景下,货箱分两种:本地 SSD 上的"热货箱"(体积小、要快)和 S3 上的"冷货箱"(几乎不限容量)。理货员重新平衡时,不能再用"本地磁盘均等"这一条老规矩——得兼顾"本地盘别压满"和"副本分散防丢",于是需求清单要相应调整。
flowchart TB
User[运维/自动控制器] -->|POST 带 goals JSON| CC[Cruise Control]
CC -->|计算搬运方案| Plan[分区副本重平衡计划]
Plan --> Exec[执行迁移
副本跨 broker 重新分布]
Exec --> K[Kafka 集群
本地 SSD + S3 分层]这张图是 Cruise Control 重平衡闭环:提交 goals(JSON)→ 计算计划 → 执行迁移 → 集群更均衡。
18.2 工程要点
Cruise Control 的 goals 与分层存储的关注点
Cruise Control 用一组 Goal 描述"什么是均衡"。常用 goals:
| Goal | 作用 | 分层存储下注意 |
|---|---|---|
RackAwareGoal | 副本跨机架分布,防机架故障丢数据 | 必须保留,数据安全底线 |
ReplicaDistributionGoal | 每个 broker 的副本数尽量均等 | 避免个别 broker 承载过多分区 |
DiskUsageDistributionGoal | 各 broker 本地磁盘使用率均衡 | 分层存储下本地盘只存热段,阈值要调小,避免本地 SSD 写满 |
NetworkInboundUsageDistributionGoal | 各 broker 入站网络流量均衡 | 重平衡本身会产生迁移流量,需限速 |
TopicReplicaDistributionGoal | 同一 topic 的副本分散 | 防单 topic 热点集中 |
步骤1:用 JSON 提交分层存储重平衡请求
Cruise Control 的 rebalance 接口接受 goals 等参数(可用 JSON body 描述本次目标),下面是一份"分层存储友好"的重平衡参数:
{
"goals": [
"RackAwareGoal",
"ReplicaDistributionGoal",
"DiskUsageDistributionGoal",
"NetworkInboundUsageDistributionGoal",
"TopicReplicaDistributionGoal"
],
"excluded_topics": [
".*_internal",
"__consumer_offsets"
],
"exclude_recently_demoted_brokers": true,
"exclude_recently_removed_brokers": true,
"allow_capacity_estimation": true,
"throttle_removal": false,
"verbose": true
}
# 步骤2:把上面的 JSON 作为请求体,POST 给 Cruise Control 的 rebalance 接口
curl -X POST "http://cruise-control.kafka.svc:9090/kafkacruisecontrol/rebalance" \
-H "Content-Type: application/json" \
-d @rebalance_goals.json \
--max-time 600
# 返回里会带 dryRun 计划;确认无误后再带 &dryRun=false 真正执行
说明:Cruise Control 官方 rebalance 接口主要用 query 参数(
?goals=...&allow_capacity_estimation=true)传 goals,但也可以用POST+ JSON body 描述更复杂的需求(如excluded_topics、是否 demote broker)。上面 JSON 是工程上"分层存储场景"推荐的 goals 组合,可直接落地。
步骤3:针对分层存储调小本地磁盘阈值
分层存储下,本地 SSD 只存"尚未上云"的热段,容量比全量小得多。如果 DiskUsageDistributionGoal 还按"磁盘快满才均衡"的老阈值,本地盘很容易写爆。要在 Cruise Control 配置里调低触发均衡的本地使用率阈值:
{
"capacity.estimation.info.configs": "log.dirs",
"disk.usage.distribution.threshold": 0.3,
"goal.utilization.balancer": "adaptive",
"default.goals": [
"RackAwareGoal",
"ReplicaDistributionGoal",
"DiskUsageDistributionGoal",
"NetworkInboundUsageDistributionGoal"
]
}
# 步骤4:查看当前集群均衡状态(是否还有 broker 本地盘偏满)
curl -s "http://cruise-control.kafka.svc:9090/kafkacruisecontrol/load" \
| jq '.brokerLoad[].diskUsager'
⚠️ 新手必踩的坑: 分层存储不能取消
RackAwareGoal和ReplicaDistributionGoal——冷数据在 S3 上虽然"无限容量",但副本元数据、消费者位移、控制器选举仍然依赖 broker 本地与跨机架分布。只盯着"本地盘均衡"会忽略"副本全落同一机架"的致命风险。重平衡时务必把RackAwareGoal放在 goals 列表最前。
考点总结
Cruise Control 通过一组 Goal 重平衡 Kafka 副本分布,接口可用 JSON body 描述 goals、excluded_topics 等需求。分层存储场景下,冷段在 S3(近乎无限),但本地 SSD 只存热段、容量有限,所以:① DiskUsageDistributionGoal 的本地使用率阈值要调小(如 0.3),防止本地盘写爆;② RackAwareGoal、ReplicaDistributionGoal 必须保留在 goals 最前,副本跨机架与跨 broker 分布是数据安全的底线,不能因为"冷数据在云上"就省略;③ 重平衡会产生迁移流量,配合 NetworkInboundUsageDistributionGoal 限速,先 dryRun 确认再执行。
十九、本地 SSD 与 S3 一致性延迟时监控并告警分区高水位
19.1 用生活类比先建立直觉
类比:Kafka 分层存储就像"书店把旧书搬去市立档案馆(S3)"。店内的书架(本地 SSD)只放近期热书。问题是:店员把书搬去档案馆的这段路上会有延迟——有时候书架上的书已经卖光/更新了,档案馆里还是旧的一本(一致性延迟)。更危险的是:如果"书架快堆满了"(分区本地保留量逼近高水位 / 磁盘容量上限),新书就上不了架,整个店停摆。
所以我们需要两路监控:① 搬运延迟——本地段上传到 S3 落后多少(一致性延迟);② 高水位——本地磁盘 / 分区保留量是不是快顶到天花板了。任一越界就告警,运维提前介入。
flowchart LR
Broker[Kafka Broker] -->|热段在本地 SSD| Local[本地磁盘]
Broker -->|异步上传冷段| S3[(S3)]
Broker -->|JMX: RemoteLogManager 指标| Exporter[Prometheus Exporter]
Exporter --> Prom[Prometheus]
Prom -->|Rule 触发| Alert[Alertmanager 告警]
Local -->|磁盘使用率>高水位| Prom这张图是监控闭环:Broker 暴露分层存储 JMX 指标 → Exporter 抓到 Prometheus → 规则判定延迟/水位 → Alertmanager 告警。
19.2 工程要点
关键监控指标:分层存储的一致性延迟与高水位
| 指标 | 含义 | 告警阈值建议 |
|---|---|---|
kafka.server:type=RemoteLogManager,name=RemoteCopyLagBytes | 还没上传到 S3 的本地段字节数 | 持续 > 1GB 告警(一致性延迟大) |
kafka.server:type=RemoteLogManager,name=RemoteCopyLagSegments | 未上传段数量 | 持续 > 50 告警 |
kafka.log:type=Log,name=Size / 单分区本地保留大小 | 单分区本地占用 | 接近 log.retention.bytes 高水位告警 |
节点磁盘使用率(node_filesystem_avail) | 本地 SSD 剩余 | > 80% 告警,> 90% 严重 |
一致性延迟的本质:本地段要等"段被关闭(rolled)“且超过
local.log.retention.ms才上传 S3。如果写入突增、段迟迟不关,或 S3 写入慢,就会出现"本地已更新、S3 还是旧的"的时间窗——消费端若读冷段可能拿到旧数据。
步骤1:Prometheus 告警规则(分区高水位 + 上传滞后)
# kafka-tiered-alerts.yaml
apiVersion: monitoring.coreos.com/v1
kind: PrometheusRule
metadata:
name: kafka-tiered-alerts
namespace: monitoring
spec:
groups:
- name: kafka-tiered-storage
rules:
# 规则1:S3 上传滞后(一致性延迟)过大
- alert: KafkaRemoteCopyLagHigh
expr: |
kafka_server_remotelogmanager_remotecopylagbytes > 1073741824
for: 10m
labels:
severity: warning
annotations:
summary: "Kafka 分层存储 S3 上传滞后超 1GB(一致性延迟大)"
description: "broker {{ $labels.instance }} 还有 {{ $value | humanize1024 }}B 未上传 S3,持续 10 分钟。"
# 规则2:本地磁盘(分区高水位)即将写满
- alert: KafkaLocalDiskHighWatermark
expr: |
100 * (1 - node_filesystem_avail_bytes{fstype=~"ext4|xfs"}
/ node_filesystem_size_bytes{fstype=~"ext4|xfs"}) > 80
for: 5m
labels:
severity: critical
annotations:
summary: "Kafka broker 本地磁盘使用率超 80%(逼近高水位)"
description: "节点 {{ $labels.instance }} 本地 SSD 使用率 {{ $value }}%,分层存储上传跟不上写入,需扩容或加速上传。"
# 规则3:单分区本地保留量逼近 retention 上限
- alert: KafkaPartitionLocalRetentionNearLimit
expr: |
kafka_log_size / on(topic,partition) kafka_config_retention_bytes > 0.9
for: 15m
labels:
severity: warning
annotations:
summary: "分区 {{ $labels.topic }}-{{ $labels.partition }} 本地保留量达上限 90%"
description: "分区本地段即将触顶,冷段上传若滞后将丢本地、读请求被迫回源 S3 变慢。"
步骤2:用脚本快速抓取"分区高水位"做即时巡检
不想等 Prometheus 也能用命令行即时看哪些分区本地快满:
#!/usr/bin/env bash
# check-partition-highwatermark.sh —— 找出本地保留量接近 retention 的分区
# 步骤1:从 JMX / 指标接口拿到分区本地大小与 retention 上限
curl -s "http://kafka-broker:7071/metrics" | grep -E "kafka_log_size|retention_bytes" > /tmp/metrics.txt
# 步骤2:简单比对,超过 90% 就打印
awk -F'[{}= ]' '
/kafka_log_size/ { size[$NF]=$(NF-1) }
/retention_bytes/ { ret[$NF]=$(NF-1) }
END {
for (p in size) {
if (p in ret && ret[p] > 0 && size[p]/ret[p] > 0.9) {
print "WARN 分区", p, "本地保留", size[p], "/", ret[p], "超过 90% 高水位"
}
}
}' /tmp/metrics.txt
⚠️ 新手必踩的坑: “分区高水位"在 Kafka 里有两个容易混淆的概念:① High Watermark(HW) = 已提交位移水位(消费者可见的最大 offset),和磁盘容量无关;② 本地保留量 / 磁盘使用率逼近上限 才是这里说的"高水位风险”(本地 SSD 写满)。监控脚本和告警规则盯的是后者——磁盘/保留量高水位,别和副本 HW 搞混。
考点总结
Kafka 分层存储下要同时盯两件事:① S3 一致性延迟(本地段上传 S3 的滞后,指标 RemoteCopyLagBytes/RemoteCopyLagSegments,持续偏大说明冷段上传跟不上写入);② 本地高水位风险(本地 SSD / 单分区保留量逼近 log.retention.bytes 或磁盘容量上限,可能写满停摆)。用 Prometheus 规则对"上传滞后 >1GB / 磁盘 >80% / 单分区保留 >90%“三类越界告警,或脚本即时巡检。注意区分”保留量高水位(磁盘风险)“和副本的”High Watermark(位移水位)"——监控指的是前者。本质:分层存储把容量压力从本地挪到云,但本地热段仍有容量上限,必须靠监控提前发现上传滞后与写满风险。
二十、自测题与动手练习
自测题(合上书能答出来,才算懂):
- Kafka 的"高吞吐"靠哪几项存储/发送设计支撑?为什么它比 RabbitMQ 更适合大数据日志场景?
- partition 内有序还是全局有序?如果想让"同一用户"的事件严格有序,Producer 该怎么发?
- 画出 Leader/Follower 与 ISR 的关系,并说明:ISR 里的副本挂了、ISR 外的 Follower 挂了,分别有什么影响?
acks=-1+min.insync.replicas=1能完全保证不丢消息吗?为什么?消费端还要做什么配合?- Consumer 是推还是拉?为什么 Kafka 选 pull 而不是 push?
动手练习(建议真做一遍):
- 在纸上画一张 Kafka 集群拓扑:1 个 topic、3 个 partition、副本因子 3,标出每个 partition 的 Leader 落在哪台 broker、Follower 在哪。
- 用你熟悉的语言写一段 Producer 代码,设置
acks=all、指定 key 并观察相同 key 是否进同一分区;再改成assign()手动消费某个分区。 - 对比"自动提交 offset"和"手动提交 + 处理完再 commit"两种写法,各写一句什么场景下会丢消息、什么场景下会重复消费。
二十一、本章小结
- Kafka 是高吞吐分布式流平台,靠分区并行、顺序写、零拷贝、拉模型把吞吐和可靠性拉满,典型用于日志采集、解耦、削峰、流处理。
- 可靠性三支柱:多副本(ISR 选主)+ acks=-1/min.insync.replicas(生产端不丢)+ 处理完再提交 offset(消费端不丢)。
- 秩序边界在 partition:partition 内严格有序,跨分区无序;保序靠"相同 key 路由到同一分区”。
- 下一篇我们可以深入 Kafka 的精确一次语义(幂等 + 事务)与消费者重平衡(rebalance)机制,它们在本章"不丢不重"的基础上进一步解决分布式一致性的最后一公里。
- Kafka 高吞吐三件套:顺序写磁盘(零随机 I/O)+ 零拷贝(sendfile)+ 批量发送(batch.size + linger.ms)。
- 可靠性核心:ISR(In-Sync Replica)选主 + acks=-1 + min.insync.replicas;三者缺一不可才能说"不丢消息"。
- 有序性边界:partition 内严格有序,跨分区无序;保序要"同一 key → 同一 partition"。
- pull vs push:Kafka 选 pull 是因为消费速率由消费者自主控制,push 模式下消费者来不及消费会导致消息堆积甚至丢失。
- 下一章我们会讲 MySQL 全体系——从索引原理到分库分表,它是企业级后端最常用的数据存储。
二十、自测题与动手练习
自测题(合上书能答出来,才算懂):
- Kafka 的"高吞吐"靠哪几项存储/发送设计支撑?为什么它比 RabbitMQ 更适合大数据日志场景?
- partition 内有序还是全局有序?如果想让"同一用户"的事件严格有序,Producer 该怎么发?
- 画出 Leader/Follower 与 ISR 的关系,并说明:ISR 里的副本挂了、ISR 外的 Follower 挂了,分别有什么影响?
acks=-1+min.insync.replicas=1能完全保证不丢消息吗?为什么?消费端还要做什么配合?- Consumer 是推还是拉?为什么 Kafka 选 pull 而不是 push?
动手练习(建议真做一遍):
- 在纸上画一张 Kafka 集群拓扑:1 个 topic、3 个 partition、副本因子 3,标出每个 partition 的 Leader 落在哪台 broker、Follower 在哪。
- 用你熟悉的语言写一段 Producer 代码,设置
acks=all、指定 key 并观察相同 key 是否进同一分区;再改成assign()手动消费某个分区。 - 对比"自动提交 offset"和"手动提交 + 处理完再 commit"两种写法,各写一句什么场景下会丢消息、什么场景下会重复消费。
精确一次(Exactly-once):消息既不丢失也不重复。需要额外开启:
① Producer 端:enable.idempotence=true(幂等 producer,防止发送端重复)
② Kafka 集群:transactional.id(事务支持)
③ Consumer 端:isolation.level=read_committed
工程实践:大多数场景用"至少一次"+“消费端幂等"就够了(数据库 UPSERT 或 Redis SETNX 防重)。精确一次开销大,只在金融级场景使用。面试加分点:提到幂等 producer 靠 producerId + epoch 实现,transaction 靠两阶段提交。