尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

flink和kafka常见面试题

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 需要多少副本确认成功后才认为消息发送成功主要有三种级别acks0Producer 发送消息后不等待 Broker 响应直接认为发送成功。性能最高但可靠性最低如果消息发送过程中 Broker 宕机消息可能丢失。acks1Producer 发送消息后只需要 Leader 副本写入成功并返回确认即可。性能和可靠性比较均衡但如果 Leader 写入成功后还没同步给 Follower 就宕机可能导致消息丢失。acksall或 -1Producer 发送消息后需要 Leader 和所有 ISR同步副本都确认写入成功后才返回成功。可靠性最高可以最大程度避免消息丢失但吞吐量会有所下降。4、Kafka 如何保证消息不丢失Producer 端通过设置 acksall要求所有 ISR同步副本都确认消息写入成功后才返回成功同时可以开启 retries 重试机制避免网络异常导致消息发送失败。Broker 端通过 **副本机制Replication**保证数据可靠性。每个 Partition 会有多个副本Leader 负责读写Follower 负责同步数据。当 Leader 宕机时可以从 ISR 副本中选举新的 Leader避免数据丢失。同时可以合理配置 min.insync.replicas保证至少多个副本同步成功。Consumer 端通过手动提交 Offset避免消息丢失。消费者处理完消息后再提交 Offset如果消费过程中失败Offset 不会提前提交重新消费时可以继续处理。5、Kafka 如何保证 Exactly OnceKafka 保证 Exactly Once精确一次语义主要依靠 幂等 Producer 事务机制 消费端 Offset 管理。首先Kafka 通过 **幂等 Producerenable.idempotencetrue**避免消息重复写入。Producer 会为每个消息分配唯一的 Producer IDPID和 Sequence NumberBroker 根据这些信息判断重复消息并丢弃保证单个 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是什么
返回列表