如何使用Flink Table API处理迟到事件?
Flink Table API处理迟到事件的可行方案
- 先明确现状:Flink Table API默认会丢弃迟到事件,但现在已经支持多种原生处理方式,不用再依赖早期的临时方案。
- 目前可用的靠谱方案:
- 给窗口设置允许延迟:定义窗口时,通过
WITH ( 'allowed-lateness' = 'XX' )参数配置允许的延迟时长,窗口在水位线超过结束时间后,会继续等你设定的时间来接收迟到事件,这段时间里窗口的计算结果会实时更新。 - 用侧输出流单独处理迟到事件:如果想把迟到事件捞出来单独处理,可以用窗口的
SIDE OUTPUT语法,把迟到事件输出到指定的侧输出表,示例SQL如下:-- 定义输入表,设置水位线 CREATE TABLE UserEvents ( user_id STRING, event_time TIMESTAMP(3), event_type STRING, WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'user_events', 'properties.bootstrap.servers' = 'localhost:9092', 'format' = 'json' ); -- 定义窗口聚合结果表 CREATE TABLE AggResults ( window_start TIMESTAMP(3), window_end TIMESTAMP(3), user_id STRING, event_count BIGINT ) WITH ( 'connector' = 'print' ); -- 定义迟到事件输出表 CREATE TABLE LateEventOutput ( user_id STRING, event_time TIMESTAMP(3), event_type STRING ) WITH ( 'connector' = 'print' ); -- 聚合并输出迟到事件到侧输出表 INSERT INTO AggResults SELECT TUMBLE_START(event_time, INTERVAL '10' MINUTE) AS window_start, TUMBLE_END(event_time, INTERVAL '10' MINUTE) AS window_end, user_id, COUNT(*) AS event_count FROM UserEvents GROUP BY TUMBLE(event_time, INTERVAL '10' MINUTE), user_id SIDE OUTPUT LATE EVENTS INTO LateEventOutput; - 关于你找到的3年前的方案:那个基于
current_watermark的方案现在已经没必要用了,Flink官方已经在Table API/SQL里内置了更直接的迟到事件处理能力,原生方案更稳定也更易维护。
- 给窗口设置允许延迟:定义窗口时,通过
内容的提问来源于stack exchange,提问作者Joseandro Luiz
相关产品推荐
相关产品推荐

