如何在Flink SQL的Kafka源无数据时推进Watermark?
Flink 1.17 SQL流式作业:无数据时Watermark无法推进问题
问题背景
基于Flink 1.17开发SQL流式作业,以Kafka为数据源,核心目标是计算amountUSD字段的最近1小时总和。已配置多项Watermark相关参数,但无新数据输入时,Watermark会停滞在最后一条事件的时间戳,无法正确更新最近1小时的总量;有新交易数据时,作业能正常输出预期结果。
作业代码
Java 主程序代码
public class StreamingJob { // ... public static void main(String[] args) throws Exception { // ... List<String> sqlStatements = ... final EnvironmentSettings settings = EnvironmentSettings.newInstance().inStreamingMode().build(); final StreamExecutionEnvironment streamEnv = StreamExecutionEnvironment.getExecutionEnvironment(); streamEnv.getConfig().setAutoWatermarkInterval(1000); streamEnv.getConfig().setLatencyTrackingInterval(10000); tableEnv.getConfig().getConfiguration().setString("pipeline.auto-watermark-interval", "5000ms"); tableEnv.getConfig().getConfiguration().setString("table.exec.source.idle-timeout", "5000ms"); tableEnv.getConfig().getConfiguration().setString("table.exec.emit.allow-lateness", "1d"); for (String statement : sqlStatements) { tableEnv.executeSql(statement); } streamEnv.execute(jobId); }
SQL 核心逻辑代码
CREATE TABLE trades_stream_by_hour ( `entityId` STRING, `poolId` STRING, `amountUSD` FLOAT, `blockNumber` BIGINT, `timestamp` BIGINT, `rowtime` as TO_TIMESTAMP(FROM_UNIXTIME(`timestamp`)), WATERMARK FOR `rowtime` AS `rowtime` - INTERVAL '5' SECONDS ) WITH ( 'connector' = 'kafka', 'topic' = 'trades', 'scan.startup.mode' = 'timestamp', 'scan.startup.timestamp-millis' = '1705851977000', -- 2小时前的 Unix 时间戳(毫秒) 'format' = 'debezium-json', 'debezium-json.schema-include' = 'true', 'debezium-json.map-null-key.mode' = 'DROP', 'debezium-json.encode.decimal-as-plain-number' = 'true' ); -- 省略 "my_sink" 表定义(JDBC 输出表) -- 省略 "pools" 表定义(JDBC 维表) CREATE VIEW last_1_hour_volumes AS SELECT *, PROCTIME() as proc_time FROM ( SELECT *, ROW_NUMBER() OVER ( PARTITION BY poolId ORDER BY window_time DESC ) AS rn FROM ( SELECT poolId, SUM(COALESCE(amountUSD, 0)) AS totalVolumeUSD, COUNT(*) AS totalVolumeTrades, MAX(`timestamp`) AS `timestamp`, MAX(`blockNumber`) AS `maxBlockNumber`, MIN(`blockNumber`) AS `minBlockNumber`, HOP_ROWTIME(`rowtime`, INTERVAL '1' MINUTE, INTERVAL '1' HOUR) as window_time FROM trades_stream_by_hour WHERE amountUSD IS NOT NULL AND amountUSD > 0 AND `rowtime` >= CAST( CURRENT_TIMESTAMP - INTERVAL '1' HOUR AS TIMESTAMP(3) ) GROUP BY poolId, HOP(`rowtime`, INTERVAL '1' MINUTE, INTERVAL '1' HOUR) ) ) WHERE rn = 1; CREATE VIEW prev_1_hour_volumes AS SELECT *, PROCTIME() as proc_time FROM ( SELECT *, ROW_NUMBER() OVER ( PARTITION BY poolId ORDER BY window_time DESC ) AS rn FROM ( SELECT poolId, SUM(COALESCE(amountUSD, 0)) AS totalVolumeUSD, COUNT(*) AS totalVolumeSwaps, MAX(`timestamp`) AS `timestamp`, MAX(`blockNumber`) AS `maxBlockNumber`, MIN(`blockNumber`) AS `minBlockNumber`, HOP_ROWTIME(`rowtime`, INTERVAL '5' MINUTE, INTERVAL '1' HOUR) as window_time FROM trades_stream_by_hour WHERE amountUSD IS NOT NULL AND amountUSD > 0 AND `rowtime` >= CAST( CURRENT_TIMESTAMP - INTERVAL '2' HOUR AS TIMESTAMP(3) ) GROUP BY poolId, HOP(`rowtime`, INTERVAL '5' MINUTE, INTERVAL '1' HOUR) ) ) WHERE rn = 12; INSERT INTO my_sink SELECT COALESCE(lv1h.poolId, pv1h.poolId) as entityId, lv1h.totalVolumeSwaps as volumeSwaps1h, lv1h.minBlockNumber as volumeMinBlock1h, lv1h.maxBlockNumber as volumeMaxBlock1h, lv1h.totalVolumeUSD as volumeUSD1h, lv1h.`timestamp` as volumeUSDTimestamp1h, pv1h.totalVolumeSwaps as prevVolumeSwaps1h, pv1h.minBlockNumber as prevVolumeMinBlock1h, pv1h.maxBlockNumber as prevVolumeMaxBlock1h, pv1h.totalVolumeUSD as prevVolumeUSD1h, pv1h.`timestamp` as prevVolumeUSDTimestamp1h, lv1h.totalVolumeUSD - pv1h.totalVolumeUSD as volumeUSDChange1h, (lv1h.totalVolumeUSD - pv1h.totalVolumeUSD) / pv1h.totalVolumeUSD as volumeUSDChangePercent1h FROM last_1_hour_volumes lv1h LEFT JOIN prev_1_hour_volumes AS pv1h ON lv1h.poolId = pv1h.poolId LEFT JOIN pools_store FOR SYSTEM_TIME AS OF lv1h.proc_time AS pool ON pool.entityId = lv1h.poolId WHERE -- 过滤掉1小时前的无效数据 (CURRENT_WATERMARK(lv1h.window_time) IS NOT NULL AND CURRENT_WATERMARK(lv1h.window_time) >= CAST(CURRENT_TIMESTAMP - INTERVAL '5' MINUTES AS TIMESTAMP(3))) OR (CURRENT_WATERMARK(pv1h.window_time) IS NOT NULL AND CURRENT_WATERMARK(pv1h.window_time) >= CAST(CURRENT_TIMESTAMP - INTERVAL '5' MINUTES AS TIMESTAMP(3)))
已尝试的无效配置
- 将
setAutoWatermarkInterval()设为5秒 - 将
pipeline.auto-watermark-interval设为5秒 - 将
table.exec.source.idle-timeout设为5秒 - 将
table.exec.emit.allow-lateness设为5秒 - 将作业并行度与Kafka分区数均设为1
内容的提问来源于stack exchange,提问作者aram.eth
相关产品推荐
相关产品推荐

