使用Kafka Connector与Flink SQL时水印行为的疑问
场景复现
我在使用Kafka Connector与Flink SQL时,对水印的分配逻辑存在疑问。已按以下顺序向Kafka主题插入消息:
1. A,2023-07-08 09:00:00,2023-07-08,1.0 2. A,2023-07-08 09:00:01,2023-07-08,2.0 3. A,2023-07-08 09:02:01,2023-07-08,2.0 4. A,2023-07-08 09:03:00,2023-07-08,4.0 5. A,2023-07-08 09:03:01,2023-07-08,5.0 6. A,2023-07-08 09:04:01,2023-07-08,6.0 7. A,2023-07-08 09:03:10,2023-07-08,5.35
创建的源表DDL:
CREATE TEMPORARY TABLE source_kafka_test ( `symbol` VARCHAR(2147483647), `msg_ts` TIMESTAMP(3), `trade_date` VARCHAR(2147483647), `trade` DOUBLE, WATERMARK FOR `msg_ts` AS `msg_ts` - INTERVAL '1' MINUTES ) WITH ( 'connector' = 'kafka', 'format' = 'csv', 'properties.bootstrap.servers' = '', 'topic' = '', 'properties.group.id' = '', 'scan.startup.mode' = 'earliest-offset' );
执行的查询语句:
create temporary view raw_data_sorted AS ( select symbol, msg_ts, trade_date, trade FROM source_kafka_test order by msg_ts ASC ); select * from raw_data_sorted;
两种执行时机的不同结果
- 先向Kafka插入所有数据,再启动Flink SQL查询
得到结果仅输出前3条数据:
symbol msg_ts trade_date trade A 2023-07-08T09:00 2023-07-08 1.0 A 2023-07-08T09:00:01 2023-07-08 2.0 A 2023-07-08T09:02:01 2023-07-08 2.0
推测水印直接取了Kafka中最后一条消息的时间2023-07-08 09:03:10,计算后水印值为2023-07-08 09:02:10,因此msg_ts晚于该时间的消息被判定为迟到事件,暂不输出。
- 先启动Flink SQL任务,再按顺序插入消息
所有7条数据都会被输出,此时水印值推进到2023-07-08 09:03:01(对应第6条消息的msg_ts减1分钟),符合预期。
核心疑问
为何水印值会因查询执行时机(先插数据还是先启动任务)不同而产生差异?这会导致二次执行查询时结果不确定,例如部分迟到事件被判定为非迟到并参与计算,或若Kafka最后一条消息为当日最后交易,窗口可能永远无法关闭(或延迟关闭)。
这本质是Flink在**批处理模式(有界流)和流处理模式(无界流)**下的水印生成逻辑差异导致的:
先插数据再启动任务:触发批处理模式
当Kafka主题中已有完整数据集,Flink将其视为有界流处理。此时水印生成逻辑会直接取整个数据集的最大事件时间,减去定义的延迟时间作为最终水印。
在你的场景中,数据集的最大msg_ts是第7条的2023-07-08 09:03:10,所以水印直接变为09:03:10 - 1min = 09:02:10。所有msg_ts <= 09:02:10的消息(前3条)会被输出,而msg_ts > 09:02:10的消息都被判定为迟到,无法输出。先启动任务再插数据:触发流处理模式
当Flink任务先启动,再逐步消费Kafka消息时,会以无界流模式运行。此时水印是随着消息的消费逐步推进的:
- 消费到前几条消息时,水印逐步提升;
- 消费到第6条
09:04:01的消息时,水印推进到09:03:01; - 后续消费到第7条
09:03:10的消息时,其msg_ts大于当前水印09:03:01,不属于迟到事件,会被正常输出。
解决方案
如果希望在有界流场景下也能得到和流处理一致的结果,可以通过以下配置强制Flink以流处理模式处理有界Kafka数据源:
在表的WITH参数中添加:
'scan.bounded.mode' = 'continuous-unbounded'
修改后的源表DDL:
CREATE TEMPORARY TABLE source_kafka_test ( `symbol` VARCHAR(2147483647), `msg_ts` TIMESTAMP(3), `trade_date` VARCHAR(2147483647), `trade` DOUBLE, WATERMARK FOR `msg_ts` AS `msg_ts` - INTERVAL '1' MINUTES ) WITH ( 'connector' = 'kafka', 'format' = 'csv', 'properties.bootstrap.servers' = '', 'topic' = '', 'properties.group.id' = '', 'scan.startup.mode' = 'earliest-offset', 'scan.bounded.mode' = 'continuous-unbounded' );
这样即使Kafka中已有预存数据,Flink也会按流处理的方式逐步推进水印,所有符合条件的消息都会被正常输出。
内容的提问来源于stack exchange,提问作者Ardelia Lortz

