几道kafka与pulsar基础面试题
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 端丢失场景
消费分为两个阶段:
- 获取元数据并从 Kafka Broker 集群拉取数据。
- 处理消息并提交 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=syncrequest.required.acks=1- 副本数量
>= 2 - 增加重试次数
- 使用
Producer.send(msg, callback)
- 异步方式:
producer.type=asyncrequest.required.acks=1queue.buffering.max.ms=5000queue.buffering.max.messages=10000queue.enqueue.timeout.ms = -1batch.num.messages=200- 通过 buffer 控制数据发送,设置阻塞模式防止数据丢失。
Consumer 端:
- 关闭自动提交 Offset,手动提交:
- 设置
enable.auto.commit = false - 使用
consumer.commitSync()或ack.acknowledge()提交。
- 设置
- 手动提交 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 数据积压怎么解决?
- 增加 Broker 节点,增加分区数量,提高并行度。
- 修改单线程消费为批量消费。
- 增加单线程消费为线程池异步消费。
- 缩短批次时间间隔。
- 控制消费速率,设置背压机制。
- 优化代码,减少 shuffle 过程。
- 将处理结果保存到 MySQL、MongoDB 或 ES,增大吞吐量。
- 使用滑动窗口控制拉取速度。
- 处理倾斜的 key,加随机数打散。
Kafka Rebalance 发生时机和分区分配策略?
Rebalance 触发时机:
- 同一 Consumer Group 中新增消费者。
- 消费者离开 Consumer Group。
- 分区数量发生变化。
- 消费者主动取消订阅。
Rebalance 分区分配策略:
- Range
- Round-Robin
- Sticky
Kafka 分区数如何确定?
分区数的确定依据:
- 根据生产者和消费者的目标吞吐量估计。
- 公式:
numPartitions = Tt / max(Tp, Tc),其中 Tt 为目标吞吐量,Tp 为 Producer 吞吐量,Tc 为 Consumer 吞吐量。
分区数过多的危害:
- 客户端/服务器端内存开销增加。
- 文件句柄开销增加。
- 增加端对端延迟。
- 降低高可用性。
数据发往 Kafka 的分区规则?
- 若指定分区,写入指定分区。
- 若未指定分区但指定 key,按 key 的 hashcode 取模写入对应分区。
- 若未指定分区和 key,采用轮询机制。
Kafka Producer Buffer Pool 的作用?
Kafka 使用内存缓冲池设计,减少 JVM GC 影响,提高发送性能和吞吐量。功能包括:
- 限制可申请的内存总量,防止 OOM。
- 持有 ProducerRecord 的引用,减少 FullGC 频率。
Kafka 时间轮的作用?
Kafka 通过时间轮处理延迟任务,减少延迟队列元素数量,提高增加删除性能。通过阻塞方式 poll 延迟队列,减少空转。使用读写锁、原子对象、synchronized 保证线程安全。
Kafka 为什么这么快?
Kafka 通过以下机制实现高性能:
- 顺序读写磁盘。
- 零拷贝技术减少数据拷贝次数。
- 批量发送和压缩消息。
- 分区和副本机制提高并行度和容错能力。
Leader Replica 选举触发时机?
- Leader Replica 失效。
- Broker 宕机。
- 新增 Broker。
- 新建分区。
- ISR 列表数量减少。
- 手动触发选举。
Leader Replica 选举策略?
- ISR 选举策略:默认从 ISR 集合中选举 Leader。
- 首选副本选举策略:优先选择首选副本作为 Leader。
- 不干净副本选举策略:从所有副本中选择 Leader,可能导致数据丢失。
Kafka 日志清理策略?
Kafka 提供两种日志清理策略:
- 日志删除:按时间、日志大小或日志起始偏移量删除不符合条件的日志分段。
- 日志压缩:针对每个消息的 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
