当前位置: 首页 > news >正文

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 需要多少副本确认成功后才认为消息发送成功,主要有三种级别:

  1. acks=0
    Producer 发送消息后不等待 Broker 响应,直接认为发送成功。性能最高,但可靠性最低,如果消息发送过程中 Broker 宕机,消息可能丢失。

  2. acks=1
    Producer 发送消息后,只需要 Leader 副本写入成功并返回确认即可。性能和可靠性比较均衡,但如果 Leader 写入成功后还没同步给 Follower 就宕机,可能导致消息丢失。

  3. 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是什么?

http://www.cnnetsun.cn/news/3824202.html

相关文章:

  • 3分钟完成网易云音乐插件革命!BetterNCM安装器完全使用指南
  • Charles抓包进阶:断点与重写实战,掌握HTTP/HTTPS流量拦截与修改
  • 随机森林算法详解——基于垃圾邮件分类案例
  • Spring Security自定义认证:从默认密码到UserDetailsService实现详解
  • 从微软XBOX裁员看技术组织架构优化:扁平化与高效协作实践
  • Unity游戏开发实战:点乘与叉乘的五大核心应用场景解析
  • 本地大模型部署实战:从CUDA环境搭建到SenseNova-U1部署全流程
  • 年薪百万的网安人,到底比你强在哪?这份「从0到1」的系统学习地图,请查收。
  • 构建AI隐私安全舱:三层纵深防御体系护航企业智能数据安全
  • 统信UOS系统安装NVIDIA官方驱动实现稳定多屏显示完整指南
  • Android编译优化:精准模块清理解决增量编译失败
  • WorkBuddy与罗与罗Skill:构建法律AI智能体工作流的工程实践
  • MySQL字符集与校对规则深度解析:从原理到实战,彻底解决乱码问题
  • Loop Engineering 没死,Graph Engineering 也没有上位
  • 从拳击手到AI金融科技创业者:蔡永军的跨界转型之路
  • 信奥赛01串问题解析:位运算与动态规划实战
  • OpenClaw AI Agent 实战:从部署到技能开发的完整指南
  • 2026年正规SEO公司怎么选:七大避坑维度+真实案例复盘+KPI对赌合同指南|详解
  • 2026年正规SEO公司怎么选:七大避坑维度+真实案例复盘+KPI对赌合同指南|指南
  • 2025最权威的十大降AI率神器横评
  • 【读论文】2020 IEEE [C] 多种基音检测算法对比研究 A comparative study of various pitch detection algorithms
  • 关于编译器报警告--scanf的返回值被忽略-程序却能正常运行的理解
  • React useState初始值写法性能优化指南
  • Kali Linux部署HexStrike AI:MCP连接失败深度排错与优化指南
  • CTFHub HTTP协议通关指南:从基础请求到实战技巧
  • 支持私有化部署的企业 Agent 方案选型指南:技术架构、安全边界与主流厂商深度测评
  • Unity Cinemachine Virtual Camera:从核心原理到第三人称镜头实战
  • 虚拟仿真、半实物仿真和实况仿真简介
  • OpenCV相机标定实战:从针孔模型到鱼眼矫正的完整指南
  • UE5 Nanite实战指南:从核心原理到资产分类启用策略