无Key的Flink Event Time窗口未关闭:并行度参数作用咨询
从分区数为1的Kafka主题读取数据,使用无Key的Event Time窗口时,代码无法输出(关闭)窗口结果。添加env.setParallelism(1);后功能恢复正常,存在以下疑问:
- 该参数在我的场景中为何是必需的?未设置时窗口为何无法关闭?
- 文档提到无Key窗口的并发度始终为1,此现象如何解释?
- 使用
TumblingProcessingTimeWindows时,无论是否设置该参数都能正常运行,原因是什么?
附相关代码
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); KafkaSource<UserModel> source = KafkaSource.<UserModel>builder() .setBootstrapServers(kafka) .setTopics("f1") .setGroupId("flink_group") .setStartingOffsets(OffsetsInitializer.earliest()) .setValueOnlyDeserializer(new JsonConverter()) .build(); DataStream<UserModel> ds = env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka Source"); WatermarkStrategy<UserModel> strategy = WatermarkStrategy.<UserModel>forBoundedOutOfOrderness(Duration.ofSeconds(20)) .withTimestampAssigner((i, timestamp) -> { return i.dt.toInstant(ZoneOffset.UTC).toEpochMilli(); }); SingleOutputStreamOperator<UserModelEx> reduce = ds.assignTimestampsAndWatermarks(strategy) .windowAll(TumblingEventTimeWindows.of(Time.seconds(10))) .reduce((acc, i) -> { acc.count += i.count; acc.dt = i.dt; System.out.println(acc.dt + " reduce:" + acc.count); return acc; }, new Rich()); reduce.print();
1. setParallelism(1)的必要性及未设置时窗口无法关闭的原因
Flink默认并行度一般大于1(通常等于CPU核心数),虽然windowAll作为全局窗口本身并发度固定为1,但上游的Kafka Source和assignTimestampsAndWatermarks算子会以默认并行度运行。
Event Time窗口的触发依赖全局水印的推进,而全局水印取的是所有并行子任务水印的最小值。你的Kafka主题只有1个分区,Kafka Source的多个并行子任务里,只有1个能消费到数据并推进水印,其他子任务因无数据流入,水印会一直停留在初始最小值。这就导致全局水印永远无法推进到窗口结束时间+乱序容忍时间的阈值,窗口自然无法关闭触发。
设置并行度为1后,所有算子只有1个并行子任务,水印能正常根据流入数据推进,窗口就能按时触发。
2. 无Key窗口并发度始终为1的文档说明与现象的解释
文档指的是**windowAll算子自身的并发度固定为1**,但上游算子(如Kafka Source、assignTimestampsAndWatermarks)的并行度不受此限制。你遇到的问题并非windowAll本身的并发问题,而是上游多并行子任务导致的水印生成不一致:部分子任务无数据无法推进水印,拖慢了全局水印的更新,最终卡住了窗口的触发逻辑。
3. TumblingProcessingTimeWindows不受并行度影响的原因
Processing Time窗口的触发逻辑依赖系统时间,不需要水印驱动。每个windowAll子任务会独立根据自身的系统时间判断窗口是否到结束时间,只要系统时间到达窗口结束点,就会关闭窗口并输出结果,和上游并行度、水印状态完全无关。因此无论并行度设为多少,Processing Time窗口都能正常运行。
内容的提问来源于stack exchange,提问作者padavan

