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

Hazelcast 5.2中EntryProcessor结合CompactSerialization报错求助

Hazelcast 5.2中EntryProcessor与CompactSerialization结合报错解决

问题描述

在Hazelcast 5.2版本中尝试将EntryProcessor与CompactSerialization结合使用时,遇到以下错误:

Hazelcast.Protocol.RemoteException: 'com.hazelcast.internal.serialization.impl.compact.DeserializedGenericRecord cannot be cast to com.hazelcast.map.EntryProcessor'

.NET客户端代码如下:

public class MyTest
{
    public int ID { get; set; }
    public bool IsGroupLevel { get; set; }
    public int ItemFormContext { get; set; }
    public MyTest() { }
    public MyTest(int id, bool isGroupLevel, int itemFormContext)
    {
        this.ID = id;
        this.IsGroupLevel = isGroupLevel;
        this.ItemFormContext = itemFormContext;
    }
}
public class MyTestSerializer : ICompactSerializer<MyTest>
{
    public string TypeName => "MyTest";

    public MyTest Read(ICompactReader reader)
    {
        var ID = reader.ReadInt32("ID");
        var IsGroupLevel = reader.ReadBoolean("IsGroupLevel");
        var ItemFormContext = reader.ReadInt32("ItemFormContext");
        return new MyTest(ID, IsGroupLevel, ItemFormContext);
    }

    public void Write(ICompactWriter writer, MyTest value)
    {
        writer.WriteInt32("ID", value.ID);
        writer.WriteBoolean("IsGroupLevel", value.IsGroupLevel);
        writer.WriteInt32("ItemFormContext", value.ItemFormContext);
    }
}

public class UpdateEntryProcessor : MyTestSerializer, IEntryProcessor<MyTest>
{
    private MyTest value;
    public UpdateEntryProcessor(MyTest value = null)
    {
        this.value = value;
    }
}

public class Program
{
    //Creating client
    var options = new HazelcastOptionsBuilder()
    .With(args)
    .WithDefault("hazelcast.networking.addresses.0", "<ip address>")   //Read Hazelcast server address from config file
    .WithDefault("hazelcast.clustername", "local")
    .WithDefault("hazelcast.socket.receive.buffer.size", "128")
    .WithDefault("hazelcast.socket.client.receive.buffer.size", "128")
    .WithDefault("hazelcast.socket.send.buffer.size", "128")
    .Build();
    options.Networking.ReconnectMode = Hazelcast.Networking.ReconnectMode.ReconnectSync;
    options.Serialization.Compact.AddSerializer(new MyTestSerializer());

    IHazelcastClient client = HazelcastClientFactory.StartNewClientAsync(options).GetAwaiter().GetResult();

    var vpmap = client.GetMapAsync<int, MyTest>("vpMap").GetAwaiter().GetResult();
    vpmap.SetAsync(1, new MyTest(1, false, 100)).GetAwaiter().GetResult();

    var r1 = vpmap.GetAsync(1).GetAwaiter().GetResult();

    //update using entry processor
    vpmap.ExecuteAsync(new UpdateEntryProcessor(r1), r1.ID).GetAwaiter().GetResult();
}

错误根源

  1. 跨语言执行限制:错误信息显示服务端是Java版本,.NET客户端定义的UpdateEntryProcessor无法在Java服务端执行——Java无法识别.NET的类和逻辑,服务端收到后只能反序列化为通用的DeserializedGenericRecord,无法转换为Java的EntryProcessor类型。
  2. 序列化类型混淆:UpdateEntryProcessor继承了MyTestSerializer,导致序列化框架将其当作MyTest类型处理,加剧了服务端的类型转换错误。
  3. 缺少核心逻辑:当前UpdateEntryProcessor未实现IEntryProcessor<MyTest>的Process方法,没有实际的更新逻辑。

解决方案

方案1:Java服务端实现EntryProcessor(推荐,适配Java集群)

如果你的Hazelcast集群是Java版本,必须在服务端用Java编写EntryProcessor:

  1. 编写Java版UpdateEntryProcessor,实现com.hazelcast.map.EntryProcessor接口,实现具体的更新逻辑。
  2. 将该类打包到服务端的类路径中,确保服务端能加载。
  3. .NET客户端调用时,通过EntryProcessor的名称或ID传递参数(参数用Compact序列化,比如MyTest类可正常在客户端和服务端之间传递)。

方案2:适配.NET版Hazelcast服务端

如果你的集群是.NET版本,修正代码如下:

1. 修正UpdateEntryProcessor类

移除不必要的继承,实现EntryProcessor的核心逻辑和自身的Compact序列化:

public class UpdateEntryProcessor : IEntryProcessor<MyTest>, ICompactSerializer<UpdateEntryProcessor>
{
    private MyTest value;

    // 无参构造函数,反序列化必需
    public UpdateEntryProcessor() { }

    public UpdateEntryProcessor(MyTest value)
    {
        this.value = value;
    }

    // 实现EntryProcessor的核心处理逻辑
    public object Process(IEntry<int, MyTest> entry)
    {
        var existingValue = entry.Value;
        // 执行你的更新操作,比如覆盖指定字段
        existingValue.IsGroupLevel = value.IsGroupLevel;
        existingValue.ItemFormContext = value.ItemFormContext;
        entry.SetValue(existingValue);
        return existingValue;
    }

    // Compact序列化配置
    public string TypeName => "UpdateEntryProcessor";

    public UpdateEntryProcessor Read(ICompactReader reader)
    {
        var testValue = reader.ReadCompact<MyTest>("value");
        return new UpdateEntryProcessor(testValue);
    }

    public void Write(ICompactWriter writer, UpdateEntryProcessor processor)
    {
        writer.WriteCompact("value", processor.value);
    }
}

2. 客户端配置添加EntryProcessor序列化器

在客户端配置中注册EntryProcessor的Compact序列化器:

options.Serialization.Compact.AddSerializer(new MyTestSerializer());
options.Serialization.Compact.AddSerializer(new UpdateEntryProcessor()); // 新增这行

额外配置说明

  • Java服务端:默认启用Compact序列化,无需额外配置,但需确保自定义EntryProcessor已加入服务端类路径。
  • .NET服务端:默认启用Compact序列化,无需额外配置,确保客户端和服务端的序列化器类型名称一致即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 11:10:41