Hazelcast中EntryProcessor实现Portable仍触发默认序列化问题求助
Hazelcast多节点下EntryProcessor备份序列化失败问题分析与修复
我来帮你拆解下这个问题的核心原因,以及对应的修复方案:
问题根源
从异常栈的最底层可以看到:
java.io.NotSerializableException: com.mycompany.common.MapValue
这个异常发生在Hazelcast序列化EntryBackupProcessorImpl(AbstractEntryProcessor的内部备份处理器)的时候,为什么会出现这种情况?主要有三个关键点:
- 默认备份处理器的序列化逻辑:你的
MyEntryProcessor虽然实现了Portable,但Hazelcast默认的备份处理器EntryBackupProcessorImpl是基于Java序列化的,它会持有你的MyEntryProcessor实例的所有状态(包括mapValue),并尝试用Java序列化来传递这些状态。 MapValue未实现序列化接口:你的MapValue类既没有实现Serializable,也没有实现Portable,导致Java序列化无法处理它,最终抛出异常。- Offloadable的影响:你启用了
OFFLOADABLE_EXECUTOR,虽然主操作卸载到了专用线程池,但备份操作依然需要传递EntryProcessor的状态,这就把序列化问题暴露了出来。
修复方案
要解决这个问题,需要从序列化支持和备份处理器自定义两个角度入手:
1. 让MapValue实现Portable接口
首先必须让MapValue支持Portable序列化(保持和EntryProcessor一致的序列化机制),示例代码如下:
package com.mycompany.common; import com.hazelcast.nio.serialization.Portable; import com.hazelcast.nio.serialization.PortableReader; import com.hazelcast.nio.serialization.PortableWriter; import java.io.IOException; public class MapValue implements Portable { // 自定义ClassId,需要和Portable工厂中的配置对应 private static final int CLASS_ID = 3; private String data; // 根据你的实际字段调整 // Portable必须有无参构造器 public MapValue() {} public MapValue(String data) { this.data = data; } // getter和setter public String getData() { return data; } public void setData(String data) { this.data = data; } @Override public int getFactoryId() { // 和MyEntryProcessor的FactoryId保持一致 return 1; } @Override public int getClassId() { return CLASS_ID; } @Override public void writePortable(PortableWriter writer) throws IOException { writer.writeString("data", data); } @Override public void readPortable(PortableReader reader) throws IOException { data = reader.readString("data"); } }
2. 自定义备份处理器并实现Portable
默认的备份处理器依赖Java序列化,我们可以自定义备份处理器,让它也实现Portable,这样备份操作就会使用Portable序列化机制。修改MyEntryProcessor,重写getBackupProcessor()方法:
// 在MyEntryProcessor类中添加以下内容 @Override public EntryBackupProcessor<String, MapValue> getBackupProcessor() { return new MyBackupProcessor(mapValue); } // 自定义备份处理器,实现Portable和EntryBackupProcessor private static class MyBackupProcessor implements EntryBackupProcessor<String, MapValue>, Portable { private MapValue mapValue; // Portable必须有无参构造器 public MyBackupProcessor() {} public MyBackupProcessor(MapValue mapValue) { this.mapValue = mapValue; } @Override public void processBackup(Entry<String, MapValue> entry) { // 和主处理器完全一致的业务逻辑 MapValue valueToSet = null; if (null == entry.getValue()) { valueToSet = mapValue; } else { MapValue valueToUpdate = entry.getValue(); valueToUpdate.setData(mapValue.getData()); valueToSet = valueToUpdate; } entry.setValue(valueToSet); } @Override public int getFactoryId() { return 1; // 共享同一个Portable工厂 } @Override public int getClassId() { return 4; // 唯一的ClassId,避免冲突 } @Override public void writePortable(PortableWriter writer) throws IOException { boolean hasMapValue = (mapValue != null); writer.writeBoolean("_has__mapValue", hasMapValue); if (hasMapValue) { writer.writePortable("mapValue", mapValue); } } @Override public void readPortable(PortableReader reader) throws IOException { if (reader.readBoolean("_has__mapValue")) { mapValue = reader.readPortable("mapValue", MapValue.class); } } }
3. 注册Portable工厂到所有节点
最后,必须在所有Hazelcast节点的配置文件(hazelcast-dev1.xml和hazelcast-dev2.xml)中注册你的Portable工厂,这样Hazelcast才能识别自定义的Portable类:
<hazelcast> <!-- 其他原有配置 --> <serialization> <portable-version>1</portable-version> <portable-factories> <portable-factory factory-id="1"> com.mycompany.common.MyPortableFactory </portable-factory> </portable-factories> </serialization> </hazelcast>
然后实现MyPortableFactory类:
package com.mycompany.common; import com.hazelcast.nio.serialization.Portable; import com.hazelcast.nio.serialization.PortableFactory; public class MyPortableFactory implements PortableFactory { @Override public Portable create(int classId) { switch (classId) { case 2: return new MyEntryProcessor(); case 3: return new MapValue(); case 4: return new MyEntryProcessor.MyBackupProcessor(); default: return null; } } }
验证步骤
- 确保所有节点的配置文件都正确注册了Portable工厂。
- 分别用
hazelcast-dev1.xml和hazelcast-dev2.xml启动两个Hazelcast节点。 - 重新执行
executeOnKeys()操作,此时备份操作会使用Portable序列化,不会再抛出NotSerializableException。
内容的提问来源于stack exchange,提问作者simpleusr
相关产品推荐
相关产品推荐

