Hazelcast EntryProcessor.process方法被调用两次的原因排查
Hazelcast EntryProcessor.process 多次调用问题
我们尝试使用Hazelcast EntryProcessor,原本认为EntryProcessor.process方法应仅在主副本所在节点执行一次。但测试发现:当Hazelcast实例数为1时,process方法仅调用一次;当实例数大于1时,process方法会被调用两次。
测试代码
public class HazelcastBackupEntryProcessorTest { public static void main(String[] args) throws InterruptedException { List<HazelcastInstance> hazelcastInstances = setupHazelcast(2); HazelcastInstance hazelcastInstance = hazelcastInstances.get(0); IMap<Integer, String> testMap; testMap = hazelcastInstance.getMap("testMap"); if (!testMap.containsKey(1)) { testMap.put(1, "Holmes"); testMap.executeOnEntries(new UpdateNameEntryProcessor("Watson")); } } private static List<HazelcastInstance> setupHazelcast(int numberOfInstances) { List<HazelcastInstance> instances = new ArrayList<>(); IntStream.range(0, numberOfInstances).forEach(i -> { Config config = new Config(); config.getJetConfig().setEnabled(true); instances.add(Hazelcast.newHazelcastInstance(config)); }); return instances; } private static class UpdateNameEntryProcessor implements EntryProcessor<Integer, String, String> { private final String name; public UpdateNameEntryProcessor(String name) { this.name = name; } @Override public String process(Entry<Integer, String> entry) { System.out.println("Coming Here"); entry.setValue(name); return null; } } }
原因分析
- 默认备份副本的执行机制:Hazelcast的IMap默认配置
backup-count=1,即每个数据条目会生成1个备份副本。调用executeOnEntries执行EntryProcessor时,默认逻辑是主副本和备份副本所在节点都会执行该EntryProcessor——这是Hazelcast为了减少主节点同步修改到备份的网络开销,直接让操作在备份节点执行,保证数据一致性。 - 单实例场景的优化:当集群只有1个实例时,主副本和备份副本存储在同一节点,Hazelcast会自动优化,只执行一次
process方法,避免重复操作。 - 未区分主备执行逻辑:你的EntryProcessor没有实现
getBackupProcessor方法来指定备份节点的执行逻辑,Hazelcast会默认将同一个EntryProcessor应用到备份副本,因此出现两次调用。
解决方案
如果需要让process仅在主节点执行,可选择以下两种方式:
- 修改IMap配置,关闭备份(不推荐生产环境,会丧失数据容错能力):
Config config = new Config(); config.getMapConfig("testMap").setBackupCount(0); - 自定义备份处理器,让备份节点执行空逻辑:
private static class UpdateNameEntryProcessor implements EntryProcessor<Integer, String, String> { private final String name; public UpdateNameEntryProcessor(String name) { this.name = name; } @Override public String process(Entry<Integer, String> entry) { System.out.println("Coming Here (Primary)"); entry.setValue(name); return null; } @Override public EntryBackupProcessor<String> getBackupProcessor() { // 备份节点不执行任何操作 return entry -> {}; } }
内容的提问来源于stack exchange,提问作者Vpal
相关产品推荐
相关产品推荐

