Flink 1.15.2水印异常无窗口输出,同任务1.13运行正常
Flink 1.15.2水印与窗口计算异常问题
- 任务从Flink 1.13迁移至1.15.2后,在CREATE TABLE DDL中为
time_ltz字段定义水印,执行1分钟滚动窗口的count(distinct userId)计算时始终无数据输出,相同任务在1.13版本可正常运行。 - 其他迁移至1.15.2的任务存在输出数据不匹配的情况,询问是否存在需要配置的水印默认设置。
相关DDL如下:
CREATE TABLE test ( eventName String, ingestion_time BIGINT, time_ltz AS TO_TIMESTAMP_LTZ(ingestion_time, 3), props ROW(userId VARCHAR, id VARCHAR, tourName VARCHAR, advertiserId VARCHAR, deviceId VARCHAR, tourId VARCHAR), WATERMARK FOR time_ltz AS time_ltz - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'test', 'scan.startup.mode' = 'latest-offset', 'properties.bootstrap.servers' = 'localhost:9092', 'properties.group.id' = 'local_test_flink_115', 'format' = 'json', 'json.ignore-parse-errors' = 'true', 'scan.topic-partition-discovery.interval' = '60000' );
附水印异常截图:
内容的提问来源于stack exchange,提问作者user9068199
相关产品推荐
相关产品推荐

