Flink自定义序列化方法使用及AtomicLongMap序列化注册失效问题
让我来逐个解答你的Flink相关问题:
Flink提供了多种自定义序列化的方式,你可以根据业务场景选择合适的方案:
基础方案:实现Serializable接口
如果你的类结构简单,直接实现java.io.Serializable接口即可,Flink的默认序列化器会处理它。不过这种方式性能偏慢,适合简单的POJO类,不适合复杂或高性能要求的场景。高效方案:基于Kryo自定义序列化
Flink默认用Kryo处理非基本类型的序列化,这也是大多数自定义场景的首选:- 编写自定义序列化器:继承
com.esotericsoftware.kryo.Serializer,实现write和read方法,定义对象的序列化/反序列化逻辑。 - 在Flink执行环境中注册你的序列化器:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.getConfig().registerTypeWithKryoSerializer(YourTargetClass.class, YourCustomSerializer.class);
要是需要强制所有类型都用Kryo序列化,可以加上
env.getConfig().enableForceKryo();。- 编写自定义序列化器:继承
深度集成方案:实现Flink TypeInformation
如果需要让你的类完全融入Flink的类型系统(比如支持更精准的优化),可以实现TypeInformation接口,并自定义对应的TypeSerializer,然后通过env.getConfig().registerTypeInformation(YourTargetClass.class, yourTypeInfo);注册。这种方式适合对性能和执行计划有极高要求的场景。
你遇到的org.apache.flink.api.common.InvalidProgramException,核心原因是RichFunction的构造函数中初始化的成员变量会被Flink序列化传输到TaskManager,但此时你注册的Kryo序列化器没正确关联到这个提前实例化的对象。下面给你两种可行的解决方案:
方案一:延迟初始化到open方法(推荐,无序列化风险)
RichFunction的open方法是在TaskManager本地执行的,不会被序列化传输。我们可以把AtomicLongMap的初始化移到这里,同时标记变量为transient避免被序列化:
public class PathAnalysis extends RichMapFunction<ApiLog, ApiLog> { // 用transient标记,告诉序列化器忽略这个变量 private transient AtomicLongMap<Object> mObjectAtomicLongMap; @Override public void open(Configuration parameters) throws Exception { super.open(parameters); // 在TaskManager本地初始化 mObjectAtomicLongMap = AtomicLongMap.create(); } @Override public ApiLog map(ApiLog value) throws Exception { // 在这里正常使用mObjectAtomicLongMap处理业务逻辑 return value; } }
这种方式完全避开了序列化问题,适合不需要跨Task共享AtomicLongMap状态的场景。
方案二:正确注册Kryo序列化器(适合需要序列化状态的场景)
如果你的业务需要将AtomicLongMap纳入Checkpoint或者需要跨节点传输,那需要确保Kryo能正确处理它:
- 先确保你的
AtomicLongMapSerializer正确实现了Kryo的序列化逻辑:public class AtomicLongMapSerializer extends Serializer<AtomicLongMap<?>> { @Override public void write(Kryo kryo, Output output, AtomicLongMap<?> object) { // 先写入map的大小,再逐个写入key和value output.writeInt(object.size()); for (Map.Entry<?, Long> entry : object.asMap().entrySet()) { kryo.writeClassAndObject(output, entry.getKey()); output.writeLong(entry.getValue()); } } @Override public AtomicLongMap<?> read(Kryo kryo, Input input, Class<? extends AtomicLongMap<?>> type) { AtomicLongMap<Object> map = AtomicLongMap.create(); int size = input.readInt(); for (int i = 0; i < size; i++) { Object key = kryo.readClassAndObject(input); long value = input.readLong(); map.put(key, value); } return map; } } - 在Flink环境中正确注册序列化器,并开启强制Kryo序列化:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 注册类型和对应的序列化器 env.getConfig().registerTypeWithKryoSerializer(AtomicLongMap.class, AtomicLongMapSerializer.class); // 强制所有非基本类型用Kryo序列化 env.getConfig().enableForceKryo(); - 修改RichFunction,不在构造函数中初始化AtomicLongMap,让Kryo负责实例化:
public class PathAnalysis extends RichMapFunction<ApiLog, ApiLog> { private AtomicLongMap<Object> mObjectAtomicLongMap; // 空构造函数,让Kryo在反序列化时创建实例 public PathAnalysis() {} @Override public void open(Configuration parameters) throws Exception { super.open(parameters); // 防止空指针,在open中做兜底初始化 if (mObjectAtomicLongMap == null) { mObjectAtomicLongMap = AtomicLongMap.create(); } } }
为什么原来的方式失效?
你在构造函数中直接实例化了mObjectAtomicLongMap,当Flink序列化PathAnalysis实例时,这个已经存在的对象会被优先处理,但此时Kryo的注册逻辑可能还没完全覆盖到这个实例的序列化,导致Flink无法找到合适的序列化器,从而抛出异常。
内容的提问来源于stack exchange,提问作者yeyunlong

