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

