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

使用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在**批处理模式(有界流)和流处理模式(无界流)**下的水印生成逻辑差异导致的:

  1. 先插数据再启动任务:触发批处理模式
    当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的消息都被判定为迟到,无法输出。

  2. 先启动任务再插数据:触发流处理模式
    当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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 08:36:14