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
相关产品推荐
相关产品推荐

