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

Hazelcast中EntryProcessor实现Portable仍触发默认序列化问题求助

Hazelcast多节点下EntryProcessor备份序列化失败问题分析与修复

我来帮你拆解下这个问题的核心原因,以及对应的修复方案:

问题根源

从异常栈的最底层可以看到:

java.io.NotSerializableException: com.mycompany.common.MapValue

这个异常发生在Hazelcast序列化EntryBackupProcessorImpl(AbstractEntryProcessor的内部备份处理器)的时候,为什么会出现这种情况?主要有三个关键点:

  1. 默认备份处理器的序列化逻辑:你的MyEntryProcessor虽然实现了Portable,但Hazelcast默认的备份处理器EntryBackupProcessorImpl是基于Java序列化的,它会持有你的MyEntryProcessor实例的所有状态(包括mapValue),并尝试用Java序列化来传递这些状态。
  2. MapValue未实现序列化接口:你的MapValue类既没有实现Serializable,也没有实现Portable,导致Java序列化无法处理它,最终抛出异常。
  3. 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;
        }
    }
}

验证步骤

  1. 确保所有节点的配置文件都正确注册了Portable工厂。
  2. 分别用hazelcast-dev1.xml和hazelcast-dev2.xml启动两个Hazelcast节点。
  3. 重新执行executeOnKeys()操作,此时备份操作会使用Portable序列化,不会再抛出NotSerializableException。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:43:24