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

使用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);

确保Flink启用Kryo作为默认序列化器:

env.getConfig().enableForceKryo();

同时避免排除RangeSet、TreeRangeSet及Instant相关类的序列化。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 14:35:45