Flink从FlinkKafkaConsumer迁移到KafkaSource后窗口不执行如何解决
问题根因
你遇到的问题是新版KafkaSource默认未开启分区空闲检测,多分区水位对齐逻辑导致全局水位无法推进,最终事件时间窗口无法触发计算。
旧版FlinkKafkaConsumer使用老的SourceFunction API,对空闲分区的处理逻辑更宽松,而新版基于FLIP-27实现的KafkaSource会严格对齐所有消费分区的水位:只要有任意一个分区没有新数据流入,该分区的水位会一直停留在Long.MIN_VALUE,全局水位就永远无法上涨超过窗口结束时间,窗口自然不会触发。
解决方法
任选一种方案即可:
方案1:将水位策略迁移到KafkaSource构造阶段(更推荐)
新KafkaSource支持按分区生成水位,天然匹配Kafka分区的事件时间进度,性能和准确性都优于在源后额外设置水位。修改代码如下:
KafkaSource<TestData> source = KafkaSource.<TestData>builder() .setBootstrapServers(propertiesForKafka.getProperty("bootstrap.servers")) .setTopics(TOPIC) .setDeserializer(new CustomDeserializer()) .setGroupId(GROUP_ID) .setStartingOffsets(OffsetsInitializer.earliest()) .build(); DataStream<TestData> testDataStreamSource = env.fromSource( source, WatermarkStrategy.<TestData>forMonotonousTimestamps() .withTimestampAssigner((event, timestamp) -> event.getTimestamp()) // 开启分区空闲检测,超过10秒无数据的分区会被排除在水位对齐逻辑外 .withIdleness(Duration.ofSeconds(10)), "Kafka Source" );
修改完成后删除原处理流水线中assignTimestampsAndWatermarks的代码,避免重复设置水位策略冲突。
方案2:在源后设置水位时添加空闲检测配置
如果业务逻辑需要在源算子之后才设置水位,直接在原流水线的水位策略中添加空闲检测即可:
testDataStreamSource .assignTimestampsAndWatermarks(WatermarkStrategy.<TestData>forMonotonousTimestamps() .withTimestampAssigner((event, timestamp) -> event.getTimestamp()) .withIdleness(Duration.ofSeconds(10))) .keyBy(TestData::getKey) .window(SlidingEventTimeWindows.of(Time.hours(3), Time.hours(1))) .process(new ProcessWindowFunction<TestData, TestDataOutput, String, TimeWindow>() { @Override public void process( .... });
内容的提问来源于stack exchange,提问作者Bagdus
相关产品推荐
相关产品推荐

