Flink 流处理

共 24 题
#

1. Flink 的两阶段提交 Sink(TwoPhaseCommitSink)如何借助事务预提交保证端到端 Exactly-Once,对 Kafka 事务有何依赖

A 两阶段提交只在 Driver 端完成,与算子无关
B TwoPhaseCommitSink 在 checkpoint 时先预提交事务,checkpoint 成功后提交,失败回滚,依赖外部系统支持事务如 Kafka 事务性 producer ✓ 正确答案
C 无需任何外部系统支持即可实现端到端 Exactly-Once
D 预提交与提交是同一操作
#

2. Flink 与 Kafka 的精确一次集成

A 只需 sink 端事务即可,源端无需处理
B 消费端用 read_uncommitted 即可
C 精确一次与 checkpoint 无关
D 源端靠 checkpoint 记录 offset 防重读,sink 端用 Kafka 事务 producer 配合 checkpoint 原子提交,二者由 checkpoint 对齐实现端到端 Exactly-Once ✓ 正确答案
#

3. Flink Checkpoint 的 Barrier 机制、对齐(aligned)与不对齐(unaligned)检查点的取舍

A 不对齐 Checkpoint 快照体积更小
B 不对齐 Checkpoint 精度更高
C Barrier 对齐不需要等待任何输入
D 对齐是等待所有输入分支的 Barrier 后再快照,保证精确一次但可能受背压拖慢;不对齐 Checkpoint 跳过对齐、快照更大但降低延迟 ✓ 正确答案
#

4. Flink 的 Exactly-Once 与 Checkpoint

A Checkpoint 通过 Barrier 对齐做一致快照并持久化,恢复时从快照重放,保证引擎级精确一次;端到端还需源端可重放与 sink 事务/幂等配合 ✓ 正确答案
B Checkpoint 只保存运行时状态,不保存 source 位置
C 开启 Checkpoint 即自动保证端到端精确一次,无需 sink 配合
D Exactly-Once 与 Checkpoint 无关
#

5. Flink 状态 TTL 如何配置以清理过期状态,状态膨胀会带来哪些 checkpoint 与内存问题、如何治理

A StateTtlConfig 可配置状态存活时间并配合 Background/增量/压缩清理,状态膨胀会增大 checkpoint 体积与内存占用,TTL 是常用治理手段 ✓ 正确答案
B 状态膨胀只影响内存,不影响 checkpoint
C 状态 TTL 不能清理任何状态
D 状态 TTL 配置后立即清理全部状态
#

6. Flink 的 CEP(Complex Event Processing)

A CEP 只能处理单事件,不能处理序列
B CEP 用 Pattern API 定义事件序列模式(如 next/followedBy/次数),通过 CEP.pattern 匹配复杂事件,常用于实时告警与风控 ✓ 正确答案
C CEP 不需要时间语义
D CEP 与状态管理无关,无法恢复
#

7. Flink 的 FlinkML 机器学习

A FlinkML 只能做离线训练
B FlinkML 是功能最全的 ML 库,超过 MLlib
C FlinkML 提供表驱动的机器学习与在线学习能力,但算法覆盖远不如 Spark MLlib 成熟,工程上常与 MLlib 分工 ✓ 正确答案
D FlinkML 与 Table API 无关
#

8. Flink 的 Kryo 与 Avro 序列化

A Avro 不支持 schema 演进
B Kryo 是 Flink 性能最好的序列化器
C Flink 优先使用基于 TypeInformation 的内置序列化器,未识别类型回退 Kryo(通用但慢),与外部系统交换结构化数据常用 Avro(schema 驱动) ✓ 正确答案
D Flink 无法处理自定义类型
#

9. Flink 的 Metric 与监控(flink-metrics)

A Flink 无法上报自定义 metric
B Flink 提供 Counter/Gauge/Histogram/Meter 四类 metric,通过 reporter(如 Prometheus/JMX)上报,checkpoint 失败、背压、watermark 进度是核心监控指标 ✓ 正确答案
C 只有 Counter 一种 metric 类型
D 监控与 checkpoint 无关
#

10. Flink 的 ProcessFunction 与 Timer

A ProcessFunction 可访问元素、状态与时间戳,并通过 TimerService 注册事件/处理时间定时器,onTimer 回调用于超时检测、状态清理等场景 ✓ 正确答案
B ProcessFunction 无法访问状态
C Timer 只能注册一次,无法重复
D Timer 与 checkpoint 无关,无法恢复
#

11. Flink 的 Savepoint 与 Checkpoint 差异

A Checkpoint 自动周期触发用于故障恢复、生命周期短,Savepoint 手动触发用于升级/迁移/扩容、要求跨版本兼容 ✓ 正确答案
B 二者完全等价,无区别
C Savepoint 自动触发,Checkpoint 手动触发
D Checkpoint 要求跨版本兼容
#

12. Flink 的 Stream 与 Batch 模式

A Batch 模式延迟更低
B Streaming 模式无法处理有界数据
C 两种模式 API 完全不同,无法复用
D Streaming 处理无界流逐事件低延迟,Batch 处理有界数据一次性执行可重度优化,Flink 批流一体用同一套 API 支持两种模式 ✓ 正确答案
#

