如何在Flink窗口流中基于SampleClass的time字段执行maxBy操作?
Flink窗口流基于自定义字段实现maxBy获取最新元素
你可以通过两种方式实现基于SampleClass的time字段执行maxBy操作,获取窗口内的最新元素:
方式一:使用带KeySelector的maxBy重载方法
Flink的maxBy算子支持传入KeySelector来指定用于比较的字段,无需直接比较整个SampleClass对象。在窗口后调用该重载方法,指定提取time字段作为比较依据即可:
SingleOutputStreamOperator<Tuple2<String, SampleClass>> sampleStream = previousStream .keyBy(value -> value.f1.getId()) .window(TumblingProcessingTimeWindows.of(Time.seconds(1))) // 传入KeySelector指定基于SampleClass的time字段比较,true表示当time相同时保留第一个元素 .maxBy(value -> value.f1.getTime(), true);
若你的Flink版本不支持lambda形式的KeySelector,可使用匿名内部类写法:
.maxBy(new KeySelector<Tuple2<String, SampleClass>, Long>() { @Override public Long getKey(Tuple2<String, SampleClass> value) throws Exception { return value.f1.getTime(); } }, true);
方式二:使用reduce算子自定义比较逻辑
如果需要更灵活的控制逻辑,可通过reduce算子手动实现保留time最大元素的逻辑:
SingleOutputStreamOperator<Tuple2<String, SampleClass>> sampleStream = previousStream .keyBy(value -> value.f1.getId()) .window(TumblingProcessingTimeWindows.of(Time.seconds(1))) .reduce((currentElement, nextElement) -> { // 比较两个元素的time字段,返回time更大的那个 if (nextElement.f1.getTime() > currentElement.f1.getTime()) { return nextElement; } else { // 若time相同,保留当前已有的元素(可根据需求调整逻辑) return currentElement; } });
两种方式均可实现需求:前者更简洁,适合仅需基于单个字段取最大值的场景;后者更灵活,可扩展处理复杂的比较或合并逻辑。
内容的提问来源于stack exchange,提问作者FelixNavidad
相关产品推荐
相关产品推荐

