Flink中如何直接获取DataStream全局总和并存储?无需遍历滚动结果
获取Flink DataStream全局最终总和的优化方案
给你两个更贴合需求的实现方式,分别适配不同的数据流场景:
1. 有界流场景:直接提取最终结果
如果你的数据流是有界的(比如示例里的固定元素集合,或者用户发送的是一批有限数据),完全不用遍历所有中间累加值,利用Flink对有界流的支持,直接取处理完成后的最后一个结果就行:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); DataStreamSource<Double> dataStream = env.fromElements(2.00, 3.00, 4.00, 11.00, 13.00, 14.00); // 全局键控做累加 SingleOutputStreamOperator<Double> sumStream = dataStream .keyBy(value -> "global-key") .reduce((a, b) -> a + b); // 执行并获取结果,迭代器会在流结束时自动终止 try (CloseableIterator<Double> iterator = sumStream.executeAndCollect()) { Double finalSum = null; while (iterator.hasNext()) { finalSum = iterator.next(); } System.out.println("最终总和:" + finalSum); // 输出47.0 } catch (Exception e) { e.printStackTrace(); }
这种写法不需要额外存中间结果,最后一次赋值的finalSum就是全局总和,因为有界流处理完所有数据后,迭代器就会停止。
2. 无界流场景:靠结束标记输出最终结果
如果你的数据流是无界的,但用户发送的数据有明确的结束信号(比如用-1.0表示这批数据发完了),可以用ProcessFunction维护一个全局累加状态,收到结束标记时才输出最终总和:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 假设数据流最后会传一个-1.0作为结束标记 DataStreamSource<Double> dataStream = env.fromElements(2.00, 3.00, 4.00, 11.00, 13.00, 14.00, -1.0); SingleOutputStreamOperator<Double> finalSumStream = dataStream .keyBy(value -> "global-key") .process(new ProcessFunction<Double, Double>() { // 定义全局累加状态 private ValueState<Double> sumState; @Override public void open(Configuration parameters) throws Exception { // 初始化状态,默认值0.0 ValueStateDescriptor<Double> descriptor = new ValueStateDescriptor<>( "sum-state", Double.class, 0.0 ); sumState = getRuntimeContext().getState(descriptor); } @Override public void processElement(Double value, Context ctx, Collector<Double> out) throws Exception { if (value == -1.0) { // 收到结束标记,输出当前累加的总和 out.collect(sumState.value()); } else { // 正常累加数值到状态里 sumState.update(sumState.value() + value); } } }); // 只取一次结果就行,因为只有结束标记会触发输出 try (CloseableIterator<Double> iterator = finalSumStream.executeAndCollect()) { if (iterator.hasNext()) { Double finalSum = iterator.next(); System.out.println("最终总和:" + finalSum); // 输出47.0 } } catch (Exception e) { e.printStackTrace(); }
这种方式只会输出一次最终结果,不会产生中间的累加值输出,完全符合你不想遍历流的需求。
补充说明
你之前用reduce得到的是滚动累加流,每来一条数据就输出当前的累加值,所以必须遍历所有中间值才能拿到最后结果。上面的两种方案要么利用有界流的特性直接取最终值,要么通过状态控制只输出一次结果,效率和简洁性都更好。
内容的提问来源于stack exchange,提问作者Muazim Wani
相关产品推荐
相关产品推荐

