flink和kafka常见面试题
1. Kafka 为什么吞吐量这么高?
Kafka 之所以吞吐量高,主要依赖以下几个设计:
首先,Kafka 采用顺序写磁盘,消息以追加方式写入 Log 文件,避免随机 I/O 带来的性能损耗,使磁盘写入效率非常高。
其次,Kafka 使用零拷贝(Zero Copy)技术,减少数据在用户态和内核态之间的复制,降低 CPU 消耗,提高网络传输效率。
同时,Kafka 支持批量发送和批量消费,将多条消息合并处理,减少网络请求和磁盘 I/O 次数,提高整体吞吐。
另外,Kafka 通过 Partition 分区机制实现水平扩展,不同 Partition 可以并行读写,充分利用多台机器和多核 CPU 的能力。
最后,Kafka 依赖操作系统 Page Cache 缓存热点数据,减少磁盘访问,提高读写性能。
总结来说,Kafka 高吞吐主要依靠:顺序写磁盘 + 零拷贝 + 批量处理 + 分区并行 + Page Cache 优化。这些设计让 Kafka 能够支撑百万级消息吞吐。
2、Kafka 为什么需要 Partition?
首先,Partition 可以实现并行处理。一个 Topic 可以拆分成多个 Partition,不同 Partition 可以分布在不同 Broker 上,由多个 Producer 和 Consumer 并行读写,从而提高系统整体吞吐量。
其次,Partition 提供了水平扩展能力。当数据量增加时,可以通过增加 Partition 数量和 Broker 节点来扩容,而不需要单机承担全部数据压力。
另外,Partition 通过 Offset 保证消息顺序。Kafka 只保证单个 Partition 内的消息有序,不保证多个 Partition 之间的全局顺序,这样可以在性能和顺序之间进行平衡。
3、kafka的ack机制
Kafka 的 ack 机制用于控制 Producer 发送消息后,Broker 需要多少副本确认成功后才认为消息发送成功,主要有三种级别:
acks=0
Producer 发送消息后不等待 Broker 响应,直接认为发送成功。性能最高,但可靠性最低,如果消息发送过程中 Broker 宕机,消息可能丢失。acks=1
Producer 发送消息后,只需要 Leader 副本写入成功并返回确认即可。性能和可靠性比较均衡,但如果 Leader 写入成功后还没同步给 Follower 就宕机,可能导致消息丢失。acks=all(或 -1)
Producer 发送消息后,需要 Leader 和所有 ISR(同步副本)都确认写入成功后才返回成功。可靠性最高,可以最大程度避免消息丢失,但吞吐量会有所下降。
4、Kafka 如何保证消息不丢失?
Producer 端通过设置 acks=all,要求所有 ISR(同步副本)都确认消息写入成功后才返回成功,同时可以开启 retries 重试机制,避免网络异常导致消息发送失败。
Broker 端通过 **副本机制(Replication)**保证数据可靠性。每个 Partition 会有多个副本,Leader 负责读写,Follower 负责同步数据。当 Leader 宕机时,可以从 ISR 副本中选举新的 Leader,避免数据丢失。同时可以合理配置 min.insync.replicas,保证至少多个副本同步成功。
Consumer 端通过手动提交 Offset避免消息丢失。消费者处理完消息后再提交 Offset,如果消费过程中失败,Offset 不会提前提交,重新消费时可以继续处理。
5、Kafka 如何保证 Exactly Once
Kafka 保证 Exactly Once(精确一次语义)主要依靠 幂等 Producer + 事务机制 + 消费端 Offset 管理。
首先,Kafka 通过 **幂等 Producer(enable.idempotence=true)**避免消息重复写入。Producer 会为每个消息分配唯一的 Producer ID(PID)和 Sequence Number,Broker 根据这些信息判断重复消息并丢弃,保证单个 Partition 内消息只写入一次。
其次,Kafka 通过 **事务机制(Transaction)**保证多步操作的原子性。例如 Producer 同时向多个 Partition 写消息,或者同时写消息和提交 Offset 时,可以通过事务保证要么全部成功,要么全部失败,避免部分成功导致数据不一致。
另外,在 Consumer 端通常采用 read_committed 模式读取事务消息,只消费已经提交的消息,同时配合 Offset 提交机制,保证消息处理和 Offset 更新的一致性。
例如 Kafka + Flink 场景中,Flink 会开启 Kafka 事务写入,并通过 Checkpoint 保存状态,当任务失败恢复时,可以回滚未提交的数据,从而实现端到端 Exactly Once。
6、Kafka 消息积压如何处理?
Kafka 消息积压定位首先通过 kafka-consumer-groups 查看 Consumer Lag,确认积压的 Topic 和 Partition;然后对比 Producer TPS 和 Consumer TPS,判断是生产过快还是消费能力不足;
如果是消费者处理能力不足,可以增加 Consumer 数量,提高并行度;同时增加 Topic 的 Partition 数量,让更多 Consumer 可以并行消费。另外可以优化消费逻辑,例如批量拉取消息、减少单条消息处理耗时、异步处理耗时任务。
如果是消费端故障或异常导致积压,需要检查消费者日志,解决异常后恢复消费;如果存在消息处理过慢,可以将耗时任务拆分,通过线程池、异步任务等方式提升吞吐。
7、Kafka Rebalance 为什么发生?
Kafka Rebalance 是消费者组重新分配 Partition 的过程,主要发生在 Consumer 加入/退出、心跳超时、消费处理超时、Partition 变化等场景。Rebalance 会导致短暂消费暂停,因此生产环境需要合理设置 session.timeout.ms、heartbeat.interval.ms、max.poll.interval.ms,并减少频繁 Rebalance。
8、Flink 架构是什么?
Flink 架构主要由 Client、Dispatcher、JobManager 和 TaskManager 组成。用户通过client提交任务 → Client 生成 JobGraph → Dispatcher 接收任务 → JobManager 调度任务、分配资源 → TaskManager 启动 Task 执行 → TaskManager 之间进行数据交换 → JobManager 通过 Checkpoint 实现状态管理和故障恢复。
9、Flink 为什么低延迟?
Flink 低延迟主要依靠 真正流式计算、事件驱动模型、Pipeline 执行、内存状态管理以及异步 Checkpoint 机制,避免了批处理等待和频繁磁盘访问,使 Flink 可以实现毫秒级实时计算。
10、Flink Checkpoint是什么?
