Kafka 分布式消息队列
27 道题- 分类
- 中间件
- 题目数
- 27 道
1 Kafka 是什么?适合解决哪些问题
答案:
Kafka 是一个分布式事件流平台。它将消息按主题(Topic)和分区(Partition)组织为可持久化、可顺序追加、可重复读取的日志;生产者(Producer)负责写入,消费者组(Consumer Group)按各自的消费位点(Offset)独立读取。它既能承担异步消息解耦,也常用于日志采集、事件总线、数据集成和流式处理。
核心特点:
- 持久化事件日志:消息在保留期内不会因某个消费者读完而删除,多个消费者组可以按各自进度重复消费、回溯消费。
- 分区并行:一个主题可拆为多个分区,读写可以分散到多个代理节点(Broker)和消费者实例上。
- 副本容错:每个分区可配置多个副本,并通过主副本(Leader)、从副本(Follower)和同步副本集合(ISR)保障可用性和数据可靠性。
- 生态集成:Kafka Connect 用于数据导入导出,Kafka Streams 用于应用内流式处理;它们建立在同一条事件日志之上。
| 适合的场景 | 原因 |
|---|---|
| 业务事件总线、服务解耦 | 生产者和消费者独立扩缩容,支持多订阅方 |
| 日志、指标、埋点汇聚 | 顺序写入和批量传输适合高吞吐数据流 |
| 变更数据捕获(CDC)、数据管道与跨系统同步 | 可持久保存事件,并支持按位点恢复和重放 |
| 实时分析与流式处理 | 可持续消费事件流并维护处理进度 |
| 不宜仅依赖 Kafka 的场景 | 原因 |
|---|---|
| 强同步远程过程调用(RPC) | 调用方需要即时结果,应使用 HTTP、gRPC 等同步协议 |
| 任意精度的延迟/定时任务 | Kafka 原生不是调度系统,需要业务侧或专用组件配合 |
| 需要跨分区全局严格顺序且吞吐很高 | 全局顺序通常意味着单分区,会限制并行能力 |
2 Kafka 的核心架构和关键概念是什么
答案:
Kafka 集群由代理节点(Broker)、客户端和控制平面组成。主题(Topic)是业务上的消息分类,分区(Partition)是实际的有序日志单元;生产者(Producer)将消息写入分区,消费者组(Consumer Group)以组为单位分摊读取,代理节点负责保存数据并复制副本。
| 概念 | 含义 | 面试中应关注的边界 |
|---|---|---|
| 代理节点(Broker) | 运行 Kafka 的服务节点,负责存储分区、处理读写请求 | 一个代理节点可承载多个分区的主副本或从副本 |
| 主题(Topic) | 消息的逻辑分类 | 主题本身不保证全局顺序,顺序由分区提供 |
| 分区(Partition) | 主题的物理分片,是追加写入的有序日志 | 一个消费者组内,同一分区同一时刻只分配给一个消费者 |
| 消息记录(Record)与位点(Offset) | 消息记录是一条消息;位点是它在某个分区中的单调递增位置 | 位点只在“主题-分区”范围内有意义,不是全局消息 ID |
| 副本(Replica) | 分区的副本集合,包含一个主副本和零个或多个从副本 | 客户端通常只向主副本写入,从副本从主副本拉取复制 |
| 生产者(Producer) | 消息生产者 | 消息键(Key)通常决定记录落到哪个分区 |
| 消费者组(Consumer Group) | 共同消费同一业务流的一组消费者 | 不同消费者组可独立消费同一主题;组内实例数超过分区数不会增加并行度 |
| 控制器(Controller)/ KRaft 仲裁组(Quorum) | 管理元数据、代理节点状态和分区主副本选举的控制平面 | 新集群应理解 KRaft;旧集群可能仍使用 ZooKeeper |
生产者(Producer)→ 主题 / 分区主副本(Leader)→ 从副本(Follower)
↓
消费者组(Consumer Group)A
消费者组(Consumer Group)B
3 Kafka 为什么能获得高吞吐量
答案:
Kafka 的高吞吐来自端到端的批处理和顺序输入/输出(I/O)设计,而不是单一参数。它将写入分散到分区,代理节点(Broker)以追加方式写日志,客户端和网络层尽可能成批传输数据。
- 分区并行:不同分区可以由不同代理节点、磁盘和消费者实例并行处理。
- 顺序追加与页缓存(Page Cache):代理节点将记录顺序追加到日志段,避免随机写;操作系统页缓存承担热数据缓存,读写路径通常不需要为每条消息单独落盘。
- 批量与压缩:生产者将多条记录合并为批次并可压缩,减少网络往返、协议开销和磁盘占用。
- 拉取式消费:消费者按批次拉取(Fetch),并可长轮询,避免代理节点为每条消息维护推送状态。
- 高效网络传输:在合适的读取路径上可利用零拷贝等机制,减少数据在用户态和内核态之间的额外复制。
高吞吐不等于低延迟或最高可靠性。更大的批次(batch)、更多压缩和更宽松的确认策略通常能提升吞吐,但会改变延迟、CPU 使用率或数据可靠性;需要结合业务的恢复点目标(RPO)、恢复时间目标(RTO)和延迟目标取舍。
4 Kafka 如何保证消息顺序性,边界是什么
答案:
Kafka 只保证单个分区(Partition)内按写入日志的位点(Offset)顺序读取,不保证一个主题(Topic)跨多个分区的全局顺序。要保证同一业务实体的顺序,应让它们使用稳定且相同的消息键(Key)路由到同一分区,例如同一 orderId 的订单事件始终使用 orderId 作为消息键。
实现要点:
- 对需要排序的业务键固定指定消息键,避免生产者随机分区。
- 在同一消费者组中,同一分区同一时刻只由一个消费者实例处理;应用内也不能再将该分区的记录无序并发执行。
- 生产者发生可重试错误时应使用幂等生产者,避免重试造成重复或乱序;跨多个生产者并发发送时,业务侧仍需定义事件版本或序列号。
- 扩容主题的分区数后,新写入记录的消息键哈希结果可能变化;若历史顺序必须连续,应提前设计分区策略或使用业务序列号校验。
常见边界:
| 目标 | 可行做法 | 代价 |
|---|---|---|
| 同一订单内有序 | 同一 orderId 路由到同一分区 | 消息键分布不均会形成热点分区 |
| 主题全局有序 | 单分区串行生产和消费 | 吞吐与消费者并行度受限 |
| 处理结果有序且不重复 | 分区串行处理 + 幂等业务写入/事务 | 实现复杂度和处理延迟上升 |
5 Kafka 主题(Topic)的分区与副本配置
答案:
主题的分区数决定最大并行度,副本因子决定容错边界;两者都应由吞吐、恢复目标和故障域数量推导,而不是套用固定数字。下面使用 Kafka 自带管理命令说明配置方式。
创建与配置示例:
bin/kafka-topics.sh --bootstrap-server broker-1:9092 \
--create --if-not-exists --topic orders \
--partitions 12 --replication-factor 3 \
--config min.insync.replicas=2 \
--config retention.ms=604800000 \
--config cleanup.policy=delete
bin/kafka-configs.sh --bootstrap-server broker-1:9092 \
--entity-type topics --entity-name orders --alter \
--add-config segment.bytes=1073741824,compression.type=producer
分区数设计原则:
| 输入 | 设计方法 | 常见误区 |
|---|---|---|
| 生产吞吐 | 用压测得到单分区稳定写入能力,再按峰值吞吐和冗余系数计算 | 只按日均流量估算,忽略峰值 |
| 消费并行度 | 一个消费者组内最多一个消费者处理一个分区 | 消费者实例多于分区不会提升并行度 |
| 热点风险 | 为高基数消息键预留分区与分片策略 | 仅增加分区,未解决热点消息键 |
| 运维成本 | 评估分区带来的文件、元数据、恢复和再均衡开销 | 以“大量分区总会更快”为前提 |
分区数约束:分区只能增加、不能减少;扩容后新消息键的路由可能变化,因此同一业务实体的历史顺序若需连续,必须提前设计分区策略。生产环境应预留增长空间,但同时控制单个代理节点(Broker)的分区和副本数量。
副本因子配置:
| 副本因子 | 配合 min.insync.replicas 的典型边界 | 适用场景 |
|---|---|---|
| 1 | 无副本容错 | 开发或可再生数据 |
| 2 | min.insync.replicas=1 时可用性较高但可靠性较弱 | 风险可接受的非关键数据 |
| 3 | min.insync.replicas=2 时可容忍一个副本失效仍安全写入 | 常见生产基线 |
| 更高 | 按故障域、恢复窗口和成本评估 | 跨机架或跨可用区的高可靠场景 |
容量估算:
保留数据量 ≈ 峰值写入字节/天 × 保留天数 ÷ 实测压缩比
物理存储量 ≈ 保留数据量 × 副本因子 × (1 + 索引、重分配和增长余量)
单 Broker 容量 ≈ 物理存储量 ÷ Broker 数,再预留副本分布不均衡空间
6 Kafka 的同步副本集合(In-Sync Replicas,ISR)机制与最小 ISR
答案:
同步副本集合(ISR)是 Kafka 高可用和数据一致性的核心机制,指与主副本(Leader)保持同步的从副本(Follower)集合。
ISR 判断条件:
Follower 进入 ISR 条件:
- Follower 持续从 Leader 拉取数据,并在 replica.lag.time.max.ms 时间窗口内保持可追赶状态
- Kafka 以时间窗口判断副本是否落后;不要把早期版本的 replica.lag.max.messages 当作当前集群的调优依据
Follower 踢出 ISR 条件:
- 超过 replica.lag.time.max.ms 未能保持同步,Leader 将其移出 ISR
ISR 收缩与扩展流程:
Leader 维护 ISR 集合 →
Follower 持续 Fetch 并追赶日志末端位点(LEO)→
超过 replica.lag.time.max.ms 仍无法保持同步 →
从 ISR 移除并提交分区状态 →
Controller 将状态变更传播给集群 →
Follower 追平后重新加入 ISR
min.insync.replicas 配置:
| minISR | 容错含义 | 风险 |
|---|---|---|
| 1 | Leader 单独在 ISR 即可写入 | 数据丢失,Leader 故障后数据不可恢复 |
| 2(推荐) | 至少 Leader + 1 个 Follower 在 ISR 才能写入 | 可容忍 1 个 ISR 节点故障 |
| replicas = 3, minISR = 3 | 所有副本均需在 ISR | 任一节点故障中断写入 |
写入一致性保障:
Producer acks=all + min.insync.replicas=2 + replicationFactor=3
→ 消息写入 Leader → Leader 等待 Follower 确认 →
至少 2 个副本(含 Leader)持久化成功 → Producer 收到 ACK
ISR 抖动处理:broker 短暂 GC pause 或网络抖动导致 Follower 频繁进出 ISR,表现为 UnderMinIsr 指标波动。通过调大 replica.lag.time.max.ms(如 60s)和确保 Broker 有充足的内存避免 GC 过频来缓解。
7 Kafka 的主副本(Leader)选举与控制器(Controller)
答案:
控制器(Controller)负责分区和副本状态的控制平面。KRaft 集群由控制器仲裁组(Controller Quorum)复制元数据并选举活跃控制器;ZooKeeper 模式属于旧架构,面试时应能说明两种模式的边界,而不要把 ZooKeeper 路径当成 KRaft 的运行机制。
控制器职责:
| 职责 | 说明 |
|---|---|
| 分区主副本选举 | 代理节点宕机时为受影响分区选举新的主副本 |
| ISR 管理 | 维护每个分区的 ISR 集合变更 |
| 主题管理 | 主题创建、删除、分区扩展 |
| 副本分配 | 新分区创建时的副本分布策略 |
| 元数据广播 | 向所有代理节点广播集群元数据变更 |
| 首选主副本选举 | 触发和完成首选副本主副本选举(Preferred Replica Leader Election) |
控制器选举的核心过程:
KRaft Controller Quorum 复制元数据日志 →
多数派选出活跃 Controller →
活跃 Controller 处理 Broker 注册、分区状态和 Leader 变更 →
活跃 Controller 故障 → 多数派重新选举 → 新 Controller 从元数据日志继续处理
旧 ZooKeeper 集群:Broker 通过临时节点竞争 Controller,
Controller 切换依赖 ZooKeeper 会话与 Watch 通知。
控制器分区主副本选举策略:
| 策略 | 行为 |
|---|---|
| OfflinePartitionLeaderElectionStrategy | 仅当没有 Leader 时选举 |
| PreferedReplicaLeaderElectionStrategy | 优先选举 Preferred Replica 为 Leader |
| NoOpLeaderElectionStrategy | 不执行自动选举 |
Broker 不可用或受控下线 → Controller 收到状态变更 →
遍历受影响的 Leader 分区 →
从可用 ISR 中选择新 Leader →
提交新的分区状态与 Leader Epoch →
向 Broker 和客户端发布元数据更新 →
副本继续追赶并恢复正常 ISR
控制器故障影响与优化:
- 控制器切换期间(数秒内),分区主副本选举、主题创建/删除等管理操作会中断。
- 正常下线应使用受控下线(controlled shutdown),让 Controller 先完成受影响分区的 Leader 切换;副本迁移仍需单独执行重分配。
- KRaft 模式下,控制器与代理节点角色可分离,控制器可独立组成仲裁组(Quorum)。
8 Kafka 的恰好一次(Exactly-Once)语义与事务
答案:
Kafka 的恰好一次语义(Exactly-Once Semantics,EOS)通过幂等生产者和事务机制,在 Kafka 的读—处理—写链路及消费位点提交范围内 避免重复可见记录。它不自动覆盖数据库、HTTP 调用、邮件等外部副作用;这些系统仍需幂等键、去重表或事务外盒(Outbox)等业务设计。
三种消息传递语义:
| 语义 | 机制 | 实现 |
|---|---|---|
| 最多一次(At-Most-Once) | acks=0,发送即忘 | 可能丢消息 |
| 至少一次(At-Least-Once) | acks=1/all,重试机制 | 可能重复 |
| 恰好一次(Exactly-Once) | 幂等 + 事务 | 不丢不重复 |
幂等生产者(Idempotent Producer):
Producer 初始化 → 向 Broker 请求 Producer ID(PID)→
Broker 分配 PID + Epoch → Producer 每条消息附带:
(PID, Epoch, SequenceNumber) →
Broker 按 (PID, Partition) 维护已提交的最大 SequenceNumber →
收到消息 → 检查:
Sequence == lastSeq+1 → 接受
Sequence <= lastSeq → 重复,丢弃,返回 ACK
Sequence > lastSeq+1 → 乱序,抛出 OutOfOrderSequenceException
# 幂等生产者配置
enable.idempotence=true
acks=all
max.in.flight.requests.per.connection=5
retries=2147483647
事务(Transactions):
Producer 配置稳定的 transactional.id → initTransactions() →
事务协调器分配或恢复 Producer ID(PID)与 Epoch →
beginTransaction() →
发送消息(关联事务)→
消息写入分区但不向消费者可见 →
sendOffsetsToTransaction() →
commitTransaction() / abortTransaction() →
Broker 写入 Commit/Abort Marker →
消费者读取到 Marker 后过滤事务消息
EOS 消费者配合:下游消费者需设置 isolation.level=read_committed,才能过滤未提交和已中止的事务消息;否则仍可能读到事务中的未提交记录。
EOS 性能代价:事务会增加协调、内部主题写入和批次提交开销,实际延迟与吞吐损失取决于事务大小、提交频率和副本配置,必须用目标负载压测后取舍。
9 Kafka 的消息压缩(Snappy / LZ4 / GZip / Zstd)
答案:
Kafka 支持由生产者(Producer)压缩、代理节点(Broker)原样存储的压缩策略,可减少网络传输和磁盘占用。
压缩算法对比:
| 算法 | 压缩率 | 压缩速度 | 解压速度 | 适用场景 |
|---|---|---|---|---|
| Snappy | 中等(2x) | 极快 | 极快 | 低延迟场景,平衡选择 |
| LZ4 | 中等(2x) | 最快 | 最快 | 低延迟、高吞吐场景首选 |
| GZip | 高(4x-5x) | 慢 | 中等 | 存储密集型、冷数据归档 |
| Zstd | 最高(4x-6x) | 中等 | 快 | 高压缩率需求,推荐生产使用 |
| uncompressed | 1x | 无开销 | 无开销 | 实时性要求极高、消息本身已压缩 |
压缩配置层级:
# Producer 端压缩(推荐)
compression.type=zstd # none / gzip / snappy / lz4 / zstd
# Broker 或 Topic 端压缩策略
compression.type=producer # producer(保留原始压缩) / uncompressed(强制解压)
# 按主题覆盖 Broker 默认值;只影响该主题后续写入的批次
bin/kafka-configs.sh --bootstrap-server broker-1:9092 \
--entity-type topics --entity-name orders --alter \
--add-config compression.type=producer
压缩在 Broker 端的处理:
Producer 压缩消息 → 封装为 MessageSet → 发送到 Broker →
Broker.compression.type=producer → 原样写入 Page Cache →
批量刷盘到日志段文件 →
Consumer Fetch 请求 → Broker 直接发送压缩数据 →
Consumer 解压消息集
选择原则:先以真实消息样本压测压缩比、端到端延迟和 CPU;消息已由图片、视频或加密载荷压缩时,重复压缩通常收益很低。Zstd 常能在压缩率与解压开销之间取得较好平衡,但不是所有低延迟链路的默认答案。
10 Kafka 的日志保留策略(按时间、大小与日志压缩)
答案:
Kafka 的日志保留策略通过清理策略及配套参数,控制消息在代理节点(Broker)上的生命周期。
三种清理策略:
| 策略 | cleanup.policy | 触发条件 | 适用场景 |
|---|---|---|---|
| 按时间删除 | delete | retention.ms / retention.minutes / retention.hours | 日志、事件流、时序数据 |
| 按大小删除 | delete | retention.bytes | 存储空间受限的场景 |
| 日志压缩 | compact | min.cleanable.dirty.ratio | 键值(KV)状态存储、变更数据捕获(CDC)、配置快照 |
主题级保留策略示例:
# 事件流:时间和大小任一条件先满足即触发删除
bin/kafka-configs.sh --bootstrap-server broker-1:9092 \
--entity-type topics --entity-name events --alter \
--add-config cleanup.policy=delete,retention.ms=604800000,retention.bytes=107374182400
# 状态主题:保留每个 Key 的最新值;删除标记需保留足够久以便慢消费者看到
bin/kafka-configs.sh --bootstrap-server broker-1:9092 \
--entity-type topics --entity-name user-state --alter \
--add-config cleanup.policy=compact,min.cleanable.dirty.ratio=0.5,delete.retention.ms=86400000
日志清理时机与机制:
日志清理线程定期扫描 → 选择脏比率较高的日志段 →
逐段读取 → 构建位点映射(每个 Key 的最新 Offset)→
保留最新 Offset 的消息,删除旧消息 → 合并为干净的新段 →
交换新旧段 → 更新索引
保留策略参数优先级:
| 优先级 | 规则 |
|---|---|
| 最高 | 已确定为活跃日志段(Active Segment)的日志段不删除 |
| 高 | min.compaction.lag.ms / max.compaction.lag.ms 内的段不参与压缩 |
| 中 | retention.ms 和 retention.bytes 同时设置时先触发的优先生效 |
| 低 | 代理节点级别的保留策略作为默认值,主题级别配置会覆盖它 |
生产环境保留策略最佳实践:
| 数据类型 | 推荐策略 | 保留周期 | 原因 |
|---|---|---|---|
| 业务事件流 | delete + 按时间清理 | 7-30 天 | 事件有时效性 |
| 状态快照(CDC) | compact | 永久 | 仅需每个 Key 的最新值 |
| 告警事件 | delete + 按大小清理 | 100GB | 存储空间受限 |
| 审计日志 | delete + time | 90 天+ | 合规要求 |
11 Kafka 的消费者组(Consumer Group)再均衡(Rebalance)机制
答案:
消费者组再均衡(Rebalance)是 Kafka 在消费者成员变化时重新分配分区所有权的协调过程。
触发再均衡(Rebalance)的条件:
| 触发条件 | 说明 |
|---|---|
| 消费者加入或离开消费者组 | 消费者启动、正常退出或崩溃(会话超时) |
| 主题分区数增加 | 新增分区需分配给消费者 |
| 订阅主题变更 | 消费者通过正则订阅匹配到新主题 |
max.poll.interval.ms 超时 | 消费者两次 poll 间隔超过阈值 |
再均衡协议阶段:
Phase 1: JoinGroup
Consumer → Group Coordinator: JoinGroupRequest(memberId, protocols)
Coordinator 选第一个 Consumer 为 Group Leader
Coordinator → Group Leader: JoinGroupResponse(所有 Consumer 信息)
Coordinator → 其他 Consumer: JoinGroupResponse(成员 ID)
Phase 2: SyncGroup
Group Leader 执行分配策略 → 生成 Partition Assignments →
Group Leader → Coordinator: SyncGroupRequest(分配方案)
Coordinator → 所有 Consumer: SyncGroupResponse(各 Consumer 的分区列表)
Consumer 收到分配方案 → 清空或保持未提交的 Offset → 开始消费
再均衡策略对比:
| 策略 | 算法 | 适用场景 |
|---|---|---|
| RangeAssignor | 按主题分区范围连续分配 | 默认策略,平衡性一般 |
| RoundRobinAssignor | 所有分区轮询均匀分配 | 消费者订阅的主题完全相同 |
| StickyAssignor | 尽量保持原有分配,减少分区迁移 | 生产推荐,减少无意义迁移 |
| CooperativeStickyAssignor | 增量式再均衡,不触发全局暂停(Stop-The-World) | 生产首选(Kafka 2.4+) |
协作式再均衡(Cooperative Rebalance):
与传统急切式再均衡(Eager Rebalance)的差异:
- 不触发全局暂停(Stop-The-World)
- 消费者不必在再均衡期间暂停所有分区消费
- 仅迁移需要重新分配的分区
- 未被重新分配的分区继续正常消费
Consumer 配置:
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor
再均衡优化措施:
# 延长 Session 超时,避免短暂网络波动触发 Rebalance
session.timeout.ms=30000 # Consumer 心跳超时(默认 45s)
# 延长心跳间隔,减少心跳开销
heartbeat.interval.ms=3000 # Consumer 向 Coordinator 发送心跳间隔
# 增加 Poll 间隔上限,防止慢速处理触发 Rebalance
max.poll.interval.ms=600000 # 两次 poll 最大间隔(默认 5min)
# 使用 Cooperative Sticky Rebalance
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor
# 增加 poll 记录数,减少 poll 调用频率
max.poll.records=500
12 Kafka 的消费者组协调器(Group Coordinator)与位点(Offset)管理
答案:
消费者组协调器(Group Coordinator)是代理节点端管理消费者组的核心组件,负责消费者成员管理、位点提交和再均衡协调。
协调器定位机制:
Consumer Group ID →
hash(groupId) % __consumer_offsets 分区数 → 目标分区 →
该分区 Leader 所在 Broker = Group Coordinator
位点存储架构(Kafka 0.9+):
旧架构(ZooKeeper 存储 Offset):
Consumer → ZK 直接写入 Offset
缺点:ZK 不适合高频写入,OOM 风险
新架构(__consumer_offsets Topic 存储 Offset):
Consumer → Group Coordinator → __consumer_offsets Topic
存储格式:Key = <GroupId, Topic, Partition>
Value = <Offset, Metadata, Timestamp>
压缩策略:compact(每个分区仅保留最新 Offset)
位点提交模式:
| 模式 | 配置 | 行为 | 风险 |
|---|---|---|---|
| 自动提交 | enable.auto.commit=true | 按 auto.commit.interval.ms 周期提交 | Consumer 故障后可能重复消费 |
| 同步手动提交 | consumer.commitSync() | 阻塞等待提交完成 | 降低吞吐 |
| 异步手动提交 | consumer.commitAsync() | 非阻塞提交,回调处理结果 | 网络异常时可能丢提交 |
| 事务内提交 | EOS 模式 | Offset 与消息在同一事务中原子提交 | 仅事务 Producer+Consumer 场景 |
协调器故障处理:
Consumer 心跳超时 → Group Coordinator 移除该 Consumer →
触发 Rebalance → 但 Coordinator 自身可能也在故障 →
Consumer 收到 NOT_COORDINATOR 错误 →
重新 FindCoordinator Request → 定位新 Coordinator →
重新 JoinGroup
消费者位点重置策略:
| 策略 | 行为 |
|---|---|
latest(默认) | 从最新位置开始(忽略历史消息) |
earliest | 从最早可用偏移量开始(重消费历史) |
none | 未找到 Offset 时抛出异常 |
KRaft 与消费者组的关系:KRaft Controller 负责集群元数据,不接管普通消费者组协调。消费者组协调仍由 Broker 端组件处理,消费位点仍存储在 __consumer_offsets 内部 Topic 中。
13 Kafka 的监控指标(JMX Exporter、Prometheus 与 Grafana)
答案:
Kafka 监控应从 JMX、客户端指标和操作系统指标三条链路采集数据。常见架构是:Kafka JVM 暴露 JMX → JMX Exporter 转为 Prometheus 指标 → Prometheus 抓取 → Grafana 展示、Alertmanager 告警。采集器的部署方式可以不同,但指标口径应保持一致。
JMX Exporter 启动示例:
export KAFKA_OPTS='-javaagent:/opt/jmx/jmx_prometheus_javaagent.jar=9404:/opt/jmx/kafka-metrics.yml'
bin/kafka-server-start.sh config/server.properties
JMX Exporter 配置(kafka-metrics-config.yml):
lowercaseOutputName: true
lowercaseOutputLabelNames: true
rules:
# Broker 核心指标
- pattern: kafka.server<type=BrokerTopicMetrics, name=(BytesInPerSec|BytesOutPerSec)><>OneMinuteRate
name: kafka_server_broker_topic_metrics_$1
type: GAUGE
- pattern: kafka.server<type=BrokerTopicMetrics, name=(MessagesInPerSec)><>OneMinuteRate
name: kafka_server_broker_topic_metrics_$1
type: GAUGE
# 网络请求
- pattern: kafka.network<type=RequestMetrics, name=(TotalTimeMs|RequestQueueTimeMs), request=(Produce|Fetch|FetchConsumer|FetchFollower)><>Count
name: kafka_network_request_metrics_$2_$1
labels:
request: "$3"
# ISR 与副本
- pattern: kafka.server<type=ReplicaManager, name=(UnderReplicatedPartitions|UnderMinIsrPartitionCount)><>Value
name: kafka_server_replica_manager_$1
- pattern: kafka.server<type=ReplicaManager, name=(IsrShrinksPerSec|IsrExpandsPerSec)><>OneMinuteRate
name: kafka_server_replica_manager_$1
# Producer 请求延迟
- pattern: kafka.network<type=RequestMetrics, name=TotalTimeMs, request=(Produce|FetchConsumer)><>(\w+)
name: kafka_network_$3_$1_$2
labels:
request: "$3"
Prometheus 抓取示例:
scrape_configs:
- job_name: kafka
static_configs:
- targets:
- broker-1.example.com:9404
- broker-2.example.com:9404
- broker-3.example.com:9404
关键监控指标分类与告警阈值:
| 类别 | 关键指标 | 告警阈值 |
|---|---|---|
| 吞吐 | kafka_server_broker_topic_metrics_bytesinpersec | > 80% 磁盘带宽上限 |
| 延迟 | kafka_network_request_metrics_totaltimems_produce | P99 > 100ms |
| ISR | kafka_server_replica_manager_underreplicatedpartitions | > 0 持续 5 分钟 |
| 消费滞后 | kafka_consumergroup_group_lag | > 10000 或增长率 > 1000/min |
| 磁盘使用 | kafka_log_log_size | > 80% 磁盘容量 |
| Controller | kafka_controller_controllerstate | 不等于 1(Broker 不是 Controller 但值 = 0 OK) |
| 网络 | kafka_network_processor_idle_percent | < 0.3 持续 10 分钟 |
| GC | jvm_gc_pause_seconds | P99 > 500ms |
Grafana Dashboard 分层:
- 总览仪表板(Overview Dashboard):集群吞吐、延迟、活跃控制器、代理节点数量
- 代理节点仪表板(Broker Dashboard):单个代理节点的请求速率、网络 I/O、磁盘 I/O、GC 行为
- 主题/分区仪表板(Topic/Partition Dashboard):分区级别读写速率、大小、ISR 状态
- 消费者组仪表板(Consumer Group Dashboard):消费滞后、提交速率、再均衡事件
14 Kafka 的性能调优(生产者、消费者与代理节点参数)
答案:
Kafka 性能调优需从生产者(Producer)、代理节点(Broker)和消费者(Consumer)三个维度分别进行。
生产者(Producer)端调优:
# 吞吐优先
batch.size=131072 # 批量大小 128KB
linger.ms=5 # 等待 5ms 攒批
compression.type=zstd # 压缩算法
buffer.memory=134217728 # 发送缓冲区 128MB
max.in.flight.requests.per.connection=5
acks=1 # Leader 确认即可
# 可靠性优先
acks=all # 所有 ISR 确认
enable.idempotence=true # 幂等生产者
max.in.flight.requests.per.connection=5
retries=2147483647 # 无限重试
delivery.timeout.ms=120000 # 超时 2 分钟
代理节点(Broker)端调优:
# server.properties;数值必须根据 CPU、磁盘、网络与压测结果确定
num.network.threads=8
num.io.threads=16
socket.send.buffer.bytes=1048576
socket.receive.buffer.bytes=1048576
socket.request.max.bytes=104857600
# 日志段与副本拉取
log.segment.bytes=1073741824
log.index.interval.bytes=4096
num.replica.fetchers=4
replica.fetch.max.bytes=10485760
compression.type=producer
# Leader 均衡应避免与大规模重分配同时进行
auto.leader.rebalance.enable=true
leader.imbalance.per.broker.percentage=10
leader.imbalance.check.interval.seconds=300
JVM 堆通常不应吞掉主机内存:Kafka 依赖操作系统页缓存处理热数据。设置 -Xms 与 -Xmx 后,还应为页缓存、直接内存、网络缓冲和操作系统预留空间;不要把示例数值当成所有机器的固定配方。
消费者(Consumer)端调优:
# 吞吐优先
fetch.min.bytes=1048576 # 每次 Fetch 最少拉取 1MB
fetch.max.wait.ms=500 # 等待 500ms 攒批
max.partition.fetch.bytes=10485760 # 每分区最大 10MB
max.poll.records=500 # 每次 Poll 最多 500 条
# 延迟优先
fetch.min.bytes=1 # 有数据就返回
fetch.max.wait.ms=50 # 最多等待 50ms
# Rebalance 优化
session.timeout.ms=30000
heartbeat.interval.ms=3000
max.poll.interval.ms=600000
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor
# Offset 管理
enable.auto.commit=false # 手动提交,避免自动提交带来的重复消费
JVM 与操作系统(OS)调优:
| 调优项 | 配置 | 说明 |
|---|---|---|
| Page Cache | vm.dirty_background_ratio=5 | 后台刷脏页阈值 |
| 文件描述符 | ulimit -n 100000 | Kafka 大量文件需求 |
| Swap | vm.swappiness=1 | 最小化 Swap |
| 磁盘调度 | noop / none | SSD 无 I/O 调度器 |
| 网络调优 | net.core.rmem_max=134217728 | 套接字接收缓冲区 |
15 Kafka 的存储规划:日志目录、JBOD 与分层存储(Tiered Storage)
答案:
Kafka 的本地日志由代理节点(Broker)的 log.dirs 管理。先根据保留量、副本因子、恢复窗口和磁盘带宽做容量规划,再决定单目录、多目录(JBOD)或分层存储;底层卷和磁盘的实现方式不改变这些 Kafka 容量约束。
多日志目录(JBOD)示例:
# server.properties
log.dirs=/data-1/kafka-logs,/data-2/kafka-logs
num.recovery.threads.per.data.dir=2
log.segment.bytes=1073741824
| 维度 | 单日志目录 | 多日志目录(JBOD) |
|---|---|---|
| I/O 并行度 | 受单块设备限制 | 可利用多块独立设备 |
| 扩容方式 | 扩容原设备或迁移数据 | 新增目录后再做分区重分配 |
| 故障影响 | 该目录上的副本不可用 | 仅受故障目录影响的副本不可用,仍依赖副本因子保障可用性 |
| 运维复杂度 | 较低 | 需监控每个目录的容量、延迟和故障状态 |
分层存储的关键边界:
活跃日志段:保留在本地磁盘,服务低延迟尾部读取
已关闭日志段:由 RemoteStorageManager 上传到远端存储
本地保留到期:本地副本可删除,历史读取转为远端读取
主题总保留到期:按主题策略删除本地与远端历史数据
Kafka 原生分层存储需要提供兼容的 RemoteStorageManager 实现;Kafka 不附带可直接用于所有对象存储的远端存储插件。Broker 侧启用 remote.log.storage.system.enable=true 后,还要为目标主题配置 remote.storage.enable=true,并分别设置 local.retention.* 与主题总保留策略。
bin/kafka-configs.sh --bootstrap-server broker-1:9092 \
--entity-type topics --entity-name audit-events --alter \
--add-config remote.storage.enable=true,local.retention.ms=86400000,retention.ms=2592000000
- 优势:降低本地热存储成本,保留更长历史并支持回溯。
- 代价:远端读取延迟、对象存储一致性与插件运维都会进入故障域;必须压测回溯和副本恢复场景。
- 限制:启用前确认 Kafka 版本与插件能力;例如压缩主题等场景存在版本和实现限制,不能直接套用到所有主题。
16 Kafka 的 SSL/TLS 加密与 SASL 认证
答案:
Kafka 的通信安全要分别处理传输加密、客户端认证、代理节点间认证和授权。生产环境常用 SASL_SSL:TLS 提供加密与服务端身份校验,SASL 提供客户端身份认证,ACL 或外部授权器决定该身份能做什么。
安全机制分工:
| 层次 | 常用机制 | 关注点 |
|---|---|---|
| 传输加密 | TLS | 证书链、主机名校验、协议与密码套件 |
| 客户端认证 | SCRAM、mTLS、OAuth | 身份生命周期与凭证轮换 |
| Broker 间认证 | TLS 或 SASL_SSL | 复制流量同样不能使用明文旁路 |
| 授权 | Kafka ACL 或外部授权器 | 最小权限和审计 |
Broker 监听器示例:
# server.properties
listeners=SASL_SSL://0.0.0.0:9093
advertised.listeners=SASL_SSL://broker-1.example.com:9093
listener.security.protocol.map=SASL_SSL:SASL_SSL
inter.broker.listener.name=SASL_SSL
ssl.keystore.location=/etc/kafka/secrets/broker.keystore.p12
ssl.keystore.password=<keystore-password>
ssl.key.password=<key-password>
ssl.truststore.location=/etc/kafka/secrets/broker.truststore.p12
ssl.truststore.password=<truststore-password>
ssl.client.auth=none
sasl.enabled.mechanisms=SCRAM-SHA-512
listener.name.sasl_ssl.scram-sha-512.sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required;
SASL 认证机制对比:
| 认证类型 | 身份来源 | 适用场景 | 注意事项 |
|---|---|---|---|
| SCRAM-SHA-512 | Kafka 用户凭证 | 服务账号、常规客户端 | 必须与 TLS 配合,避免泄露密码 |
| TLS 客户端证书 | 企业 CA / PKI | 高安全服务间调用 | 需设计签发、吊销和轮换流程 |
| OAuth 2.0 | 身份提供方 | 人员或统一身份体系 | 核对令牌刷新、JWKS 可用性与授权映射 |
| PLAIN | 外部认证服务或静态凭证 | 受控的兼容场景 | 不应在无 TLS 的网络上传输 |
客户端连接配置(SCRAM-SHA-512):
bootstrap.servers=broker-1.example.com:9093,broker-2.example.com:9093
security.protocol=SASL_SSL
sasl.mechanism=SCRAM-SHA-512
sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required \
username="my-user" \
password="<password-from-secret-store>";
ssl.truststore.location=/path/to/ca.crt
ssl.truststore.type=PEM
证书与凭证运维:私钥和密码不得写入 server.properties 或代码仓库;使用密钥管理系统或受控挂载注入。轮换前验证客户端信任链、Broker 重载或滚动重启策略,以及过期告警。advertised.listeners 中的主机名必须包含在服务端证书 SAN 中。
监听器安全策略矩阵:
plain (9092): No TLS + No Auth → 仅测试环境
tls (9093): TLS + No Auth → 加密但不认证
tls (9093): TLS + SASL → 加密 + 用户认证(生产推荐)
external: 仍使用 TLS + SASL,并让 advertised.listeners 对外可解析
17 Kafka 的访问控制列表(ACL)与基于角色的访问控制(RBAC)
答案:
Apache Kafka 原生授权模型是访问控制列表(ACL)。基于角色的访问控制(RBAC)通常由发行版、外部授权器或企业身份平台提供;面试中要明确区分“Kafka 本身的 ACL”与“平台在 ACL 之上封装的 RBAC”,不能把两者混为同一能力。
按服务身份授权示例:
# 生产者:只写 orders- 前缀主题
bin/kafka-acls.sh --bootstrap-server broker-1:9093 --command-config admin.properties \
--add --allow-principal User:order-service \
--operation Write --operation Describe \
--topic orders- --resource-pattern-type prefixed
# 消费者:读取订单主题,并只访问自己的消费者组
bin/kafka-acls.sh --bootstrap-server broker-1:9093 --command-config admin.properties \
--add --allow-principal User:order-service \
--operation Read --topic orders- --resource-pattern-type prefixed \
--group order-processor
# 事务生产者:允许使用自己的 transactional.id 前缀
bin/kafka-acls.sh --bootstrap-server broker-1:9093 --command-config admin.properties \
--add --allow-principal User:order-service \
--operation Write --transactional-id order-tx- --resource-pattern-type prefixed
ACL 资源类型与操作映射:
| 资源类型 | 常用操作 | 说明 |
|---|---|---|
topic | Read / Write / Describe / DescribeConfigs / Alter / AlterConfigs / Create / Delete | 主题(Topic)级别控制 |
group | Read / Describe | 消费者组(Consumer Group)控制 |
cluster | Describe / DescribeConfigs / Alter / IdempotentWrite / Create | 集群级别控制 |
transactionalId | Write / Describe | 事务生产者(Producer)控制 |
ACL 模式(Pattern)匹配:
| PatternType | 示例 | 匹配 |
|---|---|---|
literal | orders-topic | 精确匹配 orders-topic |
prefixed | orders- | 匹配 orders-topic、orders-dlq 等 |
超级用户(Super User)配置:超级用户应仅限于受控的集群运维身份,不应将普通应用服务加入。示例:
# server.properties
authorizer.class.name=org.apache.kafka.metadata.authorizer.StandardAuthorizer
super.users=User:kafka-admin
ACL 管理技巧:
- 每个服务使用独立认证身份,按主题、消费者组和事务 ID 最小授权。
- 对命名前缀使用
prefixed规则时,先设计好命名边界,避免一个宽前缀覆盖无关业务。 --allow-host适合边界明确的静态网络;在 NAT、弹性伸缩环境中,身份认证与网络层隔离通常更可靠。- 用
--list、审计日志和定期权限评审验证授权;管理员权限与应用权限必须分离。
18 Kafka 的 Cruise Control 自动均衡
答案:
Cruise Control 是 Kafka 生态中的独立容量与均衡服务,不是 Kafka Broker 自带的组件。它持续采集集群指标,按容量模型和优化目标计算副本、Leader 的迁移方案;确认后再通过 Kafka 的分区重分配能力执行。
它解决的不是“所有 Broker 的 CPU 必须一样高”,而是在不违反硬约束的前提下,把真正造成瓶颈的磁盘、网络、Leader 流量和副本数量拉回可接受范围。
| 目标类型 | 常见目标 | 含义 |
|---|---|---|
| 硬约束 | RackAwareGoal、容量上限 | 不满足就不应生成或执行方案,例如同一分区副本不能落在同一机架 |
| 优先目标 | DiskCapacityGoal、ReplicaCapacityGoal、Network*CapacityGoal | 避免磁盘、网络或副本数量超过节点能力 |
| 均衡目标 | LeaderBytesInDistributionGoal、ReplicaDistributionGoal | 在满足硬约束后,改善 Leader 写入流量和副本分布 |
典型操作流程:
采集指标与节点容量 → 选择目标集 → 生成 dry-run 方案 →
人工检查迁移量、硬约束和风险 → 设置限速并执行 →
观察副本同步、磁盘与客户端延迟 → 验证完成
Cruise Control 通常通过 REST API 提供集群状态和重平衡请求。不同发行版的路径、认证方式和参数会有差异,生产环境应先调用对应版本的 dry-run 接口,而不是直接执行迁移。概念上可按下面的节奏操作:
# 先查看集群状态与容量模型;实际 URL、认证和目标名以部署版本为准
curl -sS 'https://cruise-control.example.com/kafkacruisecontrol/kafka_cluster_state'
# 生成方案,只查看结果,不执行
curl -sS -X POST \
'https://cruise-control.example.com/kafkacruisecontrol/rebalance?dryrun=true'
落地时要关注:
- 指标窗口要覆盖业务高峰,否则方案可能只是在优化低负载时段。
- 容量模型必须反映真实磁盘、网卡和副本上限;错误的容量数据会让优化器给出危险方案。
- 对每次迁移设置带宽限速,优先观察
UnderReplicatedPartitions、复制延迟、磁盘水位和端到端延迟。 - 异常检测和自动修复应先在演练环境验证。节点故障、磁盘趋满时自动迁移,可能与故障恢复争抢资源。
19 Kafka 的备份与灾难恢复
答案:
Kafka 的副本机制主要应对单 Broker、磁盘或机架故障,不能代替备份:误删 Topic、错误缩短保留期、错误写入和权限配置失误都可能被复制到所有副本。灾备方案要分别覆盖数据、控制面配置和客户端切换。
| 保护对象 | 常用手段 | 恢复时的注意点 |
|---|---|---|
| 消息数据 | 跨集群复制、逻辑导出到独立存储、保留期内重放 | 明确可接受的数据窗口;异步复制存在复制延迟 |
| Topic、ACL、配额和 Broker 配置 | 版本化配置、变更记录、定期导出 | 先在目标集群重建配置,再恢复或导入数据 |
| 消费进度 | Consumer Group offset 导出或复制工具的 checkpoint | offset 只对同一份或可映射的数据有效,不能脱离数据直接恢复 |
| 证书、密钥与客户端配置 | 受控的密钥管理和独立备份 | 恢复时既要保证可用,也要避免旧凭据继续有效 |
两地灾备的常见形态:
主集群 ──异步复制──> 备用集群
↑ ↑
生产者与消费者 灾难时切换客户端入口
跨集群复制能缩短恢复时间,但默认不等于“零丢失、零重复”。切换前应记录源端最后可确认的位点和复制延迟;切换后要验证 Topic 配置、ACL、消费者位点映射以及下游幂等性。若主集群已经不可达,则以预先约定的 RPO 为准,而不是假设复制延迟一定为零。
恢复演练建议:
- 先定义故障范围、RPO、RTO 和业务优先级,区分“恢复读历史数据”和“恢复实时写入”。
- 在隔离环境重建目标 Kafka 集群及安全、Topic、ACL、配额等配置。
- 用复制链路、逻辑导出或业务重放恢复数据;随机抽取分区核对记录数量、顺序约束和关键业务字段。
- 让一组真实客户端完成鉴权、生产、消费和提交位点,再切换 DNS、服务发现或客户端 bootstrap 地址。
- 记录实际 RPO/RTO 与手工步骤,下一次按演练结果调整预案。
不要把 Broker 数据目录复制到另一套集群后直接启动,作为通用恢复方案。数据目录与 Broker ID、集群 ID、元数据版本和控制器元数据存在强关联;尤其在 KRaft 集群中,错误恢复可能造成元数据不一致。Tiered Storage 也不等于备份:它依赖远端存储插件和保留策略,若与主集群共用账户、权限或生命周期规则,仍可能同时失效。
20 Kafka 的 KRaft(无 ZooKeeper)架构
答案:
KRaft(Kafka Raft)用 Kafka 自己的 Raft 仲裁组保存和复制集群元数据,取代原先对 ZooKeeper 的依赖。它改变的是控制面,不改变普通消息仍由 Broker 日志副本保存、消费者仍由 Broker 协调的事实。
| 维度 | ZooKeeper 模式 | KRaft 模式 |
|---|---|---|
| 元数据存储 | 外部 ZooKeeper 集群 | Controller quorum 维护的元数据日志 |
| 控制器 | Broker 中选出一个活动控制器 | Controller quorum 选出活动控制器,可与 Broker 分离 |
| 运维对象 | Kafka 和 ZooKeeper 两套集群 | Kafka 的 Broker、Controller 角色及其元数据盘 |
| 一致性方式 | ZooKeeper 的读写与 Watch | Raft 日志复制和多数派确认 |
小型环境可以让同一节点同时承担 broker 和 controller 角色;生产集群通常将 Controller 与承载业务流量的 Broker 分开,避免大量副本迁移、长时间 GC 或磁盘抖动影响元数据仲裁。
角色分离的配置示意:
# Controller 节点:省略 TLS 与鉴权配置,只展示角色和仲裁关系
process.roles=controller
node.id=1
controller.quorum.voters=1@controller-1.example.com:9093,2@controller-2.example.com:9093,3@controller-3.example.com:9093
controller.listener.names=CONTROLLER
listeners=CONTROLLER://0.0.0.0:9093
inter.broker.listener.name=INTERNAL
listener.security.protocol.map=CONTROLLER:SSL,INTERNAL:SSL
metadata.log.dir=/var/lib/kafka/metadata
# Broker 节点
process.roles=broker
node.id=101
controller.quorum.voters=1@controller-1.example.com:9093,2@controller-2.example.com:9093,3@controller-3.example.com:9093
controller.listener.names=CONTROLLER
listeners=INTERNAL://0.0.0.0:9092
advertised.listeners=INTERNAL://broker-101.example.com:9092
inter.broker.listener.name=INTERNAL
listener.security.protocol.map=CONTROLLER:SSL,INTERNAL:SSL
log.dirs=/var/lib/kafka/data
三节点仲裁组需要任意两票才能写入元数据,因此可以容忍一个 Controller 故障。Controller 数量通常取奇数,例如 3 或 5;节点数、磁盘性能和跨机架布局应在创建集群时一起设计。Controller 扩缩容、角色变更和 ZooKeeper 到 KRaft 的迁移都有严格的版本前提,必须按照目标 Kafka 版本的官方迁移手册执行并先完成演练,不能把通用步骤当作可直接执行的变更脚本。
面试时应补充的风险点:
controller.quorum.voters、节点 ID 和元数据目录一旦配置错误,节点可能无法加入同一仲裁组。- Controller 的磁盘容量需求通常小于 Broker,但它保存元数据日志,仍需要稳定、持久的存储和独立的故障域。
- 不同 Kafka 发行版本对 ZooKeeper 模式与迁移路径的支持范围不同。升级前先确认目标版本支持的操作模式,而不是依赖固定的版本号或旧教程。
21 Kafka 的在线分区重分配
答案:
在线分区重分配是将分区副本从一组 Broker 迁移到另一组 Broker 的过程。它常用于扩容、下线节点、纠正副本分布或缓解磁盘倾斜;迁移期间业务通常仍可读写,但绝不是“没有影响”的操作。
触发场景:
| 场景 | 说明 |
|---|---|
| 集群扩容 | 新增代理节点后需要重新平衡分区 |
| 集群缩容 | 代理节点下线前迁移其分区 |
| 热点消除 | 将热点分区从高负载代理节点迁移 |
| 磁盘使用不均 | 均衡各代理节点的磁盘使用率 |
| 机架感知(Rack Awareness) | 修正不符合机架分布的副本 |
手动分区重分配(kafka-reassign-partitions.sh):
# 1. 生成 Topic 列表(JSON 格式)
cat > topics-to-move.json <<EOF
{
"topics": [
{"topic": "orders"},
{"topic": "payments"}
],
"version": 1
}
EOF
# 2. 生成候选方案;命令会打印建议的 reassignment JSON,需人工保存并审阅
bin/kafka-reassign-partitions.sh \
--bootstrap-server broker-1.example.com:9092 \
--topics-to-move-json-file topics-to-move.json \
--broker-list "0,1,2,3" \
--generate
# 3. 将审阅后的方案保存为 reassignment.json,再限速执行
bin/kafka-reassign-partitions.sh \
--bootstrap-server broker-1.example.com:9092 \
--reassignment-json-file reassignment.json \
--execute \
--throttle 104857600
# 4. 验证进度
bin/kafka-reassign-partitions.sh \
--bootstrap-server broker-1.example.com:9092 \
--reassignment-json-file reassignment.json \
--verify
--generate 只计算并输出候选方案,不会创建文件,也不会迁移数据。执行前应逐项检查每个分区的副本数、Broker ID、机架约束以及预计搬迁的数据量。限速参数和动态配置的具体写法随 Kafka 版本变化,以部署版本的命令帮助和运维手册为准。
重分配执行机制:
1. Controller 记录新的目标副本集合。
2. 新副本从当前 Leader 拉取数据并追赶日志。
3. 新副本追平后,Controller 切换副本集合并移除旧副本。
4. 如有 Leader 迁移,客户端会经历一次元数据刷新和短暂的角色切换。
- 先确认每个分区仍有足够 ISR,再分批迁移;不要在已有大量欠复制分区时扩大迁移范围。
- 搬迁会额外占用网络、磁盘和复制线程。应根据高峰余量设置限速,并持续观察欠复制分区、请求延迟和磁盘水位。
- 若迁移方案需要撤回,应先确认当前状态,再按已保存的原始副本分配生成回退操作。不要在不清楚副本是否追平时反复提交相互冲突的方案。
22 Kafka Connect 框架与连接器(Connector)管理
答案:
Kafka Connect 是 Kafka 的数据集成框架,用于将外部系统数据通过源连接器(Source Connector)导入 Kafka,或通过汇连接器(Sink Connector)从 Kafka 导出到外部系统。
生产通常使用分布式(Distributed)Worker。多个 Worker 通过 Kafka 内部 Topic 协调 Connector 与 Task 的分配,某个 Worker 故障后,其他 Worker 会接管其 Task;Standalone 适合本地开发或单进程场景,不提供这种协调与接管能力。
# connect-distributed.properties
bootstrap.servers=broker-1.example.com:9093,broker-2.example.com:9093
group.id=connect-orders-prod
config.storage.topic=_connect-orders-configs
offset.storage.topic=_connect-orders-offsets
status.storage.topic=_connect-orders-status
config.storage.replication.factor=3
offset.storage.replication.factor=3
status.storage.replication.factor=3
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=false
value.converter.schemas.enable=false
plugin.path=/opt/kafka-connect/plugins
config.storage.topic 应使用单分区、日志压缩(cleanup.policy=compact)的 Topic;offset 和 status Topic 也应启用压缩,并按可用性要求设置副本数。它们保存的是 Connect 的协调状态,丢失后即使外部系统和 Kafka 数据都在,也可能需要重新部署和校准连接器。
通过 REST API 创建连接器:
curl -sS -X POST 'https://connect.example.com/connectors' \
-H 'Content-Type: application/json' \
--data @- <<'JSON'
{
"name": "mysql-orders-source",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"tasks.max": "1",
"database.hostname": "mysql.example.com",
"database.port": "3306",
"database.user": "debezium",
"database.password": "${file:/etc/kafka-connect/secrets/mysql.properties:password}",
"topic.prefix": "cdc",
"database.include.list": "orders",
"table.include.list": "orders.orders",
"snapshot.mode": "initial"
}
}
JSON
示例中的 ${file:...} 依赖已启用的 Config Provider;也可以使用部署环境的密钥管理机制。密码不能直接写入 Connector 配置、日志或版本库。tasks.max 是可用的并行度上限,不保证每种连接器都会启动这么多 Task;例如部分 CDC 连接器受单个日志流或分片方式限制。
运行与提交语义:
Worker 读取配置 Topic → 任务在 Worker 间分配 →
Source Task 从外部系统读取并写入 Kafka,提交源端进度 →
Sink Task 消费 Kafka、写入目标系统、确认后再提交 Kafka 位点
源端和目标端大多不能与 Kafka 形成天然的全局事务。遇到重启、网络超时或下游部分成功时,常见结果是至少一次投递。因此目标系统要能幂等写入,或用业务唯一键处理重复;不要把 tasks.max 或 Worker 数量当作“恰好一次”的保证。
Connector 失败处理:
| 配置 | 行为 |
|---|---|
errors.tolerance=none | 任务失败立即停止(默认) |
errors.tolerance=all | 对可容忍的记录错误继续处理;必须配合告警和补偿,不能视为“错误已解决” |
errors.deadletterqueue.topic.name | 错误记录写入死信队列 |
errors.deadletterqueue.context.headers.enable=true | 死信消息附带上下文 Headers |
23 Kafka 的模式注册中心(Schema Registry)集成
答案:
模式注册中心(Schema Registry)保存 Avro、Protobuf 或 JSON Schema 的版本,并在注册新版本时检查兼容性。它是独立于 Apache Kafka Broker 的服务;Confluent Schema Registry、Apicurio Registry 等实现的部署方式和客户端包并不相同。
以 Confluent 兼容实现为例,服务端通常将模式元数据保存在 Kafka 的压缩 Topic 中:
# schema-registry.properties;安全相关参数按实际环境补齐
listeners=http://0.0.0.0:8081
kafkastore.bootstrap.servers=SSL://broker-1.example.com:9093,SSL://broker-2.example.com:9093
kafkastore.topic=_schemas
kafkastore.topic.replication.factor=3
schema.compatibility.level=BACKWARD
模式兼容性策略:
| 策略 | 检查方向 | 适用含义 |
|---|---|---|
BACKWARD | 新版本 Reader 能读取上一版本写入的数据 | 消费者先升级或需重放历史数据时常用 |
FORWARD | 旧版本 Reader 能读取新版本写入的数据 | 生产者先升级、消费者滞后时常用 |
FULL | 同时满足 BACKWARD 和 FORWARD | 生产者、消费者并行发布时更稳妥 |
*_TRANSITIVE | 对所有历史版本而不只是相邻版本检查 | 长期保留且多版本并存的 Topic |
NONE | 不做兼容性校验 | 仅适用于有额外治理手段的受控场景 |
“增加字段”或“删除字段”能否通过校验,取决于序列化格式、字段默认值、是否可空及 Reader/Writer Schema 的组合。不能只记住一张字段增删表;发布前应使用真实历史消息做兼容性测试。
客户端集成:
# 以下是 Confluent 兼容客户端的配置示例
# Producer
key.serializer=io.confluent.kafka.serializers.KafkaAvroSerializer
value.serializer=io.confluent.kafka.serializers.KafkaAvroSerializer
schema.registry.url=https://schema-registry.example.com
auto.register.schemas=false
# Consumer
key.deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer
value.deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer
schema.registry.url=https://schema-registry.example.com
specific.avro.reader=true
生产关注事项:
_schemas应启用cleanup.policy=compact,并按可用性目标设置副本数和min.insync.replicas。- 生产环境宜在 CI/CD 中注册并验证 Schema,而非让每个生产者在运行时任意注册新版本。
- 为 Schema Registry 的 TLS、认证、授权、审计和备份单独设计;Kafka 的 ACL 不会自动保护 Registry HTTP API。
- Registry 暂时不可用时,已缓存的 Schema ID 有时仍可工作;新注册、缓存未命中或新进程启动通常会失败。客户端缓存和超时策略要通过演练验证。
24 Kafka 的配额(Quota)限流机制
答案:
Kafka 的配额(Quota)机制在客户端标识(Client ID)或用户(User)粒度上对网络带宽和请求速率进行限流,防止单个客户端耗尽代理节点(Broker)资源。
Quota 类型:
| Quota 类型 | 配置参数 | 粒度 | 说明 |
|---|---|---|---|
| 网络带宽 | producer_byte_rate / consumer_byte_rate | Client ID / User | 限制生产/消费的字节速率(B/s) |
| 请求速率 | request_percentage | Client ID / User | 限制占用 I/O 线程和网络线程的时间百分比 |
| 控制面变更 | controller_mutation_rate | Client ID / User | 部分 Kafka 版本支持,用于约束创建、扩分区、删除等管理请求 |
配额可以按 User、Client ID 或二者组合配置。组合维度比单一维度更具体,但默认值和覆盖顺序会随实体配置而变化;上线前应用 --describe 检查实际生效的配置,不要只凭名称推断优先级。
动态设置 Quota(kafka-configs.sh):
# 设置 User 级别的生产、消费带宽限流
bin/kafka-configs.sh \
--bootstrap-server broker-1.example.com:9092 \
--entity-type users --entity-name my-user \
--alter --add-config producer_byte_rate=10485760,consumer_byte_rate=20971520
# 设置 Client ID 级别请求率限流
bin/kafka-configs.sh \
--bootstrap-server broker-1.example.com:9092 \
--entity-type clients --entity-name app-producer \
--alter --add-config request_percentage=30
# 查看当前 Quota 配置
bin/kafka-configs.sh \
--bootstrap-server broker-1.example.com:9092 \
--entity-type users --entity-name my-user --describe
限流时的行为:
- 超出
producer_byte_rate/consumer_byte_rate时,Broker 会延迟响应而不是立即断开连接;客户端通常表现为吞吐下降和请求延迟增加。 - 超出
request_percentage时,相关请求会被节流。它适合隔离明显失控的客户端,不是替代容量规划的手段。 quota.window.num和quota.window.size.seconds会影响采样与节流的平滑程度。阈值应从正常高峰值反推,并在压测中观察 P99 延迟。
Quota 监控:
| 指标 | 说明 |
|---|---|
| Produce / Fetch 的 throttle time | 客户端是否正在被明显节流 |
| Client ID、User 的请求与字节速率 | 找出是谁触发了配额 |
| 请求延迟、错误率与 Broker 网络利用率 | 判断阈值是否保护了集群,还是误伤正常流量 |
具体 JMX MBean 和 Prometheus 指标名取决于 Kafka 版本及导出器映射。监控面板应保留 User 与 Client ID 标签,否则很难定位被限流的业务方。
25 Kafka 的热点分区检测与处理
答案:
热点分区指流量、日志体量或处理延迟明显高于同一 Topic 其他分区的分区。它可能压满 Leader Broker,也可能让一个 Consumer 长期落后;两类问题要分开判断。
热点分区的成因:
| 成因 | 说明 |
|---|---|
| 消息键(Key)分布不均 | 特定消息键的消息量远大于其他消息键 |
| 处理能力不均 | 某些分区的业务逻辑、下游依赖或消息体更慢 |
| 代理节点硬件差异 | 主副本分布在性能不同的代理节点上 |
| 同一代理节点承载过多主副本 | 控制器分区分配策略导致主副本集中 |
检测方法:
# JMX 可查看 Topic 总吞吐;分区级吞吐通常需由 exporter、客户端或链路监控补充
bin/kafka-run-class.sh \
kafka.tools.JmxTool \
--object-name kafka.server:type=BrokerTopicMetrics,name=BytesInPerSec,topic=orders \
--reporting-interval 1000
# 通过 kafka-log-dirs.sh 查看各分区磁盘占用
bin/kafka-log-dirs.sh \
--bootstrap-server broker-1.example.com:9092 \
--describe \
--topic-list orders
# 查看 Topic 分区 Leader 分布
bin/kafka-topics.sh \
--bootstrap-server broker-1.example.com:9092 \
--describe --topic orders
Prometheus 热点分区查询:
# 以下假设 exporter 暴露了带 topic、partition 标签的分区字节计数器;实际指标名按 exporter 调整
topk(10, sum by (topic, partition) (rate(kafka_partition_bytes_in_total[5m])))
# 分区大小 Top 10
topk(10, sum by (topic, partition) (kafka_log_partition_size_bytes))
# 单 Broker 上 Leader 分区数
sum by (instance) (kafka_server_replicamanager_leadercount)
分区级指标并非 Kafka 所有 JMX 导出器的默认能力。若看不到 partition 标签,不能据 Topic 总吞吐推断具体热点分区,应补齐客户端埋点、消费位点、日志目录和下游耗时数据。
处理策略:
| 策略 | 方法 | 效果 |
|---|---|---|
| 修改分区 Key | 增加 Key 粒度或散列化 Key | 从源头均衡数据分布 |
| 增加分区数 | 提高 Topic 分区数 | 只会为更多 Key 提供落点,不能拆开一个严格有序的热 Key |
| Leader 均衡 | Preferred Leader Election 或自动 Leader 均衡 | 降低 Broker 间 Leader 倾斜,不降低单个热点分区的总流量 |
| 分区重分配 | Cruise Control 或 kafka-reassign-partitions.sh | 迁移分区到更有余量的 Broker,不改变该分区的吞吐上限 |
| Consumer 调整 | 优化处理路径、增加 Consumer 实例 | 只能改善可并行分区的消费能力,无法并行消费同一分区 |
真正由单一业务 Key 引起的热点,没有不改变语义的“Kafka 参数调优”可以消除。若业务必须保持该 Key 的全序,只能提高这一个分区及其 Leader 的处理能力,或调整业务聚合方式;若业务允许把顺序范围缩小到子 Key,才可以采用确定性的分片规则,例如 orderId + shardId。随机追加后缀会打乱同一 Key 的顺序,不能作为默认方案。
增加分区、修改分区器或改变 Key 都可能改变路由结果。上线前要明确历史数据与新数据的消费策略、顺序边界和去重方式,并通过回放压测确认结果。
26 Kafka 的生产者幂等性(Idempotent Producer)
答案:
生产者幂等性用于处理网络超时、Broker 响应丢失等重试场景:同一生产者会话把同一批记录重发时,Broker 不会再次追加日志。它是 Kafka 恰好一次语义的一部分,但不等于端到端的全局去重。
幂等性实现原理:
Producer 初始化 → 向事务协调器请求 Producer ID(PID)和 Epoch →
Producer 为每个 <PID, TopicPartition> 维护 Sequence Number(从 0 开始)→
发送消息时附加 <PID, Epoch, SequenceNumber> →
Leader Broker 为每个 <PID, TopicPartition> 维护已确认的序列状态 →
收到消息:
1. Sequence 与预期值连续 → 接受并更新状态
2. Sequence 已处理过 → 识别为重试,不重复写入
3. Sequence 跳号或 Epoch 过期 → 拒绝该请求,客户端需要处理乱序或被 fencing 的情况
幂等性配置:
enable.idempotence=true # 启用幂等
acks=all # 幂等性要求所有 ISR 确认
max.in.flight.requests.per.connection=5 # 启用幂等性时不能大于 5
新版本客户端通常会对幂等性相关的 acks、重试和 in-flight 配置做兼容性校验;不要为“无限重试”机械写入一个极大整数。重试时长还受 delivery.timeout.ms、request.timeout.ms 和业务超时约束。
Broker 端状态:
事务协调器负责分配和管理 PID / Epoch。
分区 Leader 将最近的生产者序列状态保存在日志相关的 producer state snapshot 中;
Broker 重启后可从日志与快照恢复,避免把历史重试当成新消息。
幂等性的局限:
| 场景 | 幂等性是否有效 | 说明 |
|---|---|---|
| Producer 内部重试 | 有效 | 相同 PID + SequenceNumber 被去重 |
| 普通 Producer 重启 | 不保证跨会话去重 | 新会话通常会获得新的 PID,Broker 无法把它与旧会话的业务消息自动关联 |
同一 transactional.id 的新实例接管 | 可安全 fencing 旧实例 | Epoch 会使旧实例失效,避免两个实例并发写入同一事务身份 |
| 外部数据库、HTTP 服务 | 不保证 | Kafka 的序列号无法让外部系统识别或撤销重复请求 |
与事务的配合:
幂等生产者 → 单分区 / 单 Session 内的去重保证
事务生产者 →
跨分区原子写入 +
Kafka read-process-write 链路中的 Exactly-Once +
使用稳定 transactional.id fencing 旧实例
事务生产者提交位点时需要把消费位点与输出记录放在同一个 Kafka 事务中,消费者还必须使用 isolation.level=read_committed。一旦链路跨出 Kafka,例如写数据库、调用支付接口,仍需业务幂等键、Outbox 或补偿机制。
使用边界:
transactional.id是一个稳定的生产者逻辑身份,不能被多个并发运行的实例随意共用。transaction.timeout.ms需要小于 Broker 允许的最大事务超时,并结合最长业务处理时间设置。max.in.flight.requests.per.connection在幂等性开启时可大于 1,但上限受客户端约束;不必为保序一律降为 1。
27 Kafka 常见故障排查
答案:
先按“客户端、集群控制面、分区副本、存储资源”分层缩小范围,再针对受影响的主题(Topic)、分区(Partition)和消费者组(Consumer Group)查看日志与指标。不要一开始就通过重启代理节点(Broker)或重置位点(Offset)处理问题,否则可能扩大影响或丢失现场。
场景一:代理节点(Broker)不可用或客户端无法连接
症状:客户端报 Connection to node failed、NotLeaderOrFollower、
TimeoutException,或集群中出现 Offline Partition。
排查重点:
1. 确认受影响的是单个 Broker、单个 Listener,还是整个集群。
2. 查看 Broker 日志、JVM/GC、磁盘空间和进程监听端口。
3. 核对 listeners、advertised.listeners、DNS、负载均衡和防火墙配置。
4. TLS/SASL 集群还需核对证书有效期、客户端协议与 ACL。
处理原则:
- 先恢复网络、进程或磁盘等根因,再确认 Controller 已完成 Leader 迁移。
- 不要把 advertised.listeners 配成客户端无法解析的内网地址。
场景二:副本未完全同步(Under Replicated Partitions)持续大于 0
bin/kafka-topics.sh --bootstrap-server broker-1:9092 \
--describe --under-replicated-partitions
常见原因包括代理节点宕机、代理节点间网络异常、磁盘 I/O 饱和、长时间 GC 暂停,以及从副本(Follower)拉取速度不足。应结合受影响分区的主副本/ISR、代理节点日志和磁盘延迟确认根因;恢复故障代理节点或网络后等待副本追平,必要时再评估分区迁移、容量扩展或副本拉取参数调整。
场景三:消费滞后(Consumer Lag)持续增长
bin/kafka-consumer-groups.sh --bootstrap-server broker-1:9092 \
--group order-processor --describe
先判断消费滞后是集中在少数分区还是普遍增长:前者通常是热点消息键或单分区消费阻塞,后者通常是整体处理能力不足、下游依赖变慢或再均衡频繁。消费者实例数不能超过可分配的分区数;扩容前应先确认处理耗时、提交方式和下游服务状态。若要重置位点,必须明确业务可接受的数据跳过或重复消费范围。
场景四:生产者(Producer)请求超时或写入失败
关注生产者端错误类型和代理节点端的 Produce 请求延迟、请求队列、GC、网络及磁盘指标。只要 ISR 中可确认的副本不足,写入就可能超时;盲目降低确认级别会降低可靠性,只应在业务明确可接受风险时采用。还应核对消息大小、超时与重试策略,以及代理节点的最大消息大小配置是否匹配。
场景五:磁盘空间或日志目录异常
bin/kafka-log-dirs.sh --bootstrap-server broker-1:9092 --describe
磁盘满、单块磁盘失效或日志目录不可写会影响分区副本和主副本(Leader)可用性。先确认保留策略、日志段(segment)、日志目录使用率及是否存在热点分区;对可删除数据按保留策略清理或扩容,对不可丢失数据应先恢复副本和备份,再进行迁移或故障磁盘处理。
参考资料:Apache Kafka 官方文档(kafka.apache.org/documentation)、KRaft 设计文档