fhan688 commented on code in PR #3622:
URL: https://github.com/apache/fluss/pull/3622#discussion_r3548703262
##########
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:
Agreed. Fixed the test to avoid depending on the timeout event order. It now
asserts the timed-out buckets with containsExactlyInAnyOrder, matching the
unordered iteration behavior in checkTimeout().
--
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]