You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Flink SQL的Kafka源无数据时推进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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.01 22:55:55