考察点
时间语义是流处理的根基问题,面试官想确认你是否真正理解「乱序」这件事:为什么需要事件时间、Watermark 到底在承诺什么、多并行度下 Watermark 怎么传播。追问会往 Watermark 生成策略、空闲分区、迟到数据走。
参考答案
三种时间语义
处理时间(Processing Time)是算子所在机器的系统时钟,最简单、延迟最低,但结果不可重放——同一份数据跑两次,窗口划分可能不一样,故障恢复后数据会落到不同窗口。摄入时间(Ingestion Time)是数据进入 Flink source 时的时间戳,比处理时间稳定一点,但表达不了「事件真实发生的时刻」。事件时间(Event Time)是数据里自带的业务时间戳,比如日志里的 event_time,结果确定、可重放,是严肃实时计算的默认选择。
一句话总结:用处理时间换简单,用事件时间换正确。真实项目里只有监控类的粗略指标才敢用处理时间。
Watermark 是什么
事件时间的麻烦在于乱序:数据经过网络、Kafka 分区,到达顺序和发生顺序不一致。窗口总不能无限等下去,Watermark 就是引擎对「时间进度」的声明:当前 Watermark 为 T,意味着「我认为时间戳小于等于 T 的数据已经全部到齐,之后的都算迟到」。窗口结束时间小于等于 Watermark 时就触发计算。
要诚实地说:Watermark 是个启发式估计,不是真理。设得太激进,大量数据被判迟到;设得太保守,窗口输出延迟变大。这是延迟和完整性的权衡。
生成与传播
Watermark 在 source 处生成(1.11 之后在 SourceFunction/新的 FLIP-27 Source 上配 WatermarkStrategy),常用的是有界乱序策略:forBoundedOutOfOrderness(Duration.ofSeconds(10)),含义是「我最多容忍 10 秒乱序」,实际 Watermark = 观察到的最大时间戳 - 10 秒。还有单调递增策略用于基本有序的场景。
传播规则是面试必考点:Watermark 以广播形式跟着数据流往下游走;多并行度、多输入的场景(比如 keyBy 之后、union、双流 Join),算子的 Watermark 取所有输入分区的最小值。这带来一个经典坑:Kafka 某个分区没数据,它的 Watermark 永远不前进,拖住全局 Watermark,窗口不触发。解法是 withIdleness(Duration...) 标记空闲分区,空闲分区不参与最小值计算。
迟到数据怎么办
过了 Watermark 才到的数据,窗口默认直接丢弃。可以配 allowedLateness 给窗口一个宽限期,宽限期内每来一条就重触发更新结果;再晚的进侧输出流(sideOutputLateData),落到旁路存储做补偿或审计。三层兜底:Watermark 容忍常规乱序,allowedLateness 容忍异常迟到,侧输出兜底。
可能的追问
- Watermark 设多少合适?——看数据乱序的实际分布,通常采集端打点统计 P99 乱序时长再加余量;广告日志类场景常见 5-30 秒,跨机房采集可能分钟级。
- 为什么全局 Watermark 取最小值?——保证语义正确:任何一个输入分区还没推进到 T,就不能断言 T 之前的数据到齐了,取最小是保守且正确的选择。
- 事件时间 + 事件时间窗口在 Kafka 回溯重跑时要注意什么?——重放历史数据时事件时间照常推进,行为一致,这正是事件时间的价值;但要确认下游 sink 幂等,因为窗口结果会再写一遍。