Unable to Correctly Serialize RangeSet<Instant> with Flink Serialization System

Viewed 35

I've implemented a RichFunction with following type: RichMapFunction<GeofenceEvent, OutputRangeSet>

the class OutputRangeSet has a field of type:

com.google.common.collect.RangeSet<Instant>

When this pojo is serialized using Kryo I get null fields !

So far, I tried using a TypeInfoFactory<RangeSet>:

public class InstantRangeSetTypeInfo extends TypeInfoFactory<RangeSet<Instant>> {
        @Override
        public TypeInformation<RangeSet<Instant>> createTypeInfo(Type t, Map<String, TypeInformation<?>> genericParameters) {
            TypeInformation<RangeSet<Instant>> info = TypeInformation.of(new TypeHint<RangeSet<Instant>>() {});
            return info;
        }
    
    }

That annotate my field:

public class OutputRangeSet implements Serializable {

    private String key;

    @TypeInfo(InstantRangeSetTypeInfo.class)
    private RangeSet<Instant> rangeSet;

}

Another solution (that doesn't work either) is registring a third party serializer:

env.getConfig().registerTypeWithKryoSerializer(RangeSet.class, ProtobufSerializer.class);

You can get the github project here: https://github.com/elarbikonta/tonl-events

When you run the test you can see (in debug) that the rangeSet beans I get from my RichFunction has null fields, see test method com.tonl.apps.events.IsVehicleInZoneTest#operatorChronograph :

final RangeSet<Instant> rangeSet = resultList.get(0).getRangeSet(); // rangetSet.ranges = null !

Thanks for your help

0 Answers
Related