diff --git a/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java b/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java index 16301f4612f..ebcff628e27 100644 --- a/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java +++ b/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java @@ -405,7 +405,7 @@ public class ConfigOptions { public static final ConfigOption SERVER_IO_POOL_SIZE = key("server.io-pool.size") .intType() - .defaultValue(10) + .defaultValue(2) .withDescription( "The size of the IO thread pool to run blocking operations for both coordinator and tablet servers. " + "This includes discard unnecessary snapshot files, transfer kv snapshot files, " @@ -649,6 +649,14 @@ public class ConfigOptions { "The rack for the tabletServer. This will be used in rack aware bucket assignment " + "for fault tolerance. Examples: `RACK1`, `cn-hangzhou-server10`"); + public static final ConfigOption TABLET_SERVER_REPLICA_TRANSITION_THREAD_NUM = + key("tablet-server.replica-transition-thread-num") + .intType() + .defaultValue(10) + .withDescription( + "The maximum number of replica role transitions that can run " + + "concurrently in a TabletServer."); + public static final ConfigOption TABLET_SERVER_ADVERTISED_RESOURCE_CPU_CORES = key("tablet-server.advertised-resource.cpu-cores") .doubleType() diff --git a/fluss-common/src/main/java/org/apache/fluss/config/FlussConfigUtils.java b/fluss-common/src/main/java/org/apache/fluss/config/FlussConfigUtils.java index 9a7ff3cd50e..0658dd5fde1 100644 --- a/fluss-common/src/main/java/org/apache/fluss/config/FlussConfigUtils.java +++ b/fluss-common/src/main/java/org/apache/fluss/config/FlussConfigUtils.java @@ -134,6 +134,7 @@ public static void validateTabletConfigs(Configuration conf) { "Configuration %s must be set.", ConfigOptions.TABLET_SERVER_ID.key())); } validMinValue(ConfigOptions.TABLET_SERVER_ID, serverId.get(), 0); + validMinValue(conf, ConfigOptions.TABLET_SERVER_REPLICA_TRANSITION_THREAD_NUM, 1); } public static void validateRemoteDataDirs(Configuration conf) { diff --git a/fluss-common/src/test/java/org/apache/fluss/config/FlussConfigUtilsTest.java b/fluss-common/src/test/java/org/apache/fluss/config/FlussConfigUtilsTest.java index 61a00c01d2e..fc8b3e9dd04 100644 --- a/fluss-common/src/test/java/org/apache/fluss/config/FlussConfigUtilsTest.java +++ b/fluss-common/src/test/java/org/apache/fluss/config/FlussConfigUtilsTest.java @@ -189,6 +189,14 @@ void testValidateTabletConfigs() { .isInstanceOf(IllegalConfigurationException.class) .hasMessageContaining(ConfigOptions.TABLET_SERVER_ID.key()) .hasMessageContaining("it must be greater than or equal 0"); + + conf.set(ConfigOptions.TABLET_SERVER_ID, 0); + conf.set(ConfigOptions.TABLET_SERVER_REPLICA_TRANSITION_THREAD_NUM, 0); + assertThatThrownBy(() -> validateTabletConfigs(conf)) + .isInstanceOf(IllegalConfigurationException.class) + .hasMessageContaining( + ConfigOptions.TABLET_SERVER_REPLICA_TRANSITION_THREAD_NUM.key()) + .hasMessageContaining("must be greater than or equal 1"); } @Test diff --git a/fluss-server/src/main/java/org/apache/fluss/server/replica/ReplicaManager.java b/fluss-server/src/main/java/org/apache/fluss/server/replica/ReplicaManager.java index b3357019230..441f6220bd5 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/replica/ReplicaManager.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/replica/ReplicaManager.java @@ -145,12 +145,15 @@ import java.util.Collections; import java.util.HashMap; import java.util.HashSet; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Optional; import java.util.Set; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; @@ -196,6 +199,7 @@ public class ReplicaManager implements ServerReconfigurable { private final TabletServerMetadataCache metadataCache; private final ExecutorService ioExecutor; + private final ExecutorService replicaTransitionExecutor; private final ProjectionPushdownCache projectionsCache = new ProjectionPushdownCache(); private final Lock replicaStateChangeLock = new ReentrantLock(); @@ -259,6 +263,7 @@ public ReplicaManager( ScannerManager scannerManager, Clock clock, ExecutorService ioExecutor, + ExecutorService replicaTransitionExecutor, LocalDiskManager localDiskManager, @Nullable PluginManager pluginManager) throws IOException { @@ -287,6 +292,7 @@ public ReplicaManager( scannerManager, clock, ioExecutor, + replicaTransitionExecutor, localDiskManager, pluginManager); } @@ -310,6 +316,7 @@ public ReplicaManager( ScannerManager scannerManager, Clock clock, ExecutorService ioExecutor, + ExecutorService replicaTransitionExecutor, LocalDiskManager localDiskManager, @Nullable PluginManager pluginManager) throws IOException { @@ -361,6 +368,7 @@ public ReplicaManager( this.userMetrics = userMetrics; this.clock = clock; this.ioExecutor = ioExecutor; + this.replicaTransitionExecutor = replicaTransitionExecutor; this.minInSyncReplicas = conf.get(ConfigOptions.LOG_REPLICA_MIN_IN_SYNC_REPLICAS_NUMBER); this.scannerManager = checkNotNull(scannerManager, "scannerManager"); // Historical lookup cache capacity currently uses only the first data volume. @@ -567,12 +575,18 @@ public void becomeLeaderOrFollower( inLock( replicaStateChangeLock, () -> { + Map dataByTableBucket = + new LinkedHashMap<>(); + for (NotifyLeaderAndIsrData data : notifyLeaderAndIsrDataList) { + dataByTableBucket.put(data.getTableBucket(), data); + } + // check or apply coordinator epoch. validateAndApplyCoordinatorEpoch(requestCoordinatorEpoch, "notifyLeaderAndIsr"); List replicasToBeLeader = new ArrayList<>(); List replicasToBeFollower = new ArrayList<>(); - for (NotifyLeaderAndIsrData data : notifyLeaderAndIsrDataList) { + for (NotifyLeaderAndIsrData data : dataByTableBucket.values()) { LOG.info( "Try to become leaderAndFollower for {} with isr {}, replicas: {}", data.getTableBucket(), @@ -1389,30 +1403,62 @@ private void makeLeaders( .map(NotifyLeaderAndIsrData::getTableBucket) .collect(Collectors.toSet())); + List> makeLeaderFutures = + new ArrayList<>(replicasToBeLeader.size()); for (NotifyLeaderAndIsrData data : replicasToBeLeader) { TableBucket tb = data.getTableBucket(); try { Replica replica = getReplicaOrException(tb); - // register replica to remote log manager first. - remoteLogManager.registerReplica(replica); - - // Load the latest lake progress before leader activation. Historical KV recovery - // requires its lake log end offset, while failures remain best effort for normal - // replicas. - if (replica.isDataLakeEnabled()) { - updateWithLakeTableSnapshot(replica); - } - replica.makeLeader(data); - - // start the remote log tiering tasks for leaders - remoteLogManager.startLogTiering(replica); - result.put(tb, new NotifyLeaderAndIsrResultForBucket(tb)); + makeLeaderFutures.add( + CompletableFuture.supplyAsync( + () -> makeLeader(replica, data), replicaTransitionExecutor)); } catch (Exception e) { LOG.error("Error make replica {} to leader", tb, e); result.put( tb, new NotifyLeaderAndIsrResultForBucket(tb, ApiError.fromThrowable(e))); } } + + try { + CompletableFuture.allOf( + makeLeaderFutures.toArray( + new CompletableFuture[makeLeaderFutures.size()])) + .get(); + } catch (InterruptedException e) { + makeLeaderFutures.forEach(future -> future.cancel(false)); + Thread.currentThread().interrupt(); + throw new CompletionException(e); + } catch (ExecutionException e) { + throw new CompletionException(e.getCause()); + } + for (CompletableFuture future : makeLeaderFutures) { + NotifyLeaderAndIsrResultForBucket leaderResult = future.join(); + result.put(leaderResult.getTableBucket(), leaderResult); + } + } + + private NotifyLeaderAndIsrResultForBucket makeLeader( + Replica replica, NotifyLeaderAndIsrData data) { + TableBucket tb = data.getTableBucket(); + try { + // register replica to remote log manager first. + remoteLogManager.registerReplica(replica); + + // Load the latest lake progress before leader activation. Historical KV recovery + // requires its lake log end offset, while failures remain best effort for normal + // replicas. + if (replica.isDataLakeEnabled()) { + updateWithLakeTableSnapshot(replica); + } + replica.makeLeader(data); + + // start the remote log tiering tasks for leaders + remoteLogManager.startLogTiering(replica); + return new NotifyLeaderAndIsrResultForBucket(tb); + } catch (Exception e) { + LOG.error("Error make replica {} to leader", tb, e); + return new NotifyLeaderAndIsrResultForBucket(tb, ApiError.fromThrowable(e)); + } } // NOTE: This method can be removed when fetchFromLake is deprecated diff --git a/fluss-server/src/main/java/org/apache/fluss/server/tablet/TabletServer.java b/fluss-server/src/main/java/org/apache/fluss/server/tablet/TabletServer.java index 45dfacd09ff..352d58f21ec 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/tablet/TabletServer.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/tablet/TabletServer.java @@ -189,6 +189,9 @@ public class TabletServer extends ServerBase { @GuardedBy("lock") private ExecutorService replicaStateChangeExecutor; + @GuardedBy("lock") + private ExecutorService replicaTransitionExecutor; + public TabletServer(Configuration conf) { this(conf, SystemClock.getInstance()); } @@ -285,6 +288,11 @@ protected void startServices() throws Exception { this.replicaStateChangeExecutor = Executors.newSingleThreadExecutor( new ExecutorThreadFactory("tablet-server-replica-state-change")); + this.replicaTransitionExecutor = + Executors.newFixedThreadPool( + conf.get(ConfigOptions.TABLET_SERVER_REPLICA_TRANSITION_THREAD_NUM), + new ExecutorThreadFactory( + "tablet-server-replica-transition-" + serverId)); this.scannerManager = new ScannerManager(conf, scheduler); @@ -307,6 +315,7 @@ protected void startServices() throws Exception { scannerManager, clock, ioExecutor, + replicaTransitionExecutor, localDiskManager, pluginManager); replicaManager.startup(); @@ -478,6 +487,14 @@ CompletableFuture stopServices() { exception = ExceptionUtils.firstOrSuppressed(t, exception); } + try { + if (replicaTransitionExecutor != null) { + ExecutorUtils.gracefulShutdown(5, TimeUnit.SECONDS, replicaTransitionExecutor); + } + } catch (Throwable t) { + exception = ExceptionUtils.firstOrSuppressed(t, exception); + } + try { if (replicaStateChangeExecutor != null) { ExecutorUtils.gracefulShutdown(5, TimeUnit.SECONDS, replicaStateChangeExecutor); diff --git a/fluss-server/src/test/java/org/apache/fluss/server/replica/ReplicaManagerTest.java b/fluss-server/src/test/java/org/apache/fluss/server/replica/ReplicaManagerTest.java index db4a88bc085..427eb44b503 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/replica/ReplicaManagerTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/replica/ReplicaManagerTest.java @@ -119,6 +119,7 @@ import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; import static org.apache.fluss.config.ConfigOptions.KV_FORMAT_VERSION_2; @@ -1747,21 +1748,20 @@ void becomeLeaderOrFollower() throws Exception { // make tb as leader. CompletableFuture> future = new CompletableFuture<>(); - replicaManager.becomeLeaderOrFollower( - INITIAL_COORDINATOR_EPOCH, - Collections.singletonList( - new NotifyLeaderAndIsrData( - PhysicalTablePath.of(DATA1_TABLE_PATH), - tb, + NotifyLeaderAndIsrData data = + new NotifyLeaderAndIsrData( + PhysicalTablePath.of(DATA1_TABLE_PATH), + tb, + Arrays.asList(1, 2, 3), + new LeaderAndIsr( + TABLET_SERVER_ID, + 1, Arrays.asList(1, 2, 3), - new LeaderAndIsr( - TABLET_SERVER_ID, - 1, - Arrays.asList(1, 2, 3), - Collections.emptyList(), - INITIAL_COORDINATOR_EPOCH, - INITIAL_BUCKET_EPOCH))), - future::complete); + Collections.emptyList(), + INITIAL_COORDINATOR_EPOCH, + INITIAL_BUCKET_EPOCH)); + replicaManager.becomeLeaderOrFollower( + INITIAL_COORDINATOR_EPOCH, Arrays.asList(data, data), future::complete); assertThat(future.get()).containsOnly(new NotifyLeaderAndIsrResultForBucket(tb)); assertReplicaEpochEquals( replicaManager.getReplicaOrException(tb), true, 1, INITIAL_BUCKET_EPOCH); @@ -1796,6 +1796,94 @@ void becomeLeaderOrFollower() throws Exception { replicaManager.getReplicaOrException(tb), true, 1, INITIAL_BUCKET_EPOCH); } + @Test + void testMakeLeadersWithPartialFailure() throws Exception { + TableBucket firstBucket = new TableBucket(DATA1_TABLE_ID, 1); + TableBucket secondBucket = new TableBucket(DATA1_TABLE_ID, 2); + List leaderData = + Arrays.asList( + newLeaderData(firstBucket, INITIAL_BUCKET_EPOCH), + newLeaderData(secondBucket, INITIAL_BUCKET_EPOCH)); + makeLogTableAsLeader(firstBucket.getBucket()); + replicaManager + .getReplicaOrException(firstBucket) + .makeLeader(newLeaderData(firstBucket, INITIAL_BUCKET_EPOCH + 1)); + + CompletableFuture> future = + new CompletableFuture<>(); + replicaManager.becomeLeaderOrFollower( + INITIAL_COORDINATOR_EPOCH, leaderData, future::complete); + + assertThat(future.get()) + .containsExactlyInAnyOrder( + new NotifyLeaderAndIsrResultForBucket( + firstBucket, + new ApiError( + Errors.INVALID_UPDATE_VERSION_EXCEPTION, + String.format( + "Skipped the become-leader state change for %s with a lower bucket epoch %s" + + " since the leader is already at a newer bucket epoch %s", + firstBucket, + INITIAL_BUCKET_EPOCH, + INITIAL_BUCKET_EPOCH + 1))), + new NotifyLeaderAndIsrResultForBucket(secondBucket)); + assertReplicaEpochEquals( + replicaManager.getReplicaOrException(firstBucket), + true, + INITIAL_LEADER_EPOCH, + INITIAL_BUCKET_EPOCH + 1); + assertReplicaEpochEquals( + replicaManager.getReplicaOrException(secondBucket), + true, + INITIAL_LEADER_EPOCH, + INITIAL_BUCKET_EPOCH); + } + + @Test + void testInterruptLeaderTransitionWaitWithQueuedTasks() throws Exception { + ThreadPoolExecutor transitionExecutor = (ThreadPoolExecutor) replicaTransitionExecutor; + transitionExecutor.setCorePoolSize(1); + transitionExecutor.setMaximumPoolSize(1); + CountDownLatch workerStarted = new CountDownLatch(1); + ExecutorService stateChangeExecutor = Executors.newSingleThreadExecutor(); + List queuedTasks = new ArrayList<>(); + try { + transitionExecutor.submit( + () -> { + workerStarted.countDown(); + Thread.sleep(Long.MAX_VALUE); + return null; + }); + assertThat(workerStarted.await(10, TimeUnit.SECONDS)).isTrue(); + stateChangeExecutor.submit( + () -> + replicaManager.becomeLeaderOrFollower( + INITIAL_COORDINATOR_EPOCH, + Arrays.asList( + newLeaderData( + new TableBucket(DATA1_TABLE_ID, 1), + INITIAL_BUCKET_EPOCH), + newLeaderData( + new TableBucket(DATA1_TABLE_ID, 2), + INITIAL_BUCKET_EPOCH)), + ignored -> {})); + waitUntil( + () -> transitionExecutor.getQueue().size() == 2, + Duration.ofSeconds(10), + "Leader transitions were not queued"); + + queuedTasks.addAll(transitionExecutor.shutdownNow()); + stateChangeExecutor.shutdownNow(); + assertThat(stateChangeExecutor.awaitTermination(10, TimeUnit.SECONDS)).isTrue(); + } finally { + queuedTasks.addAll(transitionExecutor.shutdownNow()); + queuedTasks.forEach(Runnable::run); + stateChangeExecutor.shutdownNow(); + stateChangeExecutor.awaitTermination(10, TimeUnit.SECONDS); + } + assertThat(replicaManager.leaderCount()).isZero(); + } + @Test void testLakeSnapshotReadFailureDoesNotFailLeaderTransition() throws Exception { TablePath tablePath = TablePath.of("test_db", "lake_table"); @@ -2454,6 +2542,21 @@ void testGetReplicaOrException() { .isInstanceOf(UnknownTableOrBucketException.class); } + private NotifyLeaderAndIsrData newLeaderData(TableBucket tableBucket, int bucketEpoch) { + List replicas = Arrays.asList(1, 2, 3); + return new NotifyLeaderAndIsrData( + PhysicalTablePath.of(DATA1_TABLE_PATH), + tableBucket, + replicas, + new LeaderAndIsr( + TABLET_SERVER_ID, + INITIAL_LEADER_EPOCH, + replicas, + Collections.emptyList(), + INITIAL_COORDINATOR_EPOCH, + bucketEpoch)); + } + private void assertReplicaEpochEquals( Replica replica, boolean isLeader, int leaderEpoch, int bucketEpoch) { assertThat(replica.isLeader()).isEqualTo(isLeader); diff --git a/fluss-server/src/test/java/org/apache/fluss/server/replica/ReplicaTestBase.java b/fluss-server/src/test/java/org/apache/fluss/server/replica/ReplicaTestBase.java index 98d714edc1e..6f0c830b713 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/replica/ReplicaTestBase.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/replica/ReplicaTestBase.java @@ -68,6 +68,7 @@ import org.apache.fluss.testutils.common.AllCallbackWrapper; import org.apache.fluss.testutils.common.ManuallyTriggeredScheduledExecutorService; import org.apache.fluss.utils.CloseableRegistry; +import org.apache.fluss.utils.ExecutorUtils; import org.apache.fluss.utils.clock.ManualClock; import org.apache.fluss.utils.concurrent.FlussScheduler; import org.apache.fluss.utils.function.FunctionWithException; @@ -159,6 +160,7 @@ public class ReplicaTestBase { protected TestCoordinatorGateway testCoordinatorGateway; private FlussScheduler scheduler; private ExecutorService ioExecutor; + protected ExecutorService replicaTransitionExecutor; // remote log related protected TestingRemoteLogStorage remoteLogStorage; @@ -219,6 +221,9 @@ public void setup(TestInfo testInfo) throws Exception { scheduler = new FlussScheduler(2); scheduler.startup(); ioExecutor = Executors.newSingleThreadExecutor(); + replicaTransitionExecutor = + Executors.newFixedThreadPool( + conf.get(ConfigOptions.TABLET_SERVER_REPLICA_TRANSITION_THREAD_NUM)); manualClock = new ManualClock(System.currentTimeMillis()); localDiskManager = LocalDiskManager.create(conf); @@ -374,6 +379,7 @@ protected ReplicaManager buildReplicaManager(CoordinatorGateway coordinatorGatew scannerManager, manualClock, ioExecutor, + replicaTransitionExecutor, localDiskManager, null); } @@ -382,6 +388,10 @@ protected ReplicaManager buildReplicaManager(CoordinatorGateway coordinatorGatew void tearDown() throws Exception { closeableRegistry.close(); + if (replicaTransitionExecutor != null) { + ExecutorUtils.gracefulShutdown(5, TimeUnit.SECONDS, replicaTransitionExecutor); + } + if (logManager != null) { logManager.shutdown(); } diff --git a/fluss-server/src/test/java/org/apache/fluss/server/replica/fetcher/ReplicaFetcherThreadTest.java b/fluss-server/src/test/java/org/apache/fluss/server/replica/fetcher/ReplicaFetcherThreadTest.java index 8a1f6fb41e5..ace1a487791 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/replica/fetcher/ReplicaFetcherThreadTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/replica/fetcher/ReplicaFetcherThreadTest.java @@ -106,6 +106,7 @@ public class ReplicaFetcherThreadTest { private ReplicaManager followerRM; private ReplicaFetcherThread followerFetcher; private ExecutorService ioExecutor; + private ExecutorService replicaTransitionExecutor; private LocalDiskManager leaderLocalDiskManager; private LocalDiskManager followerLocalDiskManager; @@ -124,6 +125,7 @@ public void setup() throws Exception { Configuration conf = new Configuration(); tb = new TableBucket(DATA1_TABLE_ID, 0); ioExecutor = Executors.newSingleThreadExecutor(); + replicaTransitionExecutor = Executors.newFixedThreadPool(2); leaderLocalDiskManager = createLocalDiskManager(leaderServerId); leaderRM = createReplicaManager(leaderServerId, leaderLocalDiskManager); followerLocalDiskManager = createLocalDiskManager(followerServerId); @@ -158,6 +160,9 @@ public void tearDown() throws Exception { if (ioExecutor != null) { ioExecutor.shutdownNow(); } + if (replicaTransitionExecutor != null) { + replicaTransitionExecutor.shutdownNow(); + } } @Test @@ -552,6 +557,7 @@ private ReplicaManager createReplicaManager(int serverId, LocalDiskManager local TestingMetricGroups.TABLET_SERVER_METRICS, manualClock, ioExecutor, + replicaTransitionExecutor, localDiskManager); replicaManager.startup(); return replicaManager; @@ -574,6 +580,7 @@ public TestingReplicaManager( TabletServerMetricGroup serverMetricGroup, Clock clock, ExecutorService ioExecutor, + ExecutorService replicaTransitionExecutor, LocalDiskManager localDiskManager) throws IOException { super( @@ -593,6 +600,7 @@ public TestingReplicaManager( new ScannerManager(conf, scheduler), clock, ioExecutor, + replicaTransitionExecutor, localDiskManager, null); }