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

如何结合Flink窗口、AggregateFunction并存储窗口聚合状态至自定义状态对象?

问题描述

我正在使用Apache Flink对Kafka流基于时间窗口进行聚合,当前实现为5分钟滚动事件时间窗口,窗口过期后数据才会保存到存储,代码结构如下:

SingleOutputStreamOperator<Aggregator> stream = source
        .forceNonParallel()
        .keyBy(Object::getKey)
        .window(TumblingEventTimeWindows.of(Time.minutes(5)))
        .aggregate(new CustomAggregator(), new KeyedWindowFunction("m5"));

我的需求是能够随时访问窗口状态以用于UI实时图表展示,但目前仅窗口过期后的数据才可被访问,用户无法看到当前正在聚合的数据,这并不理想。

我考虑了相关方案:

  • Flink是否支持通过可查询状态(Queryable State)访问窗口状态;
  • 创建自定义状态对象存储每次聚合结果,但AggregateFunction无法访问上下文无法实现;
  • ProcessWindowFunction虽能访问状态但仅在窗口过期时调用,无法满足实时需求。

现咨询:是否可以结合窗口、AggregateFunction,并将每次窗口聚合结果存储到自定义状态对象中?

解决方案

当然可以结合窗口、AggregateFunction与自定义状态实现实时访问窗口聚合结果,以下是两种可行方案:

方案一:RichAggregateFunction + 可查询状态

利用RichAggregateFunction的上下文能力,在聚合过程中将中间结果同步到可查询的键控状态,外部系统可直接查询该状态获取实时聚合值:

  1. 自定义带可查询状态的聚合函数:
public class QueryableCustomAggregator extends RichAggregateFunction<InputType, AccumulatorType, Aggregator> {
    private transient ValueState<AccumulatorType> queryableAggState;

    @Override
    public void open(Configuration parameters) throws Exception {
        ValueStateDescriptor<AccumulatorType> stateDesc = new ValueStateDescriptor<>(
                "real-time-window-agg",
                AccumulatorType.class
        );
        // 标记状态为可查询,外部通过"window-agg-query"这个名称查询
        stateDesc.setQueryable("window-agg-query");
        queryableAggState = getRuntimeContext().getState(stateDesc);
    }

    @Override
    public AccumulatorType createAccumulator() {
        return new AccumulatorType();
    }

    @Override
    public AccumulatorType add(InputType value, AccumulatorType accumulator) {
        // 执行你的核心聚合逻辑
        accumulator.merge(value);
        // 将最新聚合结果同步到可查询状态
        queryableAggState.update(accumulator);
        return accumulator;
    }

    @Override
    public Aggregator getResult(AccumulatorType accumulator) {
        // 转换为最终输出格式
        return new Aggregator(accumulator);
    }

    @Override
    public AccumulatorType merge(AccumulatorType a, AccumulatorType b) {
        a.merge(b);
        queryableAggState.update(a);
        return a;
    }
}
  1. 替换原聚合函数:
SingleOutputStreamOperator<Aggregator> stream = source
        .forceNonParallel()
        .keyBy(Object::getKey)
        .window(TumblingEventTimeWindows.of(Time.minutes(5)))
        .aggregate(new QueryableCustomAggregator(), new KeyedWindowFunction("m5"));

外部UI系统可通过Flink Queryable State Client,根据keyBy的键值查询对应窗口的实时聚合结果。

方案二:ProcessWindowFunction + 定时状态同步

如果需要更灵活的控制频率,可使用ProcessWindowFunction结合定时器,定期将窗口中间聚合结果同步到可查询状态:

  1. 自定义带定时更新的窗口处理函数:
public class RealTimeWindowProcessor extends ProcessWindowFunction<InputType, Aggregator, String, TimeWindow> {
    private transient ValueState<AccumulatorType> windowAggState;
    private transient ValueState<AccumulatorType> queryableState;

    @Override
    public void open(Configuration parameters) throws Exception {
        // 窗口内部聚合状态
        ValueStateDescriptor<AccumulatorType> windowStateDesc = new ValueStateDescriptor<>(
                "window-internal-agg",
                AccumulatorType.class
        );
        windowAggState = getRuntimeContext().getState(windowStateDesc);

        // 可查询状态
        ValueStateDescriptor<AccumulatorType> queryableDesc = new ValueStateDescriptor<>(
                "real-time-agg-state",
                AccumulatorType.class
        );
        queryableDesc.setQueryable("real-time-window-query");
        queryableState = getRuntimeContext().getState(queryableDesc);
    }

    @Override
    public void processElement(InputType value, Context context, Collector<Aggregator> out) throws Exception {
        AccumulatorType accumulator = windowAggState.value();
        if (accumulator == null) {
            accumulator = new AccumulatorType();
        }
        accumulator.update(value);
        windowAggState.update(accumulator);

        // 注册处理时间定时器,每10秒更新一次可查询状态
        long nextUpdate = context.currentProcessingTime() + 10000;
        context.timerService().registerProcessingTimeTimer(nextUpdate);
    }

    @Override
    public void onTimer(long timestamp, OnTimerContext ctx, Collector<Aggregator> out) throws Exception {
        AccumulatorType accumulator = windowAggState.value();
        if (accumulator != null) {
            queryableState.update(accumulator);
            // 可选:提前输出聚合结果到存储
            out.collect(new Aggregator(accumulator));
        }
        // 注册下一个定时器,持续更新
        ctx.timerService().registerProcessingTimeTimer(ctx.currentProcessingTime() + 10000);
    }

    @Override
    public void clear(Context context) throws Exception {
        windowAggState.clear();
        queryableState.clear();
    }
}
  1. 替换原窗口逻辑:
SingleOutputStreamOperator<Aggregator> stream = source
        .forceNonParallel()
        .keyBy(Object::getKey)
        .window(TumblingEventTimeWindows.of(Time.minutes(5)))
        .process(new RealTimeWindowProcessor());

关键注意事项

  • 可查询状态的键与keyBy的键一致,外部查询需指定对应键值;
  • 需合理配置状态TTL,避免长期运行导致内存泄漏;
  • 事件时间窗口场景下,需确保水印生成逻辑正确,避免窗口异常触发。

内容的提问来源于stack exchange,提问作者Nihad Hrnjić

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 06:25:29