前言

Kafka 面试题通常会从“为什么 Kafka 吞吐高”开始,继续追问 Topic 和 Partition 的关系、副本如何保证高可用、ISR 是什么、acks 怎么选、消费者组如何分摊消费、Offset 提交有什么风险、如何保证消息顺序、如何处理重复消费和消费积压。

如果只把 Kafka 理解成“一个消息队列”,回答很容易停留在使用层。更好的方式是把它理解成一条完整链路:

1
Producer 写入 Topic -> 根据 Key 路由到 Partition -> Leader Partition 追加日志 -> Replica 拉取同步 -> Consumer Group 按分区消费 -> 提交 Offset

本文按面试中最常见的模块梳理 Kafka 要点,适合作为面试前的复习清单。

Kafka 整体定位

Kafka 是一个分布式事件流平台,常用于消息队列、日志采集、异步解耦、削峰填谷、数据同步、流式计算和事件驱动架构。

它的核心能力来自三个设计:

  • 追加写日志:消息以顺序追加方式写入磁盘。
  • 分区机制:Topic 拆成多个 Partition,实现并行写入和消费。
  • 副本机制:Partition 有多个副本,保证高可用。

面试中可以这样概括:

Kafka 不是单纯的队列,而是基于分区日志的分布式消息系统。生产者追加写,消费者按 Offset 顺序读,分区带来并行能力,副本带来容灾能力。

核心概念

Kafka 常见概念如下:

  • Broker:Kafka 服务节点。
  • Topic:消息主题,业务上的消息分类。
  • Partition:Topic 的分区,是 Kafka 并行读写的基本单位。
  • Replica:分区副本,用于高可用。
  • Leader Replica:负责读写请求的主副本。
  • Follower Replica:从 Leader 拉取数据的副本。
  • Producer:消息生产者。
  • Consumer:消息消费者。
  • Consumer Group:消费者组,用于分摊消费。
  • Offset:消息在 Partition 中的位置。
  • Controller:负责集群元数据管理和分区 Leader 选举的角色。
flowchart TD
    A[Producer] --> B[Topic: order-event]
    B --> C[Partition 0 Leader]
    B --> D[Partition 1 Leader]
    C --> E[Replica 0 Follower]
    D --> F[Replica 1 Follower]
    C --> G[Consumer A]
    D --> H[Consumer B]

Topic 和 Partition

Topic 是逻辑概念,Partition 是物理并行单元。

一个 Topic 可以有多个 Partition,每个 Partition 内部消息有序,不同 Partition 之间没有全局顺序。

Partition 带来的好处:

  • 提高写入并行度。
  • 提高消费并行度。
  • 支持水平扩展。
  • 单个 Topic 可以分布在多个 Broker 上。

常见追问:Partition 数量是不是越多越好?

不是。Partition 太少会限制吞吐和并行消费;Partition 太多会增加文件句柄、内存、Leader 选举、Controller 元数据和恢复成本,也可能导致重平衡变慢。

消息存储

Kafka 的消息以日志形式存储在 Partition 中。每个 Partition 是一个有序、不可变、只能追加的日志序列。

1
2
3
4
Partition 0:
offset 0 -> message A
offset 1 -> message B
offset 2 -> message C

Kafka 会把日志切成多个 Segment 文件。每个 Segment 通常包括:

  • .log:消息数据文件。
  • .index:Offset 索引。
  • .timeindex:时间索引。

Kafka 读写快的原因:

  • 顺序追加写磁盘。
  • Page Cache 充分利用操作系统缓存。
  • 零拷贝减少数据复制。
  • 批量发送和压缩降低网络开销。
  • Partition 并行提升吞吐。

生产者发送流程

生产者发送消息时,会先序列化、选择分区、批量缓冲,再发送给 Broker。

sequenceDiagram
    participant App
    participant Producer
    participant Broker as Leader Broker
    participant Replica as Follower Replica

    App->>Producer: send(record)
    Producer->>Producer: 序列化和选择分区
    Producer->>Producer: 累积到 Batch
    Producer->>Broker: 发送 Produce Request
    Broker->>Broker: 追加到 Leader Log
    Broker->>Replica: 副本同步
    Broker-->>Producer: 返回 ack

