From ba02977c64e4049c7c5b830bbd91909d66a7f302 Mon Sep 17 00:00:00 2001 From: fhan Date: Fri, 4 Sep 2026 10:15:42 +0800 Subject: [PATCH 1/2] [server] Fix thread-safety bug in RebalanceManager ZooKeeper recovery path --- .../CoordinatorEventProcessor.java | 8 ++++ .../event/RecoverRebalanceEvent.java | 39 ++++++++++++++++++ .../rebalance/RebalanceManager.java | 8 ++-- .../rebalance/RebalanceManagerTest.java | 41 +++++++++++++++++++ 4 files changed, 91 insertions(+), 5 deletions(-) create mode 100644 fluss-server/src/main/java/org/apache/fluss/server/coordinator/event/RecoverRebalanceEvent.java diff --git a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessor.java b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessor.java index 3b807cd38c9..233a128c30c 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessor.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessor.java @@ -86,6 +86,7 @@ import org.apache.fluss.server.coordinator.event.NotifyLeaderAndIsrResponseReceivedEvent; import org.apache.fluss.server.coordinator.event.RebalanceEvent; import org.apache.fluss.server.coordinator.event.RebalanceTaskTimeoutEvent; +import org.apache.fluss.server.coordinator.event.RecoverRebalanceEvent; import org.apache.fluss.server.coordinator.event.RemoveServerTagEvent; import org.apache.fluss.server.coordinator.event.ResumeDropEvent; import org.apache.fluss.server.coordinator.event.RetryOfflineLeaderEvent; @@ -743,6 +744,13 @@ public void process(CoordinatorEvent event) { RebalanceEvent rebalanceEvent = (RebalanceEvent) event; completeFromCallable( rebalanceEvent.getRespCallback(), () -> processRebalance(rebalanceEvent)); + } else if (event instanceof RecoverRebalanceEvent) { + RecoverRebalanceEvent recoverRebalanceEvent = (RecoverRebalanceEvent) event; + RebalanceTask rebalanceTask = recoverRebalanceEvent.getRebalanceTask(); + rebalanceManager.registerRebalance( + rebalanceTask.getRebalanceId(), + rebalanceTask.getExecutePlan(), + rebalanceTask.getRebalanceStatus()); } else if (event instanceof CancelRebalanceEvent) { CancelRebalanceEvent cancelRebalanceEvent = (CancelRebalanceEvent) event; completeFromCallable( diff --git a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/event/RecoverRebalanceEvent.java b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/event/RecoverRebalanceEvent.java new file mode 100644 index 00000000000..b891423513a --- /dev/null +++ b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/event/RecoverRebalanceEvent.java @@ -0,0 +1,39 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.fluss.server.coordinator.event; + +import org.apache.fluss.server.zk.data.RebalanceTask; + +/** Fired during startup to recover a rebalance task that was persisted in ZooKeeper. */ +public class RecoverRebalanceEvent implements CoordinatorEvent { + + private final RebalanceTask rebalanceTask; + + public RecoverRebalanceEvent(RebalanceTask rebalanceTask) { + this.rebalanceTask = rebalanceTask; + } + + public RebalanceTask getRebalanceTask() { + return rebalanceTask; + } + + @Override + public String toString() { + return "RecoverRebalanceEvent{rebalanceTask=" + rebalanceTask + "}"; + } +} diff --git a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java index cc29b982b07..ce95dd862a0 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java @@ -29,6 +29,7 @@ import org.apache.fluss.server.coordinator.CoordinatorEventProcessor; import org.apache.fluss.server.coordinator.event.EventManager; import org.apache.fluss.server.coordinator.event.RebalanceTaskTimeoutEvent; +import org.apache.fluss.server.coordinator.event.RecoverRebalanceEvent; import org.apache.fluss.server.coordinator.rebalance.goal.Goal; import org.apache.fluss.server.coordinator.rebalance.goal.GoalOptimizer; import org.apache.fluss.server.coordinator.rebalance.model.ClusterModel; @@ -178,11 +179,8 @@ private void initialize() { try { zkClient.getRebalanceTask() .ifPresent( - rebalancePlan -> - registerRebalance( - rebalancePlan.getRebalanceId(), - rebalancePlan.getExecutePlan(), - rebalancePlan.getRebalanceStatus())); + rebalanceTask -> + eventManager.put(new RecoverRebalanceEvent(rebalanceTask))); } catch (Exception e) { LOG.error( "Failed to get rebalance plan from zookeeper, it will be treated as no" diff --git a/fluss-server/src/test/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManagerTest.java b/fluss-server/src/test/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManagerTest.java index 57731571a1e..719f7a18f93 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManagerTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManagerTest.java @@ -34,6 +34,7 @@ import org.apache.fluss.server.coordinator.event.CoordinatorEvent; import org.apache.fluss.server.coordinator.event.EventManager; import org.apache.fluss.server.coordinator.event.RebalanceTaskTimeoutEvent; +import org.apache.fluss.server.coordinator.event.RecoverRebalanceEvent; import org.apache.fluss.server.coordinator.lease.KvSnapshotLeaseManager; import org.apache.fluss.server.coordinator.remote.RemoteDirDynamicLoader; import org.apache.fluss.server.metadata.CoordinatorMetadataCache; @@ -177,6 +178,46 @@ void testRebalanceWithoutTask() throws Exception { .hasValue(new RebalanceTask(rebalanceId, COMPLETED, new HashMap<>())); } + @Test + void testStartupQueuesRecoverRebalanceEvent() throws Exception { + ManualClock clock = new ManualClock(0L); + RecordingEventManager eventManager = new RecordingEventManager(); + NoOpScheduledExecutor executor = new NoOpScheduledExecutor(); + CoordinatorEventProcessor eventProcessor = + buildCoordinatorEventProcessor(new Configuration()); + + Map plan = createRebalancePlan(2); + RebalanceTask rebalanceTask = new RebalanceTask("recover-test", NOT_STARTED, plan); + zookeeperClient.registerRebalanceTask(rebalanceTask); + + RebalanceManager manager = + new RebalanceManager( + eventProcessor, zookeeperClient, eventManager, clock, executor); + // startup() 若在 ZK 中发现 pending rebalance 任务,应入队 RecoverRebalanceEvent, + // 由协调器事件线程执行 registerRebalance,而不是在启动线程直接调用。 + manager.startup(); + + assertThat(eventManager.events).hasSize(1); + assertThat(eventManager.events.get(0)).isInstanceOf(RecoverRebalanceEvent.class); + + RecoverRebalanceEvent recoverEvent = (RecoverRebalanceEvent) eventManager.events.get(0); + assertThat(recoverEvent.getRebalanceTask()).isEqualTo(rebalanceTask); + + manager.close(); + } + + private Map createRebalancePlan(int taskCount) { + Map plan = new HashMap<>(); + for (int i = 0; i < taskCount; i++) { + TableBucket tb = new TableBucket(1L, i); + plan.put( + tb, + new RebalancePlanForBucket( + tb, 0, 0, Arrays.asList(0, 1, 2), Arrays.asList(0, 1, 2))); + } + return plan; + } + @Test void testTimeoutEnqueuesEvent() throws Exception { ManualClock clock = new ManualClock(0L); From 67d390bb88e63aefd212b9c2d8dee09e7524fa6d Mon Sep 17 00:00:00 2001 From: fhan Date: Fri, 4 Sep 2026 10:30:39 +0800 Subject: [PATCH 2/2] [test] refine comments in RebalanceManagerTest --- .../server/coordinator/rebalance/RebalanceManagerTest.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/fluss-server/src/test/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManagerTest.java b/fluss-server/src/test/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManagerTest.java index 719f7a18f93..3bf3680fda6 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManagerTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManagerTest.java @@ -193,8 +193,9 @@ void testStartupQueuesRecoverRebalanceEvent() throws Exception { RebalanceManager manager = new RebalanceManager( eventProcessor, zookeeperClient, eventManager, clock, executor); - // startup() 若在 ZK 中发现 pending rebalance 任务,应入队 RecoverRebalanceEvent, - // 由协调器事件线程执行 registerRebalance,而不是在启动线程直接调用。 + // If startup() finds a pending rebalance task in ZooKeeper, it should enqueue a + // RecoverRebalanceEvent to be processed by the coordinator event thread, instead of + // calling registerRebalance() directly on the startup thread. manager.startup(); assertThat(eventManager.events).hasSize(1);