使用Flink序列化系统无法正确序列化RangeSet<Instant>
解决Flink中Kryo序列化Guava RangeSet字段为null的问题
问题背景
在实现RichMapFunction<GeofenceEvent, OutputRangeSet>时,OutputRangeSet包含com.google.common.collect.RangeSet<Instant>字段,Kryo序列化后反序列化得到的该字段为null。已尝试以下方案但无效:
- 自定义
TypeInfoFactory<RangeSet>子类并为字段添加@TypeInfo注解 - 通过
env.getConfig().registerTypeWithKryoSerializer(RangeSet.class, ProtobufSerializer.class)注册第三方序列化器
可行解决方案
1. 注册具体实现类的Kryo序列化器
RangeSet是接口,实际序列化的是其具体实现(如TreeRangeSet),需针对实现类注册序列化器:
// 注册TreeRangeSet的Kryo序列化器 env.getConfig().registerTypeWithKryoSerializer(TreeRangeSet.class, new CollectionSerializer() { @Override protected Collection create(Kryo kryo, Input input, Class<? extends Collection> type) { return TreeRangeSet.create(); } });
若引入了kryo-guava依赖,可直接使用官方提供的RangeSetSerializer:
env.getConfig().registerTypeWithKryoSerializer(TreeRangeSet.class, RangeSetSerializer.class);
2. 为OutputRangeSet自定义Kryo序列化器
为OutputRangeSet实现专属Kryo序列化器,手动处理RangeSet<Instant>的读写逻辑:
// 为OutputRangeSet添加注解指定序列化器 @DefaultKryoSerializer(serializerClass = OutputRangeSetKryoSerializer.class) public class OutputRangeSet implements Serializable { private RangeSet<Instant> rangeSet; // getter、setter及其他逻辑 } // 自定义序列化器 public class OutputRangeSetKryoSerializer extends Serializer<OutputRangeSet> { @Override public void write(Kryo kryo, Output output, OutputRangeSet obj) { RangeSet<Instant> rangeSet = obj.getRangeSet(); // 写入Range数量 output.writeInt(rangeSet.asRanges().size()); // 逐个序列化Range的上下界 for (Range<Instant> range : rangeSet.asRanges()) { serializeBound(output, range.lowerBound()); serializeBound(output, range.upperBound()); } } @Override public OutputRangeSet read(Kryo kryo, Input input, Class<OutputRangeSet> type) { OutputRangeSet result = new OutputRangeSet(); int rangeCount = input.readInt(); RangeSet<Instant> rangeSet = TreeRangeSet.create(); // 逐个反序列化Range并添加到集合 for (int i = 0; i < rangeCount; i++) { Bound<Instant> lower = deserializeBound(input); Bound<Instant> upper = deserializeBound(input); rangeSet.add(Range.create(lower, upper)); } result.setRangeSet(rangeSet); return result; } // 序列化Range的Bound private void serializeBound(Output output, Bound<Instant> bound) throws IOException { output.writeBoolean(bound.isLowerBound()); output.writeBoolean(bound.isClosed()); output.writeLong(bound.getValue().toEpochMilli()); } // 反序列化Range的Bound private Bound<Instant> deserializeBound(Input input) throws IOException { boolean isLower = input.readBoolean(); boolean isClosed = input.readBoolean(); long epochMilli = input.readLong(); Instant instant = Instant.ofEpochMilli(epochMilli); return isLower ? Bound.lowerBound(instant, isClosed) : Bound.upperBound(instant, isClosed); } }
3. 手动指定TypeInformation
为OutputRangeSet创建自定义TypeInformation,明确RangeSet<Instant>字段的序列化逻辑:
TypeInformation<OutputRangeSet> outputRangeSetType = Types.POJO(OutputRangeSet.class, Collections.singletonMap("rangeSet", new TypeInformation<RangeSet<Instant>>() { @Override public boolean isBasicType() { return false; } @Override public boolean isTupleType() { return false; } @Override public int getArity() { return 1; } @Override public Type getTypeClass() { return RangeSet.class; } @Override public boolean isKeyType() { return false; } @Override public TypeSerializer<RangeSet<Instant>> createSerializer(ExecutionConfig config) { return new TypeSerializer<RangeSet<Instant>>() { @Override public boolean isImmutableType() { return false; } @Override public TypeSerializer<RangeSet<Instant>> duplicate() { return this; } @Override public RangeSet<Instant> createInstance() { return TreeRangeSet.create(); } @Override public RangeSet<Instant> copy(RangeSet<Instant> from) { return TreeRangeSet.create(from.asRanges()); } @Override public RangeSet<Instant> copy(RangeSet<Instant> from, RangeSet<Instant> reuse) { reuse.clear(); reuse.addAll(from.asRanges()); return reuse; } @Override public void serialize(RangeSet<Instant> value, DataOutputView target) throws IOException { target.writeInt(value.asRanges().size()); for (Range<Instant> range : value.asRanges()) { serializeBound(target, range.lowerBound()); serializeBound(target, range.upperBound()); } } @Override public RangeSet<Instant> deserialize(DataInputView source) throws IOException { RangeSet<Instant> rangeSet = TreeRangeSet.create(); int size = source.readInt(); for (int i = 0; i < size; i++) { Bound<Instant> lower = deserializeBound(source); Bound<Instant> upper = deserializeBound(source); rangeSet.add(Range.create(lower, upper)); } return rangeSet; } @Override public RangeSet<Instant> deserialize(RangeSet<Instant> reuse, DataInputView source) throws IOException { reuse.clear(); int size = source.readInt(); for (int i = 0; i < size; i++) { Bound<Instant> lower = deserializeBound(source); Bound<Instant> upper = deserializeBound(source); reuse.add(Range.create(lower, upper)); } return reuse; } @Override public void copy(DataInputView source, DataOutputView target) throws IOException { int size = source.readInt(); target.writeInt(size); for (int i = 0; i < size; i++) { // 复制Bound数据 boolean lowerIsLower = source.readBoolean(); boolean lowerIsClosed = source.readBoolean(); long lowerEpoch = source.readLong(); target.writeBoolean(lowerIsLower); target.writeBoolean(lowerIsClosed); target.writeLong(lowerEpoch); boolean upperIsLower = source.readBoolean(); boolean upperIsClosed = source.readBoolean(); long upperEpoch = source.readLong(); target.writeBoolean(upperIsLower); target.writeBoolean(upperIsClosed); target.writeLong(upperEpoch); } } @Override public int getLength() { return -1; } @Override public boolean equals(Object obj) { return obj instanceof this.getClass(); } @Override public int hashCode() { return this.getClass().hashCode(); } private void serializeBound(DataOutputView target, Bound<Instant> bound) throws IOException { target.writeBoolean(bound.isLowerBound()); target.writeBoolean(bound.isClosed()); target.writeLong(bound.getValue().toEpochMilli()); } private Bound<Instant> deserializeBound(DataInputView source) throws IOException { boolean isLower = source.readBoolean(); boolean isClosed = source.readBoolean(); long epochMilli = source.readLong(); Instant instant = Instant.ofEpochMilli(epochMilli); return isLower ? Bound.lowerBound(instant, isClosed) : Bound.upperBound(instant, isClosed); } }; } }));
在算子输出时指定该类型:
someDataStream.map(new YourRichMapFunction()).returns(outputRangeSetType);
4. 确认Flink Kryo配置
确保Flink启用Kryo作为默认序列化器:
env.getConfig().enableForceKryo();
同时避免排除RangeSet、TreeRangeSet及Instant相关类的序列化。
内容的提问来源于stack exchange,提问作者ElArbi
相关产品推荐
相关产品推荐

