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

Infinispan13设置primaryOnly=true时重平衡后新主节点未触发CacheEntryCreated如何解决?

问题根本原因

你当前配置的@Listener(primaryOnly = true, observation = Listener.Observation.POST)属于缓存条目事件监听器,仅当缓存条目本身发生创建、修改、删除等数据变更操作时才会触发事件。集群重平衡导致的缓存键主所有者切换,不会修改缓存条目本身的存储数据,因此不会触发该类监听器的事件通知,新主节点自然收不到对应提示。

解决方案

方案1:监听拓扑变更事件主动校验主所有权

你可以在原有监听器基础上增加集群拓扑变更事件的监听,重平衡完成后主动校验当前节点是否为对应缓存键的主所有者,符合条件则触发业务逻辑。
示例代码如下:

import org.infinispan.notifications.Listener;
import org.infinispan.notifications.cachelistener.annotation.CacheEntryCreated;
import org.infinispan.notifications.cachelistener.annotation.TopologyChanged;
import org.infinispan.notifications.cachelistener.event.CacheEntryCreatedEvent;
import org.infinispan.notifications.cachelistener.event.TopologyChangedEvent;
import org.infinispan.distribution.DistributionManager;
import org.infinispan.Cache;
import org.infinispan.remoting.transport.Address;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;

@Listener
public class PrimaryOwnerListener {
    private final Cache<String, Object> targetCache;
    // 若业务要求每个键仅执行一次处理逻辑,可增加已处理键的去重集合
    private final Set<String> processedKeys = ConcurrentHashMap.newKeySet();

    public PrimaryOwnerListener(Cache<String, Object> cache) {
        this.targetCache = cache;
    }

    // 原有缓存条目创建事件处理逻辑
    @CacheEntryCreated
    public void handleEntryCreated(CacheEntryCreatedEvent<String, Object> event) {
        DistributionManager dm = targetCache.getAdvancedCache().getDistributionManager();
        if (dm.getPrimaryOwner(event.getKey()).equals(targetCache.getCacheManager().getAddress())) {
            runBusinessLogic(event.getKey());
        }
    }

    // 新增拓扑变更事件监听
    @TopologyChanged
    public void handleTopologyChange(TopologyChangedEvent<String, Object> event) {
        // 仅处理重平衡完成后的后置事件
        if (!event.isPre()) {
            DistributionManager dm = targetCache.getAdvancedCache().getDistributionManager();
            Address currentNode = targetCache.getCacheManager().getAddress();
            // 遍历当前节点持有的所有缓存键
            targetCache.keySet().forEach(key -> {
                if (dm.getPrimaryOwner(key).equals(currentNode) && !processedKeys.contains(key)) {
                    runBusinessLogic(key);
                }
            });
        }
    }

    private void runBusinessLogic(String key) {
        // 你的业务处理逻辑
        System.out.printf("当前节点成为键%s的主所有者,开始处理业务%n", key);
        processedKeys.add(key);
    }
}

如果你的缓存键量级较大,可以通过TopologyChangedEvent获取新旧一致性哈希的差异,仅遍历重平衡后新分配到当前节点的键,降低遍历开销。

方案2:通过分布式独占锁实现单进程处理

如果你的核心需求是保证任意时刻仅一个进程处理指定缓存键,也可以不用感知主节点切换,直接在处理逻辑前加分布式独占锁:

// 处理业务前先获取对应键的分布式锁
boolean lockAcquired = targetCache.getAdvancedCache().lock(key);
if (lockAcquired) {
    try {
        // 执行业务逻辑,此时保证集群内仅有当前节点持有该键的锁
    } finally {
        targetCache.getAdvancedCache().unlock(key);
    }
}

该方案无需监听拓扑变更,节点掉线时锁会自动释放,其他节点可重新抢占锁执行逻辑,更适合对处理时机敏感度不高的业务场景。

注意事项

你使用的13.0.0版本完全兼容上述API,若缓存配置了持久化策略,节点启动加载持久化数据后,建议主动执行一次主所有权校验,避免重启阶段漏掉处理逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 04:06:01