改版通知

巨人肩膀网站已全新改版。若您仍依赖旧站功能或数据,欢迎联系我们,我们会协助处理。联系我们

几道kafka与pulsar基础面试题

ckckck2025年1月10日54 浏览

Kafka 可能会丢数据的场景

Producer 端丢失场景

Producer 端发送数据有 ACK 机制,可能会丢数据:

  • acks = 0:发送后自认为成功,若发生网络抖动,Producer 不会校验 ACK,数据丢失且无法重试。
  • acks = 1:消息发送到 Leader Partition 成功即认为成功。若 Leader Partition 崩溃且 Follower Partition 未同步完数据,数据丢失。
  • acks = -1 或 all:消息发送需等待 ISR 中 Leader Partition 和所有 Follower Partition 确认。可靠性最高,但若 ISR 中只剩 Leader Partition,则等同于 acks = 1。

Broker 端丢失场景

Broker 端消息存储通过异步批量刷盘,可能会丢数据:

  • Kafka 未提供同步刷盘方式,单个 Broker 可能丢失数据。
  • 通过多 Partition 多 Replica 机制最大限度保证数据不丢失,但若数据写入 PageCache 未刷盘且 Broker 宕机,极端情况下数据丢失。

Consumer 端丢失场景

消费分为两个阶段:

  1. 获取元数据并从 Kafka Broker 集群拉取数据。
  2. 处理消息并提交 Offset 记录。

可能丢数据的场景:

  • 自动提交 Offset 方式
    • 先提交 Offset 后处理消息:若处理消息时宕机,Offset 已提交,重启后从下一个 Offset 开始消费,未处理的消息丢失。
    • 先处理消息后提交 Offset:若提交前宕机,Offset 未提交,重启后重新拉取消息,不会丢失但可能重复消费。

Kafka 如何保证数据不丢失?

Broker 端:

  • Topic 副本因子个数replication.factor >= 3
  • 同步副本列表 (ISR)min.insync.replicas = 2
  • 禁用 unclean 选举unclean.leader.election.enable=false
  • 副本因子与 ISR 关系replication.factor = min.insync.replicas + 1

Producer 端:

  • 同步方式
    • producer.type=sync
    • request.required.acks=1
    • 副本数量 >= 2
    • 增加重试次数
    • 使用 Producer.send(msg, callback)
  • 异步方式
    • producer.type=async
    • request.required.acks=1
    • queue.buffering.max.ms=5000
    • queue.buffering.max.messages=10000
    • queue.enqueue.timeout.ms = -1
    • batch.num.messages=200
    • 通过 buffer 控制数据发送,设置阻塞模式防止数据丢失。

Consumer 端:

  1. 关闭自动提交 Offset,手动提交
    • 设置 enable.auto.commit = false
    • 使用 consumer.commitSync()ack.acknowledge() 提交。
  2. 手动提交 Offset 并缓存数据
    • 将拉取的数据缓存到 queue,处理完后再批量提交 Offset。

Kafka 如何保证数据 Exactly-Once?

Producer Exactly-Once:

  • enable.idempotence=true
  • 分区副本数 >= 2
  • ISR >= 2
  • 使用 ProducerID + SequenceNumber + Ack = -1 保证幂等性。

Consumer Exactly-Once:

  • 手动维护并提交偏移量:
    • 设置 enable.auto.commit=false
    • 借助外部数据库(如 Redis、MySQL)管理偏移量,确保消息处理完后再提交偏移量。

Kafka 数据积压怎么解决?

  1. 增加 Broker 节点,增加分区数量,提高并行度。
  2. 修改单线程消费为批量消费。
  3. 增加单线程消费为线程池异步消费。
  4. 缩短批次时间间隔。
  5. 控制消费速率,设置背压机制。
  6. 优化代码,减少 shuffle 过程。
  7. 将处理结果保存到 MySQL、MongoDB 或 ES,增大吞吐量。
  8. 使用滑动窗口控制拉取速度。
  9. 处理倾斜的 key,加随机数打散。

Kafka Rebalance 发生时机和分区分配策略?

Rebalance 触发时机:

  1. 同一 Consumer Group 中新增消费者。
  2. 消费者离开 Consumer Group。
  3. 分区数量发生变化。
  4. 消费者主动取消订阅。