13. Flink 的 Table API 与 SQL

A SQL 无法表达流式查询
B Table API 只能用于批处理
C Table API/SQL 把流视为动态表,执行连续查询,属声明式高层 API,可通过 toDataStream/fromDataStream 与 DataStream API 互转 ✓ 正确答案
D Table 与 DataStream 无法互转
#

14. Flink 的 Time 语义(Event/Processing/Ingestion)

A Processing Time 最符合业务语义
B Event Time 无需处理乱序
C 三种时间语义完全相同
D Event Time 取事件源头时间最符合业务语义但需 watermark 处理乱序,Processing Time 取处理时刻最简单但不确定,Ingestion Time 由源注入 ✓ 正确答案
#

15. Flink 的 Watermark 机制

A Watermark 与乱序无关,用于排序
B watermark 越大越早触发窗口
C 迟到数据无法再处理
D Watermark 表示"时间戳小于等于它的数据已到齐"的进度,用于触发 event-time 窗口,配合 allowedLateness 与 side output 处理迟到数据 ✓ 正确答案
#

16. Flink 的反压(Backpressure)机制

A 反压是 Flink 的错误,需要报错停止
B 反压只发生在源端,与算子无关
C 反压无法被检测
D 反压是下游处理慢导致上游缓冲被占满、压力沿链路反传的流控机制,可用 UI 背压面板定位瓶颈算子并针对性优化 ✓ 正确答案
#

17. Flink 窗口的触发器(Trigger)与增量聚合(AggregateFunction)在实时指标计算中的应用

A Trigger 决定窗口何时触发输出,AggregateFunction 在窗口内增量聚合只维护中间结果,内存省、性能高,适合实时指标计算 ✓ 正确答案
B 增量聚合需缓存窗口内全部数据
C Trigger 决定聚合方式,AggregateFunction 决定触发时机
D 窗口计算必须等待窗口结束才输出,无法增量
#

18. Flink CDC(Debezium 连接器)在实时入湖入仓中的应用与快照/增量切换

A Flink CDC 只能做全量同步
B Flink CDC 基于 Debezium 读取 binlog 捕获变更,先全量快照再自动切换增量 binlog,配合 checkpoint 断点续传实现实时同步 ✓ 正确答案
C CDC 需要业务表加额外字段
D CDC 不能处理增量变更
#

19. Flink on Kubernetes 的 Native 与 Standalone 部署模式有何差异,JobManager 高可用如何借助 K8s ConfigMap 与 ZooKeeper 实现

A Native 模式用 K8s API 动态管理 TaskManager,Standalone 由 Flink 自管;JobManager HA 可借助 K8s ConfigMap 存元数据选主或 ZooKeeper 选主 ✓ 正确答案
B 两种模式完全相同
C Standalone 模式深度使用 K8s 调度 Native 资源
D K8s 部署无需 HA 机制
#

20. Flink 的核心架构(JobManager/TaskManager)

A TaskManager 负责调度与切分 Stage
B JobManager 不参与 checkpoint
C 每个 TaskManager 只能运行一个算子
D JobManager 负责调度、checkpoint 协调与资源分配,TaskManager 负责执行算子与维护状态,Slot 是资源分配的基本单位 ✓ 正确答案
#

21. Flink 的状态后端(MemoryStateBackend/FsStateBackend/RocksDBStateBackend)

A RocksDBStateBackend 把状态存于本地磁盘可支撑超大状态,Memory/Fs 存于堆内存适合小状态,checkpoint 由后端决定持久化方式 ✓ 正确答案
B MemoryStateBackend 适合超大状态
C RocksDBStateBackend 状态存于 JVM 堆
D 状态后端不影响 checkpoint 方式
#

22. Flink 的窗口(Window)类型(滚动/滑动/会话)

A 滑动窗口永远不会重叠
B 滚动窗口固定大小不重叠,滑动窗口固定大小可重叠,会话窗口按不活跃间隔切分、长度不固定,适合用户会话分析 ✓ 正确答案
C 会话窗口长度固定
D 三种窗口都无法处理时间
#

23. Flink 窗口如何处理迟到数据,allowedLateness 与 sideOutputLateData 的协作机制及数据丢弃时机是怎样的

A 迟到数据一律丢弃,无法处理
B allowedLateness 允许窗口在结束后再保留一段时间并再次触发计算,超出该时间仍迟到的数据可经 sideOutputLateData 输出旁路,否则被丢弃 ✓ 正确答案
C allowedLateness 会让窗口无限期保留
D side output 是迟到数据的主要处理通道,与 allowedLateness 无关
#

24. Flink Table/SQL 与 DataStream API 的互操作(toDataStream/fromDataStream)

A Table 与 DataStream 无法互转
B 转换会丢失时间属性,无法保留 watermark
C toDataStream 只能得到 changelog 流
D fromDataStream 把流转为表,toDataStream 转追加流、toChangelogStream 转 changelog 流,需注意表类型与时间属性匹配 ✓ 正确答案