Flink中Reduce后Sink的invoke方法未触发问题咨询
问题分析与解决:Flink Sink invoke方法未被调用
问题根源
你代码中使用了GlobalWindows.create(),这种全局窗口默认没有触发器(Trigger),窗口会一直持续收集数据,永远不会触发计算和输出逻辑。这就导致Reduce后的结果无法传递到Sink,invoke方法自然不会被执行,连print()也不会有输出。
解决方案
给GlobalWindow手动添加触发器,或者替换为自带默认触发器的窗口类型,比如滚动窗口、滑动窗口等,以下是两种可行方案:
方案1:给GlobalWindow添加CountTrigger
指定窗口内元素达到一定数量时触发计算,示例中设置为每1个元素触发(结合你keyBy每个整数的逻辑,每个key对应一个元素,reduce操作实际不会改变数值,但能触发输出):
Configuration configuration = new Configuration(); StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(configuration); env.setBufferTimeout(0); List<Integer> counts = new ArrayList<>(); for (int i = 1; i <= 100; i++) { counts.add(i); } SingleOutputStreamOperator<Integer> value = env.fromCollection(counts) .keyBy(c -> c) .windowAll(GlobalWindows.create()) // 添加CountTrigger,元素数量达到1时触发窗口计算 .trigger(CountTrigger.of(1)) .reduce(new ReduceFunction<Integer>() { @Override public Integer reduce(Integer count, Integer t1) throws Exception { return count + t1; } }); value.addSink(new SinkFunction<Integer>() { @Override public void invoke(Integer value, Context context) throws Exception { log.info("sink count:" + value); } }); value.print(); env.execute("count");
方案2:替换为滚动时间窗口(更符合常规场景)
如果你的需求是基于时间或批量的窗口计算,直接替换GlobalWindow为滚动处理时间窗口,这类窗口自带默认触发器,会自动触发计算:
Configuration configuration = new Configuration(); StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(configuration); env.setBufferTimeout(0); List<Integer> counts = new ArrayList<>(); for (int i = 1; i <= 100; i++) { counts.add(i); } SingleOutputStreamOperator<Integer> value = env.fromCollection(counts) .keyBy(c -> c) // 使用1秒滚动处理时间窗口,足够收集所有批量生成的数据 .window(TumblingProcessingTimeWindow.of(Time.seconds(1))) .reduce(new ReduceFunction<Integer>() { @Override public Integer reduce(Integer count, Integer t1) throws Exception { return count + t1; } }); value.addSink(new SinkFunction<Integer>() { @Override public void invoke(Integer value, Context context) throws Exception { log.info("sink count:" + value); } }); value.print(); env.execute("count");
补充说明
- GlobalWindow仅适用于需要自定义触发逻辑的特殊场景,默认不会自动触发窗口计算;
- 时间类窗口(滚动、滑动、会话)自带默认触发器,会基于事件时间或处理时间自动触发计算和输出。
内容的提问来源于stack exchange,提问作者Lson
相关产品推荐
相关产品推荐

