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

Flink Kafka Source Connector默认Watermark策略是什么?

针对你提出的问题,结合Flink 1.11版本的实现逻辑,直接给出结论:

  • 如果你使用的是1.11版本主流的FlinkKafkaConsumer配合addSource的方式接入Kafka,且没有显式调用assignTimestampsAndWatermarks()方法传入自定义水位线策略,默认不会生成任何水位线。此时如果作业配置的时间特性为事件时间,所有依赖水位线触发的逻辑(比如窗口计算、事件时间定时器)都无法正常运行,因为作业没有水位线推进的信号。
  • 即便没有配置水位线策略,默认还是会把Kafka消息自带的时间戳(对应ConsumerRecord中的timestamp字段,具体是生产端发送时间还是Broker写入时间由Kafka Topic的message.timestamp.type配置决定)赋值给每条记录的事件时间字段,只是不会基于这个时间计算并下发水位线。
  • 如果你使用的是1.11版本开始试点的新fromSource API接入Kafka,该API要求必须显式传入WatermarkStrategy参数,不传会直接抛出参数错误,不存在默认空值的情况。

补充注意:你在源码中看到的默认值为空的配置,就是对应旧版FlinkKafkaConsumer的水位线策略字段,空值就代表不启用水位线生成逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 18:36:02