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

Flink中KeyedProcessFunction获取currentWatermark始终为Long.MIN的问题

问题原因及解决方案

以下是几个可能导致水印始终返回Long.MIN_VALUE的常见原因及对应解决方法:

  • 误用批处理执行环境
    如果你初始化的是ExecutionEnvironment(批处理环境),事件时间与水印机制默认不会工作,水印会一直停留在初始值Long.MIN_VALUE。必须切换为流处理环境:

    StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
    
  • 测试数据存在问题

    • 若events集合为空,没有数据流入算子,水印无法被推进,保持初始值。确保集合中至少包含一条有效CallEvent数据。
    • 检查CallEvent.getEventTimeStamp()返回的时间戳是否合法(比如是大于0的毫秒级时间值),如果所有事件的时间戳本身就是Long.MIN_VALUE,水印自然无法更新。
  • 调用时机不正确
    如果在CallEventKeyedProcessFunction的open()方法中调用ctx.timerService().currentWatermark(),此时还没有任何数据被处理,水印尚未更新,返回值必然是Long.MIN_VALUE。需要在processElement()方法处理数据的逻辑中调用,此时水印已经根据流入的数据完成了计算与推进。

  • 强制设置为处理时间模式
    若代码中显式设置了处理时间模式,会直接禁用事件时间与水印机制:

    // 这种设置会导致水印失效,需要移除或修改
    env.setStreamTimeCharacteristic(TimeCharacteristic.ProcessingTime);
    

    Flink 1.12及以上版本默认使用事件时间,若使用旧版本,需显式设置为事件时间:

    env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
    

内容的提问来源于stack exchange,提问作者Andra Turc

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 12:45:38