分区选择规则常见有:

  • 指定 Partition:直接写指定分区。
  • 指定 Key:对 Key 做 hash,路由到固定 Partition。
  • 无 Key:按粘性分区或轮询策略分配,提高批量效率。

如果业务要求同一个订单的消息有序,通常使用订单 ID 作为 Key,让同一订单进入同一个 Partition。

acks 和可靠性

生产者的 acks 决定写入确认级别。

acks 含义 风险
0 生产者不等 Broker 确认 吞吐高,但可能丢消息
1 Leader 写入成功即确认 Leader 宕机且副本未同步时可能丢消息
all / -1 ISR 中副本确认后再返回 可靠性高,延迟更高

常见可靠性组合:

1
2
3
4
5
acks=all
retries > 0
enable.idempotence=true
min.insync.replicas >= 2
replication.factor >= 3

这里要注意:如果 acks=all,但 min.insync.replicas=1,可靠性仍然有限。生产环境通常会配合副本数和最小同步副本数一起设计。

副本和 ISR

Kafka 的每个 Partition 可以有多个 Replica,其中一个是 Leader,其他是 Follower。生产和消费请求通常由 Leader 处理,Follower 从 Leader 拉取数据。

ISR 是 In-Sync Replicas,表示和 Leader 保持同步的副本集合。

1
2
3
Partition 0:
Leader: broker-1
ISR: broker-1, broker-2, broker-3

如果 Follower 落后太多,会被移出 ISR。写入时如果要求 acks=all,Kafka 会等待 ISR 中满足条件的副本确认。

常见追问:Leader 宕机怎么办?

Controller 会从 ISR 中选择新的 Leader。如果允许 unclean leader election,可能从非 ISR 副本中选 Leader,这样可用性提高,但可能丢数据。生产环境通常谨慎开启。

消费者组

Kafka 使用 Consumer Group 实现消息分摊消费。

同一个 Group 内,一个 Partition 同一时刻只能分配给一个 Consumer;一个 Consumer 可以消费多个 Partition。

1
2
3
Topic 有 3 个 Partition,Group 有 2 个 Consumer:
Consumer A -> Partition 0, Partition 1
Consumer B -> Partition 2

几个关键结论:

  • 同组内消费者是竞争关系,共同分摊消息。
  • 不同消费者组之间互不影响,各自都能消费完整消息。
  • 消费并行度受 Partition 数量限制。
  • Consumer 数量超过 Partition 数量时,多余 Consumer 会空闲。

Offset 提交

Offset 表示消费者在某个 Partition 消费到的位置。Kafka 会把消费者组的 Offset 存储在内部 Topic __consumer_offsets 中。

Offset 提交方式:

  • 自动提交:简单,但可能出现消息处理失败却已提交的情况。
  • 手动提交:可控,适合对可靠性要求更高的业务。

常见风险:

  • 先提交 Offset,再处理业务:业务失败会丢消息。
  • 先处理业务,再提交 Offset:提交失败或消费者宕机会重复消费。

所以 Kafka 消费端通常要接受“至少一次”语义,并通过业务幂等处理重复消息。

消息顺序性

Kafka 只能保证单 Partition 内消息有序,不能保证多个 Partition 之间全局有序。

保证局部顺序的常见方式:

  • 同一业务实体使用相同 Key。
  • 让相同 Key 路由到同一 Partition。
  • 单个 Consumer 顺序处理该 Partition。
  • 避免处理逻辑中异步乱序提交。

例如订单状态流可以用 orderId 作为 Key:

1
orderId=1001 -> Partition 2

这样同一订单的创建、支付、发货、完成事件会进入同一个 Partition,消费时保持顺序。

如果要求全局顺序,只能使用单 Partition,但会牺牲吞吐和扩展能力。面试中要明确这个取舍。

重复消费和幂等

Kafka 消费端出现重复消费是正常情况,常见原因包括:

  • 消费成功但 Offset 提交失败。
  • 消费者处理超时触发重平衡。
  • 消费者宕机后由其他消费者接管。
  • 生产者重试导致重复发送。

解决思路:

  • 业务侧做幂等。
  • 使用唯一业务 ID 去重。
  • 数据库写入使用唯一约束。
  • 状态更新要判断版本或状态流转合法性。
  • 对外部调用保存请求流水,避免重复扣款、重复发券。

面试中可以直接说:Kafka 能提升投递可靠性,但端到端不重复通常需要生产者、Kafka、消费者和业务存储共同配合。

幂等生产者和事务

Kafka 生产者支持幂等发送,开启:

1
enable.idempotence=true

幂等生产者可以避免单个 Producer 会话内,因为重试导致同一分区消息重复写入。

Kafka 事务用于实现跨多个 Partition 的原子写入,以及消费-处理-生产链路中的 Exactly Once 语义。典型场景是流处理:

1
consume input topic -> process -> produce output topic -> commit consumed offsets

要注意:Kafka 的 Exactly Once 主要解决 Kafka 内部读写链路的一致性。如果业务还写数据库、调用 HTTP 接口,仍然需要外部系统配合幂等或事务设计。

消费积压

消费积压是 Kafka 面试高频问题。积压通常表示生产速度长期大于消费速度。

常见原因:

  • 消费者处理逻辑慢。
  • Consumer 数量不足。
  • Partition 数量不足,无法增加有效并行度。
  • 下游数据库、接口或缓存变慢。
  • 单条消息处理异常反复重试。
  • Rebalance 频繁。
  • 消费端批量拉取或提交参数不合理。

排查路径:

  1. 看 Lag 是单个 Partition 高,还是所有 Partition 都高。
  2. 看生产速率是否突然增加。
  3. 看消费耗时、失败率和重试次数。
  4. 看下游依赖是否变慢。
  5. 看 Consumer Group 是否频繁 Rebalance。
  6. 看 Partition 数量和 Consumer 数量是否匹配。
  7. 看是否有热点 Key 导致某个 Partition 特别忙。

优化方式:

  • 提升单条消息处理效率。
  • 增加 Consumer 实例,但不能超过有效 Partition 并行度。
  • 增加 Partition 数量,提高并行消费能力。
  • 批量处理和批量写下游。
  • 慢操作异步化或拆分。
  • 对异常消息进入死信队列。

Rebalance

Rebalance 是消费者组内 Partition 重新分配的过程。

触发原因:

  • Consumer 加入或退出。
  • Consumer 心跳超时。
  • Topic Partition 数变化。
  • 订阅 Topic 变化。

Rebalance 期间,部分消费者可能暂停消费。如果 Rebalance 频繁,会导致吞吐下降和延迟波动。

常见优化:

  • 保证消费处理时间小于 max.poll.interval.ms
  • 合理设置 session.timeout.msheartbeat.interval.ms
  • 避免消费者频繁重启。
  • 使用静态成员减少重平衡影响。
  • 使用 cooperative sticky assignor 降低全量重分配成本。

Kafka 为什么吞吐高

Kafka 高吞吐不是单点原因,而是一组设计共同作用:

  • Partition 并行读写。
  • 顺序追加写日志。
  • 操作系统 Page Cache。
  • 批量发送和批量拉取。
  • 零拷贝。
  • 压缩减少网络传输。
  • Consumer 自己维护 Offset,Broker 负担轻。
  • 磁盘顺序 IO 性能高。

面试回答时不要只说“顺序写磁盘”,还要结合批量、Page Cache、零拷贝和分区并行。

Kafka 和 RabbitMQ 的区别

Kafka 和 RabbitMQ 都能做消息系统,但设计目标不同。

对比项 Kafka RabbitMQ
核心模型 分区日志 Exchange + Queue
消费方式 按 Offset 拉取 Broker 推送或消费者拉取
消息保留 按时间或大小保留,可重复消费 消费确认后通常删除
吞吐能力 高吞吐、适合流式数据 路由灵活、延迟较低
顺序性 单 Partition 有序 单队列有序
典型场景 日志、埋点、数据管道、流处理 业务异步、任务队列、复杂路由

面试中可以总结:Kafka 更像可持久化的分布式提交日志,RabbitMQ 更像传统消息代理。

常见性能优化

生产端:

  • 使用批量发送。
  • 调整 batch.sizelinger.ms
  • 开启压缩,例如 lz4、snappy、zstd。
  • 合理设置 acks 和重试。
  • 使用业务 Key 保证必要顺序。

Broker:

  • 合理规划 Partition 数量。
  • 设置合适副本数。
  • 监控磁盘 IO、网络、Page Cache 和请求延迟。
  • 避免单 Broker 承担过多 Leader。
  • 关注 Controller、ISR Shrink、Under Replicated Partitions。

消费端:

  • 批量拉取和批量处理。
  • 控制单次处理时间。
  • 业务处理做幂等。
  • 增加 Consumer 并行度。
  • 避免频繁 Rebalance。

常见排查思路

如果面试官问“Kafka 消息延迟变高怎么排查”,可以按以下路径回答:

  1. 看是生产延迟、Broker 延迟,还是消费延迟。
  2. 看 Topic、Partition 和 Consumer Group 的 Lag。
  3. 看是否只有某些 Partition 延迟高,判断热点 Key。
  4. 看 Broker 磁盘 IO、网络、CPU、请求队列和 Page Cache。
  5. 看 ISR 是否频繁收缩,是否有副本同步慢。
  6. 看 Consumer 是否处理慢、提交慢或频繁 Rebalance。
  7. 看下游数据库、缓存、接口是否拖慢消费。
  8. 看是否有大批量消息、消息体变大或压缩配置变化。

如果是丢消息问题,可以重点检查生产者 acks、重试、幂等、Broker 副本配置、min.insync.replicas、unclean leader election、消费者 Offset 提交顺序和业务处理结果。

高频面试题

Kafka 为什么快?

因为 Kafka 使用顺序追加写日志、Page Cache、零拷贝、批量发送、压缩和 Partition 并行读写。它避免了大量随机 IO,并把吞吐建立在顺序 IO 和批处理之上。

Kafka 如何保证消息不丢?

生产端设置 acks=all、开启重试和幂等;Broker 设置足够副本数和 min.insync.replicas;消费端处理成功后再提交 Offset;业务侧做好幂等和失败重试。端到端可靠性需要多层配合。

Kafka 如何保证顺序?

Kafka 只能保证单 Partition 内有序。把同一业务实体的消息用相同 Key 发送到同一 Partition,并由消费者顺序处理。如果要求全局有序,只能使用单 Partition,但吞吐会受限。

Kafka 会重复消费吗?

会。消费成功但 Offset 提交失败、消费者宕机、重平衡、生产者重试都可能导致重复。实际项目通常通过业务幂等、唯一键、去重表或状态机来处理。

Consumer 数量越多越好吗?

不是。同一个消费者组内,一个 Partition 同时只能被一个 Consumer 消费。如果 Consumer 数量超过 Partition 数量,多出的 Consumer 会空闲。

什么是 ISR?

ISR 是和 Leader 保持同步的副本集合。Leader 宕机后通常从 ISR 中选新 Leader。ISR 越健康,Kafka 的高可用和数据可靠性越好。

消费积压怎么处理?

先判断积压范围和原因:生产突增、消费者慢、下游慢、Partition 不足、热点 Key、异常消息或频繁 Rebalance。然后再做扩容消费者、增加 Partition、优化处理逻辑、批量写下游、隔离异常消息等处理。

总结

Kafka 面试的主线可以围绕四个问题展开:

  1. 消息怎么写:Producer、Partition、Leader、acks、ISR。
  2. 消息怎么存:追加日志、Segment、Page Cache、副本。
  3. 消息怎么读:Consumer Group、Partition 分配、Offset。
  4. 线上怎么治理:顺序性、重复消费、积压、Rebalance、可靠性和性能优化。

把这条链路讲清楚,再结合幂等、事务、消费积压和 Broker 排查,Kafka 相关问题就能从“会用消息队列”升级成“理解分布式事件流系统”。