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

如何用Flink实现实时过滤器并计算键值与总和的比率?

解决Flink实时计算键值占全局总和比率的方案

Hey, great question! This is a classic challenge with state sharing in Flink's keyed operators—since each keyed operator instance only holds state for its assigned keys, you can't directly compute the global sum across all keys. But we've got two solid approaches to solve this, tailored to your example data and requirements (including filtering out values of 0).

方案1:广播状态(Broadcast State)实现并行计算

这个方案适合数据量较大、需要并行处理的场景,核心思路是:

  • 先计算全局实时总和并广播给所有Keyed算子
  • 每个Keyed算子维护自身键的累加值,结合广播的全局总和计算比率

代码实现(Java)

首先定义输入数据的POJO:

public class Record {
    private String key;
    private int value;

    // 构造方法、getter、setter
    public Record(String key, int value) {
        this.key = key;
        this.value = value;
    }

    public String getKey() { return key; }
    public int getValue() { return value; }
}

然后是主处理流程:

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

// 1. 输入流:先过滤掉值为0的记录
DataStream<Record> inputStream = env.fromElements(
        new Record("k1", 1),
        new Record("k2", 3),
        new Record("k1", 1),
        new Record("k2", 5),
        new Record("k3", 0)
).filter(record -> record.getValue() != 0);

// 2. 计算全局实时总和流:用全局滚动窗口实时更新总和
DataStream<Integer> totalSumStream = inputStream
        .map(Record::getValue)
        .windowAll(TumblingProcessingTimeWindows.of(Time.milliseconds(1)))
        .sum(0);

// 3. 定义广播状态描述符,用于传递全局总和
MapStateDescriptor<String, Integer> broadcastStateDesc = new MapStateDescriptor<>(
        "total-sum-state",
        BasicTypeInfo.STRING_TYPE_INFO,
        BasicTypeInfo.INT_TYPE_INFO
);

// 4. 将全局总和流广播出去
BroadcastStream<Integer> broadcastSumStream = totalSumStream
        .map(sum -> Tuple2.of("total", sum))
        .broadcast(broadcastStateDesc);

// 5. Keyed流处理:维护键的累加值,结合广播总和计算比率
inputStream
        .keyBy(Record::getKey)
        .connect(broadcastSumStream)
        .process(new KeyedBroadcastProcessFunction<String, Record, Tuple2<String, Integer>, Tuple2<String, Double>>() {

            // 维护当前键的累加状态
            private ValueState<Integer> keyAccumulator;

            @Override
            public void open(Configuration parameters) throws Exception {
                keyAccumulator = getRuntimeContext().getState(new ValueStateDescriptor<>(
                        "key-accumulator",
                        Integer.class,
                        0
                ));
            }

            @Override
            public void processElement(Record value, ReadOnlyContext ctx, Collector<Tuple2<String, Double>> out) throws Exception {
                // 更新当前键的累加值
                int currentAccum = keyAccumulator.value() + value.getValue();
                keyAccumulator.update(currentAccum);

                // 从广播状态获取最新全局总和
                Integer totalSum = ctx.getBroadcastState(broadcastStateDesc).get("total");
                if (totalSum != null && totalSum > 0) {
                    double ratio = (double) currentAccum / totalSum;
                    out.collect(Tuple2.of(value.getKey(), ratio));
                }
            }

            @Override
            public void processBroadcastElement(Tuple2<String, Integer> value, Context ctx, Collector<Tuple2<String, Double>> out) throws Exception {
                // 更新广播状态中的全局总和
                ctx.getBroadcastState(broadcastStateDesc).put(value.f0, value.f1);
            }
        })
        .print();

env.execute("Real-Time Key Ratio Calculator");

关键说明

  • 全局窗口的选择:这里用TumblingProcessingTimeWindows是为了实时触发总和更新,你也可以用GlobalWindow结合自定义Trigger实现更高效的触发逻辑
  • 广播状态的作用:确保所有Keyed算子实例都能拿到最新的全局总和,解决状态隔离问题
  • 过滤逻辑:提前过滤值为0的记录,避免影响总和计算和比率结果

方案2:全局ProcessFunction实现简单聚合

如果你的数据量不大,不需要高并行度,这个方案更简洁——直接在全局ProcessFunction中维护所有键的累加状态和全局总和。

代码实现(Java)

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1); // 全局状态建议单并行度,避免分布式状态冲突

// 输入流:过滤值为0的记录
DataStream<Record> inputStream = env.fromElements(
        new Record("k1", 1),
        new Record("k2", 3),
        new Record("k1", 1),
        new Record("k2", 5),
        new Record("k3", 0)
).filter(record -> record.getValue() != 0);

inputStream
        .process(new ProcessFunction<Record, Tuple2<String, Double>>() {

            // 维护所有键的累加状态
            private MapState<String, Integer> keyAccumulators;
            // 维护全局总和状态
            private ValueState<Integer> totalSum;

            @Override
            public void open(Configuration parameters) throws Exception {
                keyAccumulators = getRuntimeContext().getMapState(new MapStateDescriptor<>(
                        "key-accumulators",
                        String.class,
                        Integer.class
                ));
                totalSum = getRuntimeContext().getState(new ValueStateDescriptor<>(
                        "total-sum",
                        Integer.class,
                        0
                ));
            }

            @Override
            public void processElement(Record value, Context ctx, Collector<Tuple2<String, Double>> out) throws Exception {
                String key = value.getKey();
                int val = value.getValue();

                // 更新当前键的累加值
                int currentKeyAccum = keyAccumulators.getOrDefault(key, 0) + val;
                keyAccumulators.put(key, currentKeyAccum);

                // 更新全局总和
                int currentTotal = totalSum.value() + val;
                totalSum.update(currentTotal);

                // 计算并输出比率
                if (currentTotal > 0) {
                    double ratio = (double) currentKeyAccum / currentTotal;
                    out.collect(Tuple2.of(key, ratio));
                }
            }
        })
        .print();

env.execute("Simple Real-Time Ratio Calculator");

关键说明

  • 并行度限制:全局状态在多并行度下会面临分布式状态同步问题,所以建议设为1
  • 状态维护:用MapState存储所有键的累加值,ValueState存储全局总和,逻辑直观易懂

方案对比与结果验证

两种方案都能完美匹配你的需求:

  • 输入k1,1 → 总和1,k1累加1 → 输出(k1, 1.0)
  • 输入k2,3 → 总和4,k2累加3 → 输出(k2, 0.75)
  • 输入k1,1 → 总和5,k1累加2 → 输出(k1, 0.4)
  • 输入k2,5 → 总和10,k2累加8 → 输出(k2, 0.8)
  • 输入k3,0 → 被过滤,无输出

选择方案时:

  • 数据量大、需要并行处理 → 选广播状态方案
  • 数据量小、追求实现简洁 → 选全局ProcessFunction方案

内容的提问来源于stack exchange,提问作者Richard Grossman

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 04:12:21