Rebalance 分区分配策略:

  • Range
  • Round-Robin
  • Sticky

Kafka 分区数如何确定?

分区数的确定依据:

  • 根据生产者和消费者的目标吞吐量估计。
  • 公式:numPartitions = Tt / max(Tp, Tc),其中 Tt 为目标吞吐量,Tp 为 Producer 吞吐量,Tc 为 Consumer 吞吐量。

分区数过多的危害:

  1. 客户端/服务器端内存开销增加。
  2. 文件句柄开销增加。
  3. 增加端对端延迟。
  4. 降低高可用性。

数据发往 Kafka 的分区规则?

  • 若指定分区,写入指定分区。
  • 若未指定分区但指定 key,按 key 的 hashcode 取模写入对应分区。
  • 若未指定分区和 key,采用轮询机制。

Kafka Producer Buffer Pool 的作用?

Kafka 使用内存缓冲池设计,减少 JVM GC 影响,提高发送性能和吞吐量。功能包括:

  • 限制可申请的内存总量,防止 OOM。
  • 持有 ProducerRecord 的引用,减少 FullGC 频率。

Kafka 时间轮的作用?

Kafka 通过时间轮处理延迟任务,减少延迟队列元素数量,提高增加删除性能。通过阻塞方式 poll 延迟队列,减少空转。使用读写锁、原子对象、synchronized 保证线程安全。

Kafka 为什么这么快?

Kafka 通过以下机制实现高性能:

  • 顺序读写磁盘。
  • 零拷贝技术减少数据拷贝次数。
  • 批量发送和压缩消息。
  • 分区和副本机制提高并行度和容错能力。

Leader Replica 选举触发时机?

  1. Leader Replica 失效。
  2. Broker 宕机。
  3. 新增 Broker。
  4. 新建分区。
  5. ISR 列表数量减少。
  6. 手动触发选举。

Leader Replica 选举策略?

  1. ISR 选举策略:默认从 ISR 集合中选举 Leader。
  2. 首选副本选举策略:优先选择首选副本作为 Leader。
  3. 不干净副本选举策略:从所有副本中选择 Leader,可能导致数据丢失。

Kafka 日志清理策略?

Kafka 提供两种日志清理策略:

  1. 日志删除:按时间、日志大小或日志起始偏移量删除不符合条件的日志分段。
  2. 日志压缩:针对每个消息的 key 进行整合,保留最新版本。

Kafka 零拷贝

Kafka 使用 Linux 的 sendFile 技术,省去进程切换和一次数据拷贝,提升性能。传统方式发生 4 次拷贝,2 次 DMA 和 2 次 CPU,而零拷贝方式只发生 2 次上下文切换和 3 次数据拷贝。

Kafka 数据存储结构?

Kafka 数据存储结构包括:

  • 分区(Partition)
  • 分段(Segment)
  • 索引文件(Index File)
  • 日志文件(Log File)

Kafka、RabbitMQ、RocketMQ、Pulsar 选型建议

Apache Kafka

  • 优点:高吞吐量、低延迟,适合大数据处理。
  • 缺点:复杂,学习曲线陡峭。
  • 适用场景:日志聚合、流式数据处理、实时数据分析。

RabbitMQ

  • 优点:支持多种消息协议,易于管理。
  • 缺点:单节点吞吐量较低。
  • 适用场景:微服务间通信、复杂消息路由。

RocketMQ

  • 优点:高性能,支持事务消息。
  • 缺点:社区相对较小。
  • 适用场景:高并发、高吞吐量场景。

Apache Pulsar

  • 优点:支持多租户、地理分布。
  • 缺点:配置和运维复杂。
  • 适用场景:跨数据中心的分布式消息处理。

Pulsar Bundle 理解

Pulsar 将 Namespace 拆分为 Bundle,Bundle 是 Namespace 的子集。Topic 通过哈希分配到 Bundle,Bundle 再分配到不同的 Broker。Bundle 的分配由 Load Manager 负责,确保负载均衡。

Pulsar 存算分离和读写分离

Pulsar 采用存储计算分离架构,Broker 无状态,BookKeeper 负责持久化存储。读写分离通过 Journal 和 Memtable 实现,写入时数据实时写入 Journal,读取时优先从 Memtable 读取,提升性能。

end