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

使用EvictingQueue计算滚动平均时触发FlinkRuntimeException问题排查

问题:Flink中使用Guava EvictingQueue维护滚动平均时出现序列化NPE

在Flink的KeyedBroadcastProcessFunction中,尝试用Guava 29的EvictingQueue维护最近N条记录的count属性滚动平均,简化代码如下:

public class PacingController extends KeyedBroadcastProcessFunction<> {

    private ValueState<EvictingQueue<Double>> lastNCountAvg;

    public void open(Configuration conf) {
        ValueStateDescriptor<EvictingQueue<Double>> descriptor =
                new ValueStateDescriptor<>(
                    "last-n-count",
                    TypeInformation.of(new TypeHint<>() {}));
        lastNCountAvg = getRuntimeContext().getState(descriptor);
    }

    public void processElement(..., ReadOnlyContext readOnlyContext,
                           Collector<BidPacingThresholdPerASUOutput> collector) throws Exception {    
        if(lastNCountAvg.value() == null){      
            lastNCountAvg.update(EvictingQueue.create(10)); // 初始化队列
        }
        // 其他业务逻辑
    }

    public void onTimer(
        long timestamp,
        OnTimerContext context,
        Collector<> out) throws Exception {
        // 更新最新count值
        EvictingQueue<Double> lastNCountAvgQueue = lastNCountAvg.value();
        lastNCountAvgQueue.add(count);
        lastNCountAvg.update(lastNCountAvgQueue); // 此处抛出异常
    }
}

已确保队列空时会初始化,但添加新条目更新状态时出现如下异常:

Caused by: TimerException{org.apache.flink.util.FlinkRuntimeException: Error while adding data to RocksDB}
    ... 14 more
Caused by: org.apache.flink.util.FlinkRuntimeException: Error while adding data to RocksDB
    at org.apache.flink.contrib.streaming.state.RocksDBValueState.update(RocksDBValueState.java:109)
    at inmarket.bidpacing.streaming.pacing.RollingAverageController.onTimer(PacingController.java:410)
    at org.apache.flink.streaming.api.operators.co.CoBroadcastWithKeyedOperator.onProcessingTime(CoBroadcastWithKeyedOperator.java:152)
    at org.apache.flink.streaming.api.operators.InternalTimerServiceImpl.onProcessingTime(InternalTimerServiceImpl.java:284)
    at org.apache.flink.streaming.runtime.tasks.StreamTask.invokeProcessingTimeCallback(StreamTask.java:1432)
    ... 13 more
Caused by: java.lang.NullPointerException
    at com.google.common.collect.ForwardingCollection.size(ForwardingCollection.java:65)
    at com.esotericsoftware.kryo.serializers.CollectionSerializer.write(CollectionSerializer.java:65)
    at com.esotericsoftware.kryo.serializers.CollectionSerializer.write(CollectionSerializer.java:22)
    at com.esotericsoftware.kryo.Kryo.writeClassAndObject(Kryo.java:599)
    at org.apache.flink.api.java.typeutils.runtime.kryo.KryoSerializer.serialize(KryoSerializer.java:316)
    at org.apache.flink.contrib.streaming.state.AbstractRocksDBState.serializeValueInternal(AbstractRocksDBState.java:158)
    at org.apache.flink.contrib.streaming.state.AbstractRocksDBState.serializeValue(AbstractRocksDBState.java:180)
    at org.apache.flink.contrib.streaming.state.AbstractRocksDBState.serializeValue(AbstractRocksDBState.java:168)
    at org.apache.flink.contrib.streaming.state.RocksDBValueState.update(RocksDBValueState.java:107)
    ... 17 more

询问是否为已知问题,以及无需手动维护滚动列表的替代方案。


原因分析

这并非Flink的已知问题,而是Guava EvictingQueue与Flink Kryo序列化机制的兼容性问题:

  • EvictingQueue继承自Guava的ForwardingCollection,依赖内部的delegate字段实现集合操作。
  • Kryo序列化时,无法正确保留EvictingQueue的内部delegate引用;反序列化后delegate为null,调用size()方法时触发NPE。
  • 即使初始化时队列正常,当状态被RocksDB持久化/恢复后,序列化的队列对象内部结构损坏,导致后续操作失败。

替代方案

1. 使用Apache Commons Collections的CircularFifoQueue

CircularFifoQueue是专门实现固定大小的循环队列,添加元素时会自动移除最旧的元素,且原生支持Java序列化,与Flink状态序列化兼容:

// 替换EvictingQueue为CircularFifoQueue
private ValueState<CircularFifoQueue<Double>> lastNCountAvg;

// 初始化时替换为
lastNCountAvg.update(new CircularFifoQueue<>(10));

2. 利用Flink滑动窗口计算滚动平均

如果你的滚动平均是基于时间窗口(比如最近10分钟的平均),直接使用Flink的滑动窗口API更简洁,无需手动维护队列:

// 在KeyedStream上定义滑动窗口
keyedStream
    .window(SlidingProcessingTimeWindows.of(Time.minutes(10), Time.seconds(1)))
    .aggregate(new AverageAggregate());

// 自定义AggregateFunction实现平均计算
public class AverageAggregate implements AggregateFunction<YourInput, Tuple2<Double, Integer>, Double> {
    @Override
    public Tuple2<Double, Integer> createAccumulator() {
        return Tuple2.of(0.0, 0);
    }

    @Override
    public Tuple2<Double, Integer> add(YourInput input, Tuple2<Double, Integer> accumulator) {
        return Tuple2.of(accumulator.f0 + input.getCount(), accumulator.f1 + 1);
    }

    @Override
    public Double getResult(Tuple2<Double, Integer> accumulator) {
        return accumulator.f0 / accumulator.f1;
    }

    @Override
    public Tuple2<Double, Integer> merge(Tuple2<Double, Integer> a, Tuple2<Double, Integer> b) {
        return Tuple2.of(a.f0 + b.f0, a.f1 + b.f1);
    }
}

3. 自定义可序列化的固定大小队列

如果不想引入第三方依赖,可以自己实现一个简单的固定大小队列,确保序列化正常:

public class FixedSizeQueue<E> extends ArrayList<E> {
    private final int maxSize;

    public FixedSizeQueue(int maxSize) {
        this.maxSize = maxSize;
    }

    @Override
    public boolean add(E e) {
        boolean added = super.add(e);
        while (size() > maxSize) {
            remove(0);
        }
        return added;
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 17:27:14