Flink Kafka Source Connector默认Watermark策略是什么?
Flink 1.11 Kafka消费者默认水位线策略说明
针对你提出的问题,结合Flink 1.11版本的实现逻辑,直接给出结论:
- 如果你使用的是1.11版本主流的
FlinkKafkaConsumer配合addSource的方式接入Kafka,且没有显式调用assignTimestampsAndWatermarks()方法传入自定义水位线策略,默认不会生成任何水位线。此时如果作业配置的时间特性为事件时间,所有依赖水位线触发的逻辑(比如窗口计算、事件时间定时器)都无法正常运行,因为作业没有水位线推进的信号。 - 即便没有配置水位线策略,默认还是会把Kafka消息自带的时间戳(对应ConsumerRecord中的timestamp字段,具体是生产端发送时间还是Broker写入时间由Kafka Topic的
message.timestamp.type配置决定)赋值给每条记录的事件时间字段,只是不会基于这个时间计算并下发水位线。 - 如果你使用的是1.11版本开始试点的新
fromSourceAPI接入Kafka,该API要求必须显式传入WatermarkStrategy参数,不传会直接抛出参数错误,不存在默认空值的情况。
补充注意:你在源码中看到的默认值为空的配置,就是对应旧版
FlinkKafkaConsumer的水位线策略字段,空值就代表不启用水位线生成逻辑。
内容的提问来源于stack exchange,提问作者bbaja42
相关产品推荐
相关产品推荐

