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

Flink 1.14 Java实现DataStream按key取最新值计算中位数方案

Flink按Name分组、同Place取最新值计算滑动窗口中位数方案

你猜测的MapState是实现这个需求的核心,直接用内置的聚合函数没法满足同Place只保留最新值的规则,自定义窗口处理函数结合状态即可实现,完全可以支撑业务场景的中位数计算需求,不用被“流上不能算中位数”的说法限制——只要你的状态规模可控即可。

核心逻辑

  • 按Name做keyBy符合分组要求,后续所有状态都是Keyed级别的,不同Name的数据完全隔离
  • 用MapState<String, Tuple2<Long, Integer>>存储每个Place对应的最新消息:Map的key是Place维度值,value存该Place消息的事件时间戳和Number字段值,同Place新消息到达时直接覆盖旧值即可自动实现去重
  • 维护有效数值集合,每次Place值更新时同步替换集合里的旧Number,窗口触发时对集合排序即可算出中位数
  • 事件时间场景必须先配置合理的Watermark策略,避免乱序数据导致新值被旧值覆盖

可直接落地的实现代码

首先定义对应实体类(如果用Tuple类型传数据可以跳过,直接按位置取字段即可):

// 流消息实体类,对应(Name, Place, Number, Time)格式
public class Message {
    private String name;
    private String place;
    private Integer number;
    private Long time; // 事件时间,存毫秒级时间戳
    // 补全构造方法、getter/setter
}

// 中位数计算结果输出实体类
public class MedianResult {
    private String name;
    private Long windowEnd;
    private Double median;
    private Integer validCount; // 参与计算的有效数值个数
    // 补全构造方法、getter/setter
}

流处理核心逻辑:

// 1. 配置事件时间Watermark,示例允许2秒乱序,可根据业务实际乱序程度调整
WatermarkStrategy<Message> watermarkStrategy = WatermarkStrategy
        .<Message>forBoundedOutOfOrderness(Duration.ofSeconds(2))
        .withTimestampAssigner((msg, recordTimestamp) -> msg.getTime());

DataStream<Message> timedStream = inputStream.assignTimestampsAndWatermarks(watermarkStrategy);

// 2. 分组开窗计算中位数
timedStream
        .keyBy(Message::getName)
        .window(SlidingEventTimeWindows.of(Time.seconds(30), Time.seconds(1)))
        // 可根据需求配置允许迟到时间,示例允许10秒迟到数据计入窗口
        .allowedLateness(Time.seconds(10))
        .process(new ProcessWindowFunction<Message, MedianResult, String, TimeWindow>() {

            // 存储Place对应的最新(时间戳, Number值)映射
            private transient MapState<String, Tuple2<Long, Integer>> placeLatestState;
            // 存储当前窗口内所有有效Number值,用于计算中位数
            private transient ListState<Integer> validNumState;

            @Override
            public void open(Configuration parameters) throws Exception {
                MapStateDescriptor<String, Tuple2<Long, Integer>> placeStateDesc = new MapStateDescriptor<>(
                        "place-latest",
                        String.class,
                        TypeInformation.of(new TypeHint<Tuple2<Long, Integer>>() {})
                );
                placeLatestState = getRuntimeContext().getMapState(placeStateDesc);

                ListStateDescriptor<Integer> numStateDesc = new ListStateDescriptor<>(
                        "valid-nums",
                        Integer.class
                );
                validNumState = getRuntimeContext().getListState(numStateDesc);
            }

            @Override
            public void process(String name, Context context, Iterable<Message> elements, Collector<MedianResult> out) throws Exception {
                // 遍历窗口内所有元素,更新去重状态
                for (Message msg : elements) {
                    String place = msg.getPlace();
                    Long msgTime = msg.getTime();
                    Integer num = msg.getNumber();

                    Tuple2<Long, Integer> existing = placeLatestState.get(place);
                    // 仅当当前Place无记录,或新消息时间晚于已有记录时间时更新
                    if (existing == null || msgTime > existing.f0) {
                        // 先移除旧的无效数值
                        if (existing != null) {
                            List<Integer> tempNums = new ArrayList<>();
                            for (Integer validNum : validNumState.get()) {
                                if (!validNum.equals(existing.f1)) {
                                    tempNums.add(validNum);
                                }
                            }
                            validNumState.update(tempNums);
                        }
                        // 写入新的有效记录
                        placeLatestState.put(place, Tuple2.of(msgTime, num));
                        validNumState.add(num);
                    }
                }

                // 排序后计算中位数
                List<Integer> sortedNums = new ArrayList<>();
                validNumState.get().forEach(sortedNums::add);
                Collections.sort(sortedNums);
                double median;
                int size = sortedNums.size();
                if (size % 2 == 1) {
                    median = sortedNums.get(size / 2);
                } else {
                    median = (sortedNums.get(size/2 -1) + sortedNums.get(size/2)) / 2.0;
                }

                out.collect(new MedianResult(name, context.window().getEnd(), median, size));
            }
        });

用你给出的Jonah测试数据验证逻辑:

  • 第一条(Jonah, Mars, 1, 1:00):写入Mars对应记录,有效数值为[1]
  • 第二条(Jonah, Mars, 2, 1:01):Mars记录时间更新,删除旧值1,写入新值2,有效数值为[2]
  • 第三条(Jonah, Moon, 3, 1:02):新增Moon记录,有效数值为[2,3]
  • 第四条(Jonah, Earth, 4, 1:03):新增Earth记录,有效数值为[2,3,4]
  • 排序后取中位数结果为3,完全符合规则要求。

优化提示

  • 性能优化:如果单Key下Place维度超过10万级,用List排序的性能会下降,可以替换为双堆结构(大顶堆存小于中位数的值、小顶堆存大于等于中位数的值)维护中位数,插入和查询的时间复杂度从O(nlogn)降到O(logn),只需要把ListState替换为两个对应堆结构的ValueState即可。绝大多数业务场景下Place维度不会达到这个量级,直接用List排序足够用
  • 状态膨胀问题:Flink窗口会在生命周期结束后自动清理对应窗口的所有状态,只要你设置的窗口长度、允许迟到时间合理,不会出现状态无限增长的问题
  • 乱序/迟到数据:除了调整Watermark乱序阈值、allowedLateness时间外,还可以用侧输出流收集超过迟到容忍时间的数据,后续做离线补算即可

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 13:57:12