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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 05:17:11