如何让Ignite所有缓存的Rebalancing事件仅触发一次?
解决Ignite多缓存Rebalancing事件重复触发的问题
嘿,作为踩过Ignite不少坑的老玩家,我完全懂你这种被三次重复事件烦到的心情!毕竟每个缓存的Rebalancing是独立触发的,节点变动时三个缓存自然会各发一次事件。不过有两种靠谱的方法能让你只触发一次逻辑:
方法1:监听节点拓扑事件(推荐)
既然节点加入/离开是Rebalancing的根源,不如直接监听节点级别的拓扑事件,而不是每个缓存的Rebalance事件。这样不管有多少缓存,节点变动时只会触发一次你的逻辑:
Ignite ignite = Ignition.start(); // 注册本地监听,监听节点加入和离开事件 ignite.events().localListen(evt -> { // 这里写你要执行的全局逻辑 System.out.println("节点拓扑变动,执行一次全局处理"); return true; // 返回true保持监听 }, EVT_NODE_JOINED, EVT_NODE_LEFT);
这种方法最省心,因为拓扑事件是全局触发的,完全避开了多缓存的重复问题。如果需要在Rebalance完成后执行逻辑,可以结合EVT_CACHE_REBALANCE_FINISHED事件,再配合全局锁做去重。
方法2:共享监听器+去重逻辑
如果你一定要监听Rebalance事件,可以创建一个共享的监听器实例,并在内部加去重机制,确保同一节点变动周期内只执行一次:
import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import org.apache.ignite.events.CacheRebalanceStartedEvent; import org.apache.ignite.events.CacheRebalancingStoppedEvent; import org.apache.ignite.lang.CacheRebalancingListener; // 自定义共享监听器 class GlobalRebalanceListener implements CacheRebalancingListener<Object, Object> { // 用原子布尔值标记是否已经触发过逻辑 private final AtomicBoolean isTriggered = new AtomicBoolean(false); // 调度器用于重置标记,避免下次节点变动无法触发 private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(); @Override public void onRebalancingStarted(CacheRebalanceStartedEvent<Object, Object> evt) { // CAS操作:只有当标记为false时才执行逻辑 if (isTriggered.compareAndSet(false, true)) { // 这里写你的全局处理逻辑 System.out.println("Rebalance启动,执行一次全局处理"); // 重置标记的时间建议根据你的Rebalance平均耗时调整 scheduler.schedule(() -> isTriggered.set(false), 15, TimeUnit.SECONDS); } } @Override public void onRebalancingStopped(CacheRebalancingStoppedEvent<Object, Object> evt) { // 如果需要在Rebalance结束时执行逻辑,也可以在这里加类似的去重 } } // 给所有缓存注册同一个监听器实例 Ignite ignite = Ignition.start(); GlobalRebalanceListener sharedListener = new GlobalRebalanceListener(); ignite.cache("cache1").rebalancing().addListener(sharedListener); ignite.cache("cache2").rebalancing().addListener(sharedListener); ignite.cache("cache3").rebalancing().addListener(sharedListener);
注意事项
- 如果你是在分布式集群中运行,建议用Ignite的
AtomicLong或者分布式锁(IgniteLock)来做全局去重,避免多个节点同时触发逻辑。 - 方法2中的重置时间要根据你的缓存大小、网络情况调整,确保Rebalance完全结束后标记才会重置,避免遗漏下一次节点变动的事件。
内容的提问来源于stack exchange,提问作者Hyun
相关产品推荐
相关产品推荐

