Copilot commented on code in PR #3622:
URL: https://github.com/apache/fluss/pull/3622#discussion_r3548338709
##########
fluss-server/src/main/java/org/apache/fluss/server/DynamicServerConfig.java:
##########
@@ -110,6 +112,8 @@ class DynamicServerConfig {
void register(ServerReconfigurable serverReconfigurable) {
serverReconfigures.put(serverReconfigurable.getClass(),
serverReconfigurable);
+ serverReconfigurable.validate(currentConfig);
+ serverReconfigurable.reconfigure(currentConfig);
}
Review Comment:
`register()` mutates and reads `currentConfig` without the `ReadWriteLock`,
even though `updateDynamicConfig()` updates `currentConfig` under the write
lock. This creates a data race (late registration can observe a
stale/inconsistent `currentConfig`) and can also leave a failing
`ServerReconfigurable` in `serverReconfigures` because it is added to the map
before `validate()` succeeds.
##########
fluss-server/src/test/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessorTest.java:
##########
@@ -2241,7 +2241,14 @@ private static void drainPendingNotifyTriggers(
private int countInProgressRebalanceTasks(TableBucket... buckets) {
int count = 0;
for (TableBucket tb : buckets) {
- if
(eventProcessor.getRebalanceManager().getRebalancePlanForBucket(tb) != null) {
+ RebalanceStatus status =
+ eventProcessor
+ .getRebalanceManager()
+ .listRebalanceProgress(null)
+ .progressForBucketMap()
+ .get(tb)
+ .status();
+ if (!RebalanceStatus.FINAL_STATUSES.contains(status)) {
count++;
}
Review Comment:
`countInProgressRebalanceTasks(...)` assumes `listRebalanceProgress(null)`
is non-null and that `progressForBucketMap().get(tb)` returns a non-null entry.
If a bucket is absent (e.g., after cancel or if the progress map is cleared),
this will throw an NPE and make the test brittle.
##########
fluss-server/src/test/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManagerTest.java:
##########
@@ -308,6 +349,385 @@ void testTimeoutTreatsTaskAsCompleted() throws Exception {
manager.close();
}
+ @Test
+ void testRegisterRebalanceRespectsMaxInflightTasks() throws Exception {
+ ManualClock clock = new ManualClock(0L);
+ RecordingEventManager eventManager = new RecordingEventManager();
+ NoOpScheduledExecutor executor = new NoOpScheduledExecutor();
+ Configuration conf = new Configuration();
+ conf.set(ConfigOptions.COORDINATOR_REBALANCE_MAX_INFLIGHT_TASKS, 2);
+ RecordingCoordinatorEventProcessor eventProcessor =
+ buildRecordingCoordinatorEventProcessor(conf);
+ RebalanceManager manager =
+ new RebalanceManager(
+ eventProcessor, zookeeperClient, eventManager, clock,
conf, executor);
+
+ Map<TableBucket, RebalancePlanForBucket> plan = createRebalancePlan(5);
+ List<TableBucket> buckets = new ArrayList<>(plan.keySet());
+ zookeeperClient.registerRebalanceTask(
+ new RebalanceTask("inflight-test", NOT_STARTED, plan));
+
+ manager.registerRebalance("inflight-test", plan, NOT_STARTED);
+
+ assertThat(eventProcessor.executedPlans).hasSize(2);
+ assertThat(countStatus(manager,
RebalanceStatus.REBALANCING)).isEqualTo(2);
+ assertThat(countStatus(manager,
RebalanceStatus.NOT_STARTED)).isEqualTo(3);
+
+ manager.finishRebalanceTask(buckets.get(0), COMPLETED);
+
+ assertThat(eventProcessor.executedPlans).hasSize(3);
+
assertThat(eventProcessor.executedPlans.get(2).getTableBucket()).isEqualTo(buckets.get(2));
+ assertThat(countStatus(manager,
RebalanceStatus.REBALANCING)).isEqualTo(2);
+ assertThat(countStatus(manager,
RebalanceStatus.NOT_STARTED)).isEqualTo(2);
+
+ manager.close();
+ }
+
+ @Test
+ void testRebalanceRoundLimitsActivatedBuckets() throws Exception {
+ ManualClock clock = new ManualClock(0L);
+ RecordingEventManager eventManager = new RecordingEventManager();
+ NoOpScheduledExecutor executor = new NoOpScheduledExecutor();
+ Configuration conf = new Configuration();
+ conf.set(ConfigOptions.COORDINATOR_REBALANCE_MAX_INFLIGHT_TASKS, 2);
+ conf.set(ConfigOptions.COORDINATOR_REBALANCE_MAX_BUCKETS_PER_ROUND, 2);
+ RecordingCoordinatorEventProcessor eventProcessor =
+ buildRecordingCoordinatorEventProcessor(conf);
+ RebalanceManager manager =
+ new RebalanceManager(
+ eventProcessor, zookeeperClient, eventManager, clock,
conf, executor);
+
+ Map<TableBucket, RebalancePlanForBucket> plan = createRebalancePlan(5);
+ List<TableBucket> buckets = new ArrayList<>(plan.keySet());
+ zookeeperClient.registerRebalanceTask(new RebalanceTask("round-test",
NOT_STARTED, plan));
+
+ manager.registerRebalance("round-test", plan, NOT_STARTED);
+
+ assertThat(manager.getMaxBucketsPerRound()).isEqualTo(2);
+ assertThat(eventProcessor.executedPlans).hasSize(2);
+ assertThat(countStatus(manager,
RebalanceStatus.REBALANCING)).isEqualTo(2);
+ assertThat(countStatus(manager,
RebalanceStatus.NOT_STARTED)).isEqualTo(3);
+
+ manager.finishRebalanceTask(buckets.get(0), COMPLETED);
+
+ assertThat(eventProcessor.executedPlans).hasSize(2);
+ assertThat(countStatus(manager,
RebalanceStatus.REBALANCING)).isEqualTo(1);
+ assertThat(countStatus(manager,
RebalanceStatus.NOT_STARTED)).isEqualTo(3);
+
+ manager.finishRebalanceTask(buckets.get(1), COMPLETED);
+
+ assertThat(eventProcessor.executedPlans).hasSize(4);
+
assertThat(eventProcessor.executedPlans.get(2).getTableBucket()).isEqualTo(buckets.get(2));
+
assertThat(eventProcessor.executedPlans.get(3).getTableBucket()).isEqualTo(buckets.get(3));
+ assertThat(countStatus(manager,
RebalanceStatus.REBALANCING)).isEqualTo(2);
+ assertThat(countStatus(manager,
RebalanceStatus.NOT_STARTED)).isEqualTo(1);
+
+ manager.close();
+ }
+
+ @Test
+ void testRebalanceRoundWorksWithLowerMaxInflightTasks() throws Exception {
+ ManualClock clock = new ManualClock(0L);
+ RecordingEventManager eventManager = new RecordingEventManager();
+ NoOpScheduledExecutor executor = new NoOpScheduledExecutor();
+ Configuration conf = new Configuration();
+ conf.set(ConfigOptions.COORDINATOR_REBALANCE_MAX_INFLIGHT_TASKS, 1);
+ conf.set(ConfigOptions.COORDINATOR_REBALANCE_MAX_BUCKETS_PER_ROUND, 2);
+ RecordingCoordinatorEventProcessor eventProcessor =
+ buildRecordingCoordinatorEventProcessor(conf);
+ RebalanceManager manager =
+ new RebalanceManager(
+ eventProcessor, zookeeperClient, eventManager, clock,
conf, executor);
+
+ Map<TableBucket, RebalancePlanForBucket> plan = createRebalancePlan(3);
+ List<TableBucket> buckets = new ArrayList<>(plan.keySet());
+ zookeeperClient.registerRebalanceTask(
+ new RebalanceTask("round-lower-inflight-test", NOT_STARTED,
plan));
+
+ manager.registerRebalance("round-lower-inflight-test", plan,
NOT_STARTED);
+
+ assertThat(eventProcessor.executedPlans).hasSize(1);
+
assertThat(eventProcessor.executedPlans.get(0).getTableBucket()).isEqualTo(buckets.get(0));
+ assertThat(countStatus(manager,
RebalanceStatus.REBALANCING)).isEqualTo(1);
+ assertThat(countStatus(manager,
RebalanceStatus.NOT_STARTED)).isEqualTo(2);
+
+ manager.finishRebalanceTask(buckets.get(0), COMPLETED);
+
+ assertThat(eventProcessor.executedPlans).hasSize(2);
+
assertThat(eventProcessor.executedPlans.get(1).getTableBucket()).isEqualTo(buckets.get(1));
+ assertThat(countStatus(manager,
RebalanceStatus.REBALANCING)).isEqualTo(1);
+ assertThat(countStatus(manager,
RebalanceStatus.NOT_STARTED)).isEqualTo(1);
+
+ manager.finishRebalanceTask(buckets.get(1), COMPLETED);
+
+ assertThat(eventProcessor.executedPlans).hasSize(3);
+
assertThat(eventProcessor.executedPlans.get(2).getTableBucket()).isEqualTo(buckets.get(2));
+ assertThat(countStatus(manager,
RebalanceStatus.REBALANCING)).isEqualTo(1);
+ assertThat(countStatus(manager, RebalanceStatus.NOT_STARTED)).isZero();
+
+ manager.close();
+ }
+
+ @Test
+ void testZeroMaxBucketsPerRoundKeepsMaxInflightBehavior() throws Exception
{
+ ManualClock clock = new ManualClock(0L);
+ RecordingEventManager eventManager = new RecordingEventManager();
+ NoOpScheduledExecutor executor = new NoOpScheduledExecutor();
+ Configuration conf = new Configuration();
+ conf.set(ConfigOptions.COORDINATOR_REBALANCE_MAX_INFLIGHT_TASKS, 2);
+ conf.set(ConfigOptions.COORDINATOR_REBALANCE_MAX_BUCKETS_PER_ROUND, 0);
+ RecordingCoordinatorEventProcessor eventProcessor =
+ buildRecordingCoordinatorEventProcessor(conf);
+ RebalanceManager manager =
+ new RebalanceManager(
+ eventProcessor, zookeeperClient, eventManager, clock,
conf, executor);
+
+ Map<TableBucket, RebalancePlanForBucket> plan = createRebalancePlan(5);
+ List<TableBucket> buckets = new ArrayList<>(plan.keySet());
+ zookeeperClient.registerRebalanceTask(
+ new RebalanceTask("unlimited-round-test", NOT_STARTED, plan));
+
+ manager.registerRebalance("unlimited-round-test", plan, NOT_STARTED);
+
+ assertThat(manager.getMaxBucketsPerRound()).isZero();
+ assertThat(eventProcessor.executedPlans).hasSize(2);
+ assertThat(countStatus(manager,
RebalanceStatus.REBALANCING)).isEqualTo(2);
+ assertThat(countStatus(manager,
RebalanceStatus.NOT_STARTED)).isEqualTo(3);
+
+ manager.finishRebalanceTask(buckets.get(0), COMPLETED);
+
+ assertThat(eventProcessor.executedPlans).hasSize(3);
+
assertThat(eventProcessor.executedPlans.get(2).getTableBucket()).isEqualTo(buckets.get(2));
+ assertThat(countStatus(manager,
RebalanceStatus.REBALANCING)).isEqualTo(2);
+ assertThat(countStatus(manager,
RebalanceStatus.NOT_STARTED)).isEqualTo(2);
+
+ manager.close();
+ }
+
+ @Test
+ void testMaxBucketsPerRoundSmallerThanMaxInflightCapsScheduling() throws
Exception {
+ ManualClock clock = new ManualClock(0L);
+ RecordingEventManager eventManager = new RecordingEventManager();
+ NoOpScheduledExecutor executor = new NoOpScheduledExecutor();
+ Configuration conf = new Configuration();
+ conf.set(ConfigOptions.COORDINATOR_REBALANCE_MAX_INFLIGHT_TASKS, 3);
+ conf.set(ConfigOptions.COORDINATOR_REBALANCE_MAX_BUCKETS_PER_ROUND, 1);
+ RecordingCoordinatorEventProcessor eventProcessor =
+ buildRecordingCoordinatorEventProcessor(conf);
+ RebalanceManager manager =
+ new RebalanceManager(
+ eventProcessor, zookeeperClient, eventManager, clock,
conf, executor);
+
+ Map<TableBucket, RebalancePlanForBucket> plan = createRebalancePlan(3);
+ List<TableBucket> buckets = new ArrayList<>(plan.keySet());
+ zookeeperClient.registerRebalanceTask(
+ new RebalanceTask("round-smaller-than-inflight-test",
NOT_STARTED, plan));
+
+ manager.registerRebalance("round-smaller-than-inflight-test", plan,
NOT_STARTED);
+
+ assertThat(eventProcessor.executedPlans).hasSize(1);
+
assertThat(eventProcessor.executedPlans.get(0).getTableBucket()).isEqualTo(buckets.get(0));
+ assertThat(countStatus(manager,
RebalanceStatus.REBALANCING)).isEqualTo(1);
+ assertThat(countStatus(manager,
RebalanceStatus.NOT_STARTED)).isEqualTo(2);
+
+ manager.finishRebalanceTask(buckets.get(0), COMPLETED);
+
+ assertThat(eventProcessor.executedPlans).hasSize(2);
+
assertThat(eventProcessor.executedPlans.get(1).getTableBucket()).isEqualTo(buckets.get(1));
+ assertThat(countStatus(manager,
RebalanceStatus.REBALANCING)).isEqualTo(1);
+ assertThat(countStatus(manager,
RebalanceStatus.NOT_STARTED)).isEqualTo(1);
+
+ manager.close();
+ }
+
+ @Test
+ void testNegativeMaxBucketsPerRoundIsRejected() {
+ ManualClock clock = new ManualClock(0L);
+ RecordingEventManager eventManager = new RecordingEventManager();
+ NoOpScheduledExecutor executor = new NoOpScheduledExecutor();
+ Configuration conf = new Configuration();
+ conf.set(ConfigOptions.COORDINATOR_REBALANCE_MAX_BUCKETS_PER_ROUND,
-1);
+ RecordingCoordinatorEventProcessor eventProcessor =
+ buildRecordingCoordinatorEventProcessor(new Configuration());
+
+ assertThatThrownBy(
+ () ->
+ new RebalanceManager(
+ eventProcessor,
+ zookeeperClient,
+ eventManager,
+ clock,
+ conf,
+ executor))
+ .hasMessageContaining(
+
ConfigOptions.COORDINATOR_REBALANCE_MAX_BUCKETS_PER_ROUND.key());
+ }
+
+ @Test
+ void testIncreaseMaxInflightRebalanceTasksStartsMoreTasks() throws
Exception {
+ ManualClock clock = new ManualClock(0L);
+ RecordingEventManager eventManager = new RecordingEventManager();
+ NoOpScheduledExecutor executor = new NoOpScheduledExecutor();
+ Configuration conf = new Configuration();
+ conf.set(ConfigOptions.COORDINATOR_REBALANCE_MAX_INFLIGHT_TASKS, 1);
+ RecordingCoordinatorEventProcessor eventProcessor =
+ buildRecordingCoordinatorEventProcessor(conf);
+ RebalanceManager manager =
+ new RebalanceManager(
+ eventProcessor, zookeeperClient, eventManager, clock,
conf, executor);
+
+ Map<TableBucket, RebalancePlanForBucket> plan = createRebalancePlan(4);
+ zookeeperClient.registerRebalanceTask(
+ new RebalanceTask("increase-inflight-test", NOT_STARTED,
plan));
+ manager.registerRebalance("increase-inflight-test", plan, NOT_STARTED);
+ assertThat(eventProcessor.executedPlans).hasSize(1);
+
+ Configuration newConfig = new Configuration(conf);
+ newConfig.set(ConfigOptions.COORDINATOR_REBALANCE_MAX_INFLIGHT_TASKS,
3);
+ manager.reconfigure(newConfig);
+
+ assertThat(manager.getMaxInflightRebalanceTasks()).isEqualTo(1);
+ assertThat(eventProcessor.executedPlans).hasSize(1);
+ applyLatestMaxInflightTasksChangedEvent(eventManager, manager, 3);
+
+ assertThat(manager.getMaxInflightRebalanceTasks()).isEqualTo(3);
+ assertThat(eventProcessor.executedPlans).hasSize(3);
+ assertThat(countStatus(manager,
RebalanceStatus.REBALANCING)).isEqualTo(3);
+ assertThat(countStatus(manager,
RebalanceStatus.NOT_STARTED)).isEqualTo(1);
+
+ manager.close();
+ }
+
+ @Test
+ void testZeroMaxInflightRebalanceTasksPausesAndResumesScheduling() throws
Exception {
+ ManualClock clock = new ManualClock(0L);
+ RecordingEventManager eventManager = new RecordingEventManager();
+ NoOpScheduledExecutor executor = new NoOpScheduledExecutor();
+ Configuration conf = new Configuration();
+ conf.set(ConfigOptions.COORDINATOR_REBALANCE_MAX_INFLIGHT_TASKS, 0);
+ RecordingCoordinatorEventProcessor eventProcessor =
+ buildRecordingCoordinatorEventProcessor(conf);
+ RebalanceManager manager =
+ new RebalanceManager(
+ eventProcessor, zookeeperClient, eventManager, clock,
conf, executor);
+
+ Map<TableBucket, RebalancePlanForBucket> plan = createRebalancePlan(4);
+ zookeeperClient.registerRebalanceTask(
+ new RebalanceTask("paused-inflight-test", NOT_STARTED, plan));
+ manager.registerRebalance("paused-inflight-test", plan, NOT_STARTED);
+
+ assertThat(manager.getMaxInflightRebalanceTasks()).isZero();
+ assertThat(eventProcessor.executedPlans).isEmpty();
+ assertThat(countStatus(manager,
RebalanceStatus.NOT_STARTED)).isEqualTo(4);
+ assertThat(manager.hasInProgressRebalance()).isTrue();
+
+ Configuration newConfig = new Configuration(conf);
+ newConfig.set(ConfigOptions.COORDINATOR_REBALANCE_MAX_INFLIGHT_TASKS,
2);
+ manager.reconfigure(newConfig);
+
+ assertThat(manager.getMaxInflightRebalanceTasks()).isZero();
+ assertThat(eventProcessor.executedPlans).isEmpty();
+ applyLatestMaxInflightTasksChangedEvent(eventManager, manager, 2);
+
+ assertThat(manager.getMaxInflightRebalanceTasks()).isEqualTo(2);
+ assertThat(eventProcessor.executedPlans).hasSize(2);
+ assertThat(countStatus(manager,
RebalanceStatus.REBALANCING)).isEqualTo(2);
+ assertThat(countStatus(manager,
RebalanceStatus.NOT_STARTED)).isEqualTo(2);
+
+ manager.close();
+ }
+
+ @Test
+ void
testDecreaseMaxInflightRebalanceTasksToZeroDoesNotCancelRunningTasks() throws
Exception {
+ ManualClock clock = new ManualClock(0L);
+ RecordingEventManager eventManager = new RecordingEventManager();
+ NoOpScheduledExecutor executor = new NoOpScheduledExecutor();
+ Configuration conf = new Configuration();
+ conf.set(ConfigOptions.COORDINATOR_REBALANCE_MAX_INFLIGHT_TASKS, 3);
+ RecordingCoordinatorEventProcessor eventProcessor =
+ buildRecordingCoordinatorEventProcessor(conf);
+ RebalanceManager manager =
+ new RebalanceManager(
+ eventProcessor, zookeeperClient, eventManager, clock,
conf, executor);
+
+ Map<TableBucket, RebalancePlanForBucket> plan = createRebalancePlan(5);
+ List<TableBucket> buckets = new ArrayList<>(plan.keySet());
+ zookeeperClient.registerRebalanceTask(
+ new RebalanceTask("decrease-inflight-test", NOT_STARTED,
plan));
+ manager.registerRebalance("decrease-inflight-test", plan, NOT_STARTED);
+ assertThat(eventProcessor.executedPlans).hasSize(3);
+
+ Configuration newConfig = new Configuration(conf);
+ newConfig.set(ConfigOptions.COORDINATOR_REBALANCE_MAX_INFLIGHT_TASKS,
0);
+ manager.reconfigure(newConfig);
+
+ assertThat(manager.getMaxInflightRebalanceTasks()).isEqualTo(3);
+ applyLatestMaxInflightTasksChangedEvent(eventManager, manager, 0);
+
+ assertThat(manager.getMaxInflightRebalanceTasks()).isZero();
+ assertThat(eventProcessor.executedPlans).hasSize(3);
+ manager.finishRebalanceTask(buckets.get(0), COMPLETED);
+ manager.finishRebalanceTask(buckets.get(1), COMPLETED);
+ assertThat(eventProcessor.executedPlans).hasSize(3);
+
+ manager.finishRebalanceTask(buckets.get(2), COMPLETED);
+
+ assertThat(eventProcessor.executedPlans).hasSize(3);
+ assertThat(countStatus(manager,
RebalanceStatus.NOT_STARTED)).isEqualTo(2);
+ assertThat(countStatus(manager, RebalanceStatus.REBALANCING)).isZero();
+
+ newConfig.set(ConfigOptions.COORDINATOR_REBALANCE_MAX_INFLIGHT_TASKS,
1);
+ manager.reconfigure(newConfig);
+
+ assertThat(eventProcessor.executedPlans).hasSize(3);
+ applyLatestMaxInflightTasksChangedEvent(eventManager, manager, 1);
+
+ assertThat(eventProcessor.executedPlans).hasSize(4);
+
assertThat(eventProcessor.executedPlans.get(3).getTableBucket()).isEqualTo(buckets.get(3));
+ assertThat(countStatus(manager,
RebalanceStatus.REBALANCING)).isEqualTo(1);
+
+ manager.close();
+ }
+
+ @Test
+ void testTimeoutEnqueuesEventsForAllInflightTasks() throws Exception {
+ ManualClock clock = new ManualClock(0L);
+ RecordingEventManager eventManager = new RecordingEventManager();
+ NoOpScheduledExecutor executor = new NoOpScheduledExecutor();
+ Configuration conf = new Configuration();
+ conf.set(ConfigOptions.COORDINATOR_REBALANCE_MAX_INFLIGHT_TASKS, 2);
+ RecordingCoordinatorEventProcessor eventProcessor =
+ buildRecordingCoordinatorEventProcessor(conf);
+ RebalanceManager manager =
+ new RebalanceManager(
+ eventProcessor, zookeeperClient, eventManager, clock,
conf, executor);
+
+ Map<TableBucket, RebalancePlanForBucket> plan = createRebalancePlan(3);
+ List<TableBucket> buckets = new ArrayList<>(plan.keySet());
+ zookeeperClient.registerRebalanceTask(
+ new RebalanceTask("timeout-all-inflight-test", NOT_STARTED,
plan));
+ manager.registerRebalance("timeout-all-inflight-test", plan,
NOT_STARTED);
+ assertThat(eventProcessor.executedPlans).hasSize(2);
+
+ clock.advanceTime(Duration.ofMillis(130_000));
+ manager.checkTimeout();
+
+ assertThat(eventManager.events).hasSize(2);
+ assertThat(((RebalanceTaskTimeoutEvent)
eventManager.events.get(0)).getTableBucket())
+ .isEqualTo(buckets.get(0));
+ assertThat(((RebalanceTaskTimeoutEvent)
eventManager.events.get(1)).getTableBucket())
+ .isEqualTo(buckets.get(1));
Review Comment:
This test assumes `checkTimeout()` enqueues timeout events in a stable order
(bucket 0 then bucket 1). After the change, `checkTimeout()` iterates a
`HashMap` snapshot of `inflightRebalanceTaskStartMs`, so the enqueue order is
not deterministic when multiple tasks are in-flight; this can make the test
flaky.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]