第 197 题:延迟数据的处理,Side Output与Allowed Lateness?
题目
延迟数据的处理,Side Output与Allowed Lateness?
完整讲解
一、延迟数据的问题
事件时间下,数据可能晚于其事件时间到达(网络、乱序、重试),即延迟数据(Late Data)。若窗口已按 watermark 触发并输出,延迟数据若直接丢弃会漏算;若重新触发会重复输出或需更新已有结果,需明确策略。
二、Allowed Lateness
- 含义:在 watermark 超过窗口 end 后,仍允许一段时间内到达的、属于该窗口的延迟数据参与计算。在这段「允许延迟」内,每来一条延迟数据可再次触发窗口输出(即 late output / 延迟输出),或更新侧输出。
- 实现:窗口状态在「允许延迟」期内保留,延迟数据到达时重新聚合或触发;超过允许延迟的数据可进 Side Output 或丢弃。允许延迟越长,状态保留越久、资源占用越大,需与业务 SLA 权衡。
三、Side Output
- 含义:将不符合主输出路径的数据(如超过 Allowed Lateness 的延迟数据、异常数据)单独输出到一条或多条侧流,便于后续补算、监控或人工处理。
- 用法:在 ProcessFunction 或 WindowFunction 中,对「过晚」或异常记录
ctx.output(sideOutputTag, record);主流继续正常输出,侧流可接单独 sink 或下游作业做补偿。
四、小结
- Allowed Lateness:在窗口结束后仍接受延迟数据并产生 late 输出,平衡完整性与状态成本。Side Output:把无法纳入主流的数据导出,用于补偿或审计。二者常配合:超过允许延迟的进 Side Output。
面试要点
- 能说明延迟数据的来源与影响(漏算、重复);能解释 Allowed Lateness 的含义与 late 输出。
- 能说清 Side Output 的用途(延迟/异常数据单独输出)及与主流的配合。
- 能简述状态保留时间与允许延迟的权衡;能写出 Flink 中 Side Output 的用法。
记忆要点
- 延迟数据:晚于 watermark 到达;Allowed Lateness=窗口结束后仍接受延迟并产生 late 输出。
- Side Output=非主路径数据单独输出(过晚数据、异常);二者配合,超允许延迟进 Side Output。