You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 12:31:00