From 2aad05fef6bbdc5b090914752b3e10a19245d80d Mon Sep 17 00:00:00 2001 From: morazow Date: Tue, 21 Jul 2026 17:47:46 +0200 Subject: [PATCH 1/3] [server][rpc] Support DescribeTabletServers API for safe scaling and upgrades Add a read-only DescribeTabletServers RPC to the Coordinator, following the pattern of GetClusterHealth (#3399). It returns, per tablet server, the same four counters getClusterHealth() reports cluster-wide, scoped to that server: numReplicas, inSyncReplicas, numLeaderReplicas and activeLeaderReplicas. This lets operational tooling (e.g. the Fluss Kubernetes Operator): - gate scale-in on a server being empty (numReplicas == 0), - gate rolling upgrades on a per-server green predicate (inSyncReplicas == numReplicas && activeLeaderReplicas == numLeaderReplicas), - report per-server tablet load in cluster status. The counters are computed from in-memory CoordinatorContext state on the coordinator event thread via AccessContextEvent. Tablet servers forward the request to the coordinator. Live and shutting-down servers are always reported, even with zero replicas, and dead servers that still hold assigned replicas remain visible so a scale-in gate cannot treat them as drained. Closes #3570 --- .../org/apache/fluss/client/admin/Admin.java | 23 ++ .../apache/fluss/client/admin/FlussAdmin.java | 7 + .../client/admin/TabletServerDescription.java | 130 ++++++++++ .../client/utils/ClientRpcMessageUtils.java | 16 ++ .../fluss/client/admin/FlussAdminITCase.java | 81 ++++++ .../sink/testutils/TestAdminAdapter.java | 6 + .../rpc/gateway/AdminReadOnlyGateway.java | 11 + .../apache/fluss/rpc/protocol/ApiKeys.java | 3 +- fluss-rpc/src/main/proto/FlussApi.proto | 15 ++ .../rpc/TestingTabletGatewayService.java | 8 + .../coordinator/CoordinatorService.java | 78 ++++++ .../fluss/server/tablet/TabletService.java | 25 ++ .../DescribeTabletServersTest.java | 239 ++++++++++++++++++ .../coordinator/TestCoordinatorGateway.java | 8 + .../tablet/TestTabletServerGateway.java | 8 + 15 files changed, 657 insertions(+), 1 deletion(-) create mode 100644 fluss-client/src/main/java/org/apache/fluss/client/admin/TabletServerDescription.java create mode 100644 fluss-server/src/test/java/org/apache/fluss/server/coordinator/DescribeTabletServersTest.java diff --git a/fluss-client/src/main/java/org/apache/fluss/client/admin/Admin.java b/fluss-client/src/main/java/org/apache/fluss/client/admin/Admin.java index 5d749b6d432..dcc368bfefb 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/admin/Admin.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/admin/Admin.java @@ -799,6 +799,29 @@ CompletableFuture registerProducerOffsets( */ CompletableFuture getClusterHealth(); + /** + * Describe the replica load of each tablet server in the cluster asynchronously. + * + *

The returned list contains {@link TabletServerDescription} counters per tablet server. + * Live tablet servers are always included, even when they host no replicas. + * + *

This API is designed for operational tooling (e.g., a Kubernetes operator): + * + *

    + *
  • scale-in safety gate - a tablet server may only be removed once its {@code numReplicas} + * is {@code 0}, since terminating a non-empty server causes under-replication or data + * unavailability; + *
  • rolling-upgrade gate - only proceed to the next server when the previous one is green + * again ({@code inSyncReplicas == numReplicas && activeLeaderReplicas == + * numLeaderReplicas}); + *
  • reporting per-server tablet load in cluster status. + *
+ * + * @return a {@link CompletableFuture} that completes with the per tablet server replica loads. + * @since 1.0 + */ + CompletableFuture> describeTabletServers(); + /** * List per-bucket remote log manifest entries for a table or partition scope. * diff --git a/fluss-client/src/main/java/org/apache/fluss/client/admin/FlussAdmin.java b/fluss-client/src/main/java/org/apache/fluss/client/admin/FlussAdmin.java index a0909ceec73..dc9fcfa2b4e 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/admin/FlussAdmin.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/admin/FlussAdmin.java @@ -67,6 +67,7 @@ import org.apache.fluss.rpc.messages.DatabaseExistsResponse; import org.apache.fluss.rpc.messages.DeleteProducerOffsetsRequest; import org.apache.fluss.rpc.messages.DescribeClusterConfigsRequest; +import org.apache.fluss.rpc.messages.DescribeTabletServersRequest; import org.apache.fluss.rpc.messages.DropAclsRequest; import org.apache.fluss.rpc.messages.DropDatabaseRequest; import org.apache.fluss.rpc.messages.DropTableRequest; @@ -942,6 +943,12 @@ public CompletableFuture getClusterHealth() { .thenApply(ClientRpcMessageUtils::toClusterHealth); } + @Override + public CompletableFuture> describeTabletServers() { + return gateway.describeTabletServers(new DescribeTabletServersRequest()) + .thenApply(ClientRpcMessageUtils::toTabletServerDescriptions); + } + @VisibleForTesting public AdminGateway getAdminGateway() { return gateway; diff --git a/fluss-client/src/main/java/org/apache/fluss/client/admin/TabletServerDescription.java b/fluss-client/src/main/java/org/apache/fluss/client/admin/TabletServerDescription.java new file mode 100644 index 00000000000..a4605e38302 --- /dev/null +++ b/fluss-client/src/main/java/org/apache/fluss/client/admin/TabletServerDescription.java @@ -0,0 +1,130 @@ +/* + * 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.client.admin; + +import org.apache.fluss.annotation.PublicEvolving; + +import java.util.Objects; + +/** + * Per tablet server replica load returned by {@link Admin#describeTabletServers()}. + * + *

It reports the replica counters of a single tablet server: how many replicas it hosts, how + * many of those are in sync, how many buckets it leads, and how many of those leaders are active. + * The {@link ClusterHealth} aggregates the same counters across all tablet servers. + * + *

Operational tooling can derive from it: + * + *

    + *
  • a scale-in safety gate: the server is empty and safe to remove iff {@code numReplicas == + * 0}; + *
  • a per-server health (green) predicate for rolling upgrades: {@code inSyncReplicas == + * numReplicas && activeLeaderReplicas == numLeaderReplicas}. + *
+ * + *

Note: a bucket without an elected leader is attributed to no server's {@code + * numLeaderReplicas}, so the per-server values sum to the cluster-wide {@link + * ClusterHealth#getNumLeaderReplicas()} only when every bucket has an elected leader. Such buckets + * still surface through {@code inSyncReplicas < numReplicas} on the assigned servers. + * + * @since 1.0 + */ +@PublicEvolving +public final class TabletServerDescription { + + private final int serverId; + private final int numReplicas; + private final int inSyncReplicas; + private final int numLeaderReplicas; + private final int activeLeaderReplicas; + + public TabletServerDescription( + int serverId, + int numReplicas, + int inSyncReplicas, + int numLeaderReplicas, + int activeLeaderReplicas) { + this.serverId = serverId; + this.numReplicas = numReplicas; + this.inSyncReplicas = inSyncReplicas; + this.numLeaderReplicas = numLeaderReplicas; + this.activeLeaderReplicas = activeLeaderReplicas; + } + + public int getServerId() { + return serverId; + } + + /** Number of replicas assigned to this tablet server. */ + public int getNumReplicas() { + return numReplicas; + } + + /** Number of assigned replicas that are in the ISR of their bucket. */ + public int getInSyncReplicas() { + return inSyncReplicas; + } + + /** Number of hosted replicas that are the elected leader of their bucket. */ + public int getNumLeaderReplicas() { + return numLeaderReplicas; + } + + /** Number of leader replicas on this server that are currently active. */ + public int getActiveLeaderReplicas() { + return activeLeaderReplicas; + } + + @Override + public boolean equals(Object o) { + if (this == o) { + return true; + } + if (!(o instanceof TabletServerDescription)) { + return false; + } + TabletServerDescription that = (TabletServerDescription) o; + return serverId == that.serverId + && numReplicas == that.numReplicas + && inSyncReplicas == that.inSyncReplicas + && numLeaderReplicas == that.numLeaderReplicas + && activeLeaderReplicas == that.activeLeaderReplicas; + } + + @Override + public int hashCode() { + return Objects.hash( + serverId, numReplicas, inSyncReplicas, numLeaderReplicas, activeLeaderReplicas); + } + + @Override + public String toString() { + return "TabletServerDescription{" + + "serverId=" + + serverId + + ", numReplicas=" + + numReplicas + + ", inSyncReplicas=" + + inSyncReplicas + + ", numLeaderReplicas=" + + numLeaderReplicas + + ", activeLeaderReplicas=" + + activeLeaderReplicas + + '}'; + } +} diff --git a/fluss-client/src/main/java/org/apache/fluss/client/utils/ClientRpcMessageUtils.java b/fluss-client/src/main/java/org/apache/fluss/client/utils/ClientRpcMessageUtils.java index 70f940c3cd2..24ac21aca8f 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/utils/ClientRpcMessageUtils.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/utils/ClientRpcMessageUtils.java @@ -21,6 +21,7 @@ import org.apache.fluss.client.admin.ClusterHealthStatus; import org.apache.fluss.client.admin.OffsetSpec; import org.apache.fluss.client.admin.ProducerOffsetsResult; +import org.apache.fluss.client.admin.TabletServerDescription; import org.apache.fluss.client.lookup.LookupBatch; import org.apache.fluss.client.lookup.PrefixLookupBatch; import org.apache.fluss.client.metadata.AcquireKvSnapshotLeaseResult; @@ -55,6 +56,7 @@ import org.apache.fluss.rpc.messages.AlterDatabaseRequest; import org.apache.fluss.rpc.messages.AlterTableRequest; import org.apache.fluss.rpc.messages.CreatePartitionRequest; +import org.apache.fluss.rpc.messages.DescribeTabletServersResponse; import org.apache.fluss.rpc.messages.DropPartitionRequest; import org.apache.fluss.rpc.messages.GetClusterHealthResponse; import org.apache.fluss.rpc.messages.GetFileSystemSecurityTokenResponse; @@ -898,4 +900,18 @@ private static ClusterHealthStatus toClusterHealthStatus(int pbStatus) { return ClusterHealthStatus.UNKNOWN; } } + + public static List toTabletServerDescriptions( + DescribeTabletServersResponse resp) { + return resp.getTabletServersList().stream() + .map( + load -> + new TabletServerDescription( + load.getServerId(), + load.getNumReplicas(), + load.getInSyncReplicas(), + load.getNumLeaderReplicas(), + load.getActiveLeaderReplicas())) + .collect(Collectors.toList()); + } } diff --git a/fluss-client/src/test/java/org/apache/fluss/client/admin/FlussAdminITCase.java b/fluss-client/src/test/java/org/apache/fluss/client/admin/FlussAdminITCase.java index 150b84e5679..f784706c37f 100644 --- a/fluss-client/src/test/java/org/apache/fluss/client/admin/FlussAdminITCase.java +++ b/fluss-client/src/test/java/org/apache/fluss/client/admin/FlussAdminITCase.java @@ -2860,4 +2860,85 @@ void testClusterHealthDuringRollingUpgrade() throws Exception { assertThat(afterRecovery.getNumLeaderReplicas()) .isEqualTo(afterRecovery.getActiveLeaderReplicas()); } + + @Test + void testDescribeTabletServersDuringRollingUpgrade() throws Exception { + TablePath tablePath = TablePath.of("test_db", "describe_tablet_servers_table"); + TableDescriptor tableDescriptor = + TableDescriptor.builder().schema(DEFAULT_SCHEMA).distributedBy(3, "id").build(); + long tableId = createTable(tablePath, tableDescriptor, true); + waitAllReplicasReady(tableId, 3); + + // Phase 1: Cluster is healthy - every live server is reported, hosts replicas of the + // created table (replication factor 3 on 3 servers) and is green. + List servers = admin.describeTabletServers().get(); + assertThat(servers).extracting(TabletServerDescription::getServerId).contains(0, 1, 2); + for (TabletServerDescription server : servers) { + assertThat(server.getNumReplicas()).isGreaterThan(0); + assertThat(isServerGreen(server)).isTrue(); + } + + // The per-server counters must sum up to the cluster-wide health counters. + ClusterHealth health = admin.getClusterHealth().get(); + assertThat(servers.stream().mapToInt(TabletServerDescription::getNumReplicas).sum()) + .isEqualTo(health.getNumReplicas()); + assertThat(servers.stream().mapToInt(TabletServerDescription::getInSyncReplicas).sum()) + .isEqualTo(health.getInSyncReplicas()); + assertThat(servers.stream().mapToInt(TabletServerDescription::getNumLeaderReplicas).sum()) + .isEqualTo(health.getNumLeaderReplicas()); + assertThat( + servers.stream() + .mapToInt(TabletServerDescription::getActiveLeaderReplicas) + .sum()) + .isEqualTo(health.getActiveLeaderReplicas()); + + // Phase 2: Stop one tablet server (simulate server crash during rolling upgrade). It must + // still be reported with its assigned replicas, but no longer green - an operator must + // not treat it as safe to remove. + int stoppedServerId = 0; + FLUSS_CLUSTER_EXTENSION.stopTabletServer(stoppedServerId); + FLUSS_CLUSTER_EXTENSION.assertHasTabletServerNumber(2); + + for (int bucket = 0; bucket < 3; bucket++) { + TableBucket tb = new TableBucket(tableId, bucket); + FLUSS_CLUSTER_EXTENSION.waitUntilReplicaShrinkFromIsr(tb, stoppedServerId); + } + + TabletServerDescription stopped = + getTabletServerDescription(admin.describeTabletServers().get(), stoppedServerId); + assertThat(stopped.getNumReplicas()).isGreaterThan(0); + assertThat(stopped.getInSyncReplicas()).isLessThan(stopped.getNumReplicas()); + + // Phase 3: Restart the server and wait until every server is green again. + FLUSS_CLUSTER_EXTENSION.startTabletServer(stoppedServerId); + FLUSS_CLUSTER_EXTENSION.assertHasTabletServerNumber(3); + + for (int bucket = 0; bucket < 3; bucket++) { + TableBucket tb = new TableBucket(tableId, bucket); + FLUSS_CLUSTER_EXTENSION.waitUntilReplicaExpandToIsr(tb, stoppedServerId); + } + + waitUntil( + () -> + admin.describeTabletServers().get().stream() + .allMatch(FlussAdminITCase::isServerGreen), + Duration.ofMinutes(1), + "All tablet servers should become green again after server restart"); + } + + private static TabletServerDescription getTabletServerDescription( + List servers, int serverId) { + return servers.stream() + .filter(server -> server.getServerId() == serverId) + .findFirst() + .orElseThrow( + () -> + new AssertionError( + "no description reported for tablet server " + serverId)); + } + + private static boolean isServerGreen(TabletServerDescription server) { + return server.getInSyncReplicas() == server.getNumReplicas() + && server.getActiveLeaderReplicas() == server.getNumLeaderReplicas(); + } } diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/sink/testutils/TestAdminAdapter.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/sink/testutils/TestAdminAdapter.java index e8120c83d26..a0ca360b11f 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/sink/testutils/TestAdminAdapter.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/sink/testutils/TestAdminAdapter.java @@ -26,6 +26,7 @@ import org.apache.fluss.client.admin.OffsetSpec; import org.apache.fluss.client.admin.ProducerOffsetsResult; import org.apache.fluss.client.admin.RegisterResult; +import org.apache.fluss.client.admin.TabletServerDescription; import org.apache.fluss.client.metadata.ActiveKvSnapshots; import org.apache.fluss.client.metadata.KvSnapshotMetadata; import org.apache.fluss.client.metadata.KvSnapshots; @@ -326,6 +327,11 @@ public CompletableFuture getClusterHealth() { throw new UnsupportedOperationException("Not implemented in TestAdminAdapter"); } + @Override + public CompletableFuture> describeTabletServers() { + throw new UnsupportedOperationException("Not implemented in TestAdminAdapter"); + } + @Override public CompletableFuture> listRemoteLogManifests( long tableId, @Nullable Long partitionId) { diff --git a/fluss-rpc/src/main/java/org/apache/fluss/rpc/gateway/AdminReadOnlyGateway.java b/fluss-rpc/src/main/java/org/apache/fluss/rpc/gateway/AdminReadOnlyGateway.java index 574e1a510dd..00de401eddd 100644 --- a/fluss-rpc/src/main/java/org/apache/fluss/rpc/gateway/AdminReadOnlyGateway.java +++ b/fluss-rpc/src/main/java/org/apache/fluss/rpc/gateway/AdminReadOnlyGateway.java @@ -22,6 +22,8 @@ import org.apache.fluss.rpc.messages.DatabaseExistsResponse; import org.apache.fluss.rpc.messages.DescribeClusterConfigsRequest; import org.apache.fluss.rpc.messages.DescribeClusterConfigsResponse; +import org.apache.fluss.rpc.messages.DescribeTabletServersRequest; +import org.apache.fluss.rpc.messages.DescribeTabletServersResponse; import org.apache.fluss.rpc.messages.GetClusterHealthRequest; import org.apache.fluss.rpc.messages.GetClusterHealthResponse; import org.apache.fluss.rpc.messages.GetDatabaseInfoRequest; @@ -202,4 +204,13 @@ CompletableFuture describeClusterConfigs( */ @RPC(api = ApiKeys.GET_CLUSTER_HEALTH) CompletableFuture getClusterHealth(GetClusterHealthRequest request); + + /** + * Describe the replica load of each tablet server in the cluster. + * + * @return the describe tablet servers response. + */ + @RPC(api = ApiKeys.DESCRIBE_TABLET_SERVERS) + CompletableFuture describeTabletServers( + DescribeTabletServersRequest request); } diff --git a/fluss-rpc/src/main/java/org/apache/fluss/rpc/protocol/ApiKeys.java b/fluss-rpc/src/main/java/org/apache/fluss/rpc/protocol/ApiKeys.java index b03e1ec63ab..67587c582f7 100644 --- a/fluss-rpc/src/main/java/org/apache/fluss/rpc/protocol/ApiKeys.java +++ b/fluss-rpc/src/main/java/org/apache/fluss/rpc/protocol/ApiKeys.java @@ -106,7 +106,8 @@ public enum ApiKeys { SCAN_KV(1061, 0, 0, PUBLIC), GET_CLUSTER_HEALTH(1062, 0, 0, PUBLIC), LIST_REMOTE_LOG_MANIFESTS(1063, 0, 0, PUBLIC), - LIST_KV_SNAPSHOTS(1064, 0, 0, PUBLIC); + LIST_KV_SNAPSHOTS(1064, 0, 0, PUBLIC), + DESCRIBE_TABLET_SERVERS(1065, 0, 0, PUBLIC); private static final Map ID_TO_TYPE = Arrays.stream(ApiKeys.values()) diff --git a/fluss-rpc/src/main/proto/FlussApi.proto b/fluss-rpc/src/main/proto/FlussApi.proto index 57e96092d4b..258adc65246 100644 --- a/fluss-rpc/src/main/proto/FlussApi.proto +++ b/fluss-rpc/src/main/proto/FlussApi.proto @@ -807,8 +807,23 @@ message GetClusterHealthResponse { required int32 status = 5; // PbClusterHealthStatus: GREEN=0, YELLOW=1, RED=2, UNKNOWN=3 } +message DescribeTabletServersRequest { +} + +message DescribeTabletServersResponse { + repeated PbTabletServerLoad tablet_servers = 1; +} + // --------------- Inner classes ---------------- +message PbTabletServerLoad { + required int32 server_id = 1; + required int32 num_replicas = 2; + required int32 in_sync_replicas = 3; + required int32 num_leader_replicas = 4; + required int32 active_leader_replicas = 5; +} + message PbApiVersion { required int32 api_key = 1; required int32 min_version = 2; diff --git a/fluss-rpc/src/test/java/org/apache/fluss/rpc/TestingTabletGatewayService.java b/fluss-rpc/src/test/java/org/apache/fluss/rpc/TestingTabletGatewayService.java index f465bb4a69a..f146cffab71 100644 --- a/fluss-rpc/src/test/java/org/apache/fluss/rpc/TestingTabletGatewayService.java +++ b/fluss-rpc/src/test/java/org/apache/fluss/rpc/TestingTabletGatewayService.java @@ -23,6 +23,8 @@ import org.apache.fluss.rpc.messages.DatabaseExistsResponse; import org.apache.fluss.rpc.messages.DescribeClusterConfigsRequest; import org.apache.fluss.rpc.messages.DescribeClusterConfigsResponse; +import org.apache.fluss.rpc.messages.DescribeTabletServersRequest; +import org.apache.fluss.rpc.messages.DescribeTabletServersResponse; import org.apache.fluss.rpc.messages.FetchLogRequest; import org.apache.fluss.rpc.messages.FetchLogResponse; import org.apache.fluss.rpc.messages.GetClusterHealthRequest; @@ -272,4 +274,10 @@ public CompletableFuture getClusterHealth( GetClusterHealthRequest request) { return null; } + + @Override + public CompletableFuture describeTabletServers( + DescribeTabletServersRequest request) { + return null; + } } diff --git a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorService.java b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorService.java index ef6a735677f..5b99dc3b4ee 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorService.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorService.java @@ -90,6 +90,8 @@ import org.apache.fluss.rpc.messages.CreateTableResponse; import org.apache.fluss.rpc.messages.DeleteProducerOffsetsRequest; import org.apache.fluss.rpc.messages.DeleteProducerOffsetsResponse; +import org.apache.fluss.rpc.messages.DescribeTabletServersRequest; +import org.apache.fluss.rpc.messages.DescribeTabletServersResponse; import org.apache.fluss.rpc.messages.DropAclsRequest; import org.apache.fluss.rpc.messages.DropAclsResponse; import org.apache.fluss.rpc.messages.DropDatabaseRequest; @@ -123,6 +125,7 @@ import org.apache.fluss.rpc.messages.PbProducerTableOffsets; import org.apache.fluss.rpc.messages.PbTableBucket; import org.apache.fluss.rpc.messages.PbTableOffsets; +import org.apache.fluss.rpc.messages.PbTabletServerLoad; import org.apache.fluss.rpc.messages.PrepareLakeTableSnapshotRequest; import org.apache.fluss.rpc.messages.PrepareLakeTableSnapshotResponse; import org.apache.fluss.rpc.messages.RebalanceRequest; @@ -202,6 +205,7 @@ import java.util.Map; import java.util.Optional; import java.util.Set; +import java.util.TreeMap; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutorService; import java.util.function.Supplier; @@ -1496,6 +1500,80 @@ static GetClusterHealthResponse computeClusterHealth(CoordinatorContext ctx) { return response; } + @Override + public CompletableFuture describeTabletServers( + DescribeTabletServersRequest request) { + if (authorizer != null) { + authorizer.authorize(currentSession(), OperationType.DESCRIBE, Resource.cluster()); + } + + AccessContextEvent event = + new AccessContextEvent<>(CoordinatorService::describeTabletServers); + eventManagerSupplier.get().put(event); + return event.getResultFuture(); + } + + @VisibleForTesting + static DescribeTabletServersResponse describeTabletServers(CoordinatorContext ctx) { + // sorted by server id for a deterministic response order + Map loads = new TreeMap<>(); + // report live and shutting-down servers even if they host no replicas, so that + // an evacuated server explicitly shows zero replicas + for (int serverId : ctx.liveOrShuttingDownTabletServers()) { + getOrCreateLoad(loads, serverId); + } + + for (TableBucket tb : ctx.getAllBuckets()) { + countBucketReplicas(ctx, tb, loads); + } + + DescribeTabletServersResponse response = new DescribeTabletServersResponse(); + response.addAllTabletServers(loads.values()); + return response; + } + + private static void countBucketReplicas( + CoordinatorContext ctx, TableBucket tb, Map loads) { + for (int serverId : ctx.getAssignment(tb)) { + PbTabletServerLoad load = getOrCreateLoad(loads, serverId); + load.setNumReplicas(load.getNumReplicas() + 1); + } + + Optional laiOpt = ctx.getBucketLeaderAndIsr(tb); + if (!laiOpt.isPresent()) { + return; + } + + for (int serverId : laiOpt.get().isr()) { + PbTabletServerLoad load = getOrCreateLoad(loads, serverId); + load.setInSyncReplicas(load.getInSyncReplicas() + 1); + } + + int leader = laiOpt.get().leader(); + if (leader == LeaderAndIsr.NO_LEADER) { + return; + } + + PbTabletServerLoad load = getOrCreateLoad(loads, leader); + load.setNumLeaderReplicas(load.getNumLeaderReplicas() + 1); + if (ctx.isLeaderActive(tb)) { + load.setActiveLeaderReplicas(load.getActiveLeaderReplicas() + 1); + } + } + + private static PbTabletServerLoad getOrCreateLoad( + Map loads, int serverId) { + return loads.computeIfAbsent( + serverId, + id -> + new PbTabletServerLoad() + .setServerId(id) + .setNumReplicas(0) + .setInSyncReplicas(0) + .setNumLeaderReplicas(0) + .setActiveLeaderReplicas(0)); + } + @VisibleForTesting public DataLakeFormat getDataLakeFormat() { return lakeCatalogDynamicLoader.getLakeCatalogContainer().getDataLakeFormat(); diff --git a/fluss-server/src/main/java/org/apache/fluss/server/tablet/TabletService.java b/fluss-server/src/main/java/org/apache/fluss/server/tablet/TabletService.java index a18b6289d5e..949935d3a10 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/tablet/TabletService.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/tablet/TabletService.java @@ -37,6 +37,8 @@ import org.apache.fluss.rpc.entity.ResultForBucket; import org.apache.fluss.rpc.gateway.CoordinatorGateway; import org.apache.fluss.rpc.gateway.TabletServerGateway; +import org.apache.fluss.rpc.messages.DescribeTabletServersRequest; +import org.apache.fluss.rpc.messages.DescribeTabletServersResponse; import org.apache.fluss.rpc.messages.FetchLogRequest; import org.apache.fluss.rpc.messages.FetchLogResponse; import org.apache.fluss.rpc.messages.GetClusterHealthRequest; @@ -513,6 +515,29 @@ public CompletableFuture getClusterHealth( return coordinatorGateway.getClusterHealth(request); } + @Override + public CompletableFuture describeTabletServers( + DescribeTabletServersRequest request) { + // Tablet servers don't own the cluster-wide replica view; we forward the call to the + // coordinator over the internal listener. + if (authorizer != null) { + authorizer.authorize(currentSession(), OperationType.DESCRIBE, Resource.cluster()); + } + + if (metadataCache.getCoordinatorServer(interListenerName) == null) { + // Fail fast during the startup window before the tablet has received its first + // UpdateMetadataRequest from the coordinator. The supplier inside coordinatorGateway + // would otherwise block-and-retry, which is not what callers want. + CompletableFuture failed = new CompletableFuture<>(); + failed.completeExceptionally( + new StaleMetadataException( + "Tablet server has not yet received coordinator metadata; tablet" + + " server loads are unavailable.")); + return failed; + } + return coordinatorGateway.describeTabletServers(request); + } + @Override public CompletableFuture scanKv(ScanKvRequest request) { ScanKvResponse response = new ScanKvResponse(); diff --git a/fluss-server/src/test/java/org/apache/fluss/server/coordinator/DescribeTabletServersTest.java b/fluss-server/src/test/java/org/apache/fluss/server/coordinator/DescribeTabletServersTest.java new file mode 100644 index 00000000000..a83cb7ebc1f --- /dev/null +++ b/fluss-server/src/test/java/org/apache/fluss/server/coordinator/DescribeTabletServersTest.java @@ -0,0 +1,239 @@ +/* + * 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; + +import org.apache.fluss.cluster.Endpoint; +import org.apache.fluss.cluster.ServerType; +import org.apache.fluss.metadata.TableBucket; +import org.apache.fluss.rpc.messages.DescribeTabletServersResponse; +import org.apache.fluss.rpc.messages.GetClusterHealthResponse; +import org.apache.fluss.rpc.messages.PbTabletServerLoad; +import org.apache.fluss.server.metadata.ServerInfo; +import org.apache.fluss.server.zk.ZkEpoch; +import org.apache.fluss.server.zk.data.LeaderAndIsr; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.util.Arrays; +import java.util.Collections; + +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for {@link CoordinatorService#describeTabletServers(CoordinatorContext)}. */ +class DescribeTabletServersTest { + + private CoordinatorContext ctx; + + @BeforeEach + void setUp() { + ctx = new CoordinatorContext(ZkEpoch.INITIAL_EPOCH); + ctx.setLiveTabletServers( + Arrays.asList(makeServerInfo(0), makeServerInfo(1), makeServerInfo(2))); + } + + @Test + void testEmptyClusterReportsLiveServersWithZeroCounters() { + DescribeTabletServersResponse resp = CoordinatorService.describeTabletServers(ctx); + + assertThat(resp.getTabletServersList()) + .extracting(PbTabletServerLoad::getServerId) + .containsExactly(0, 1, 2); + assertLoad(resp, 0, 0, 0, 0, 0); + assertLoad(resp, 1, 0, 0, 0, 0); + assertLoad(resp, 2, 0, 0, 0, 0); + } + + @Test + void testAllInSyncAllLeadersActive() { + TableBucket tb = new TableBucket(1L, 0); + ctx.updateBucketReplicaAssignment(tb, Arrays.asList(0, 1, 2)); + ctx.putBucketLeaderAndIsr( + tb, new LeaderAndIsr(0, 1, Arrays.asList(0, 1, 2), Collections.emptyList(), 0, 1)); + + DescribeTabletServersResponse resp = CoordinatorService.describeTabletServers(ctx); + + assertLoad(resp, 0, 1, 1, 1, 1); + assertLoad(resp, 1, 1, 1, 0, 0); + assertLoad(resp, 2, 1, 1, 0, 0); + } + + @Test + void testReplicaOutOfIsrOnlyAffectsThatServer() { + TableBucket tb = new TableBucket(1L, 0); + ctx.updateBucketReplicaAssignment(tb, Arrays.asList(0, 1, 2)); + ctx.putBucketLeaderAndIsr( + tb, new LeaderAndIsr(0, 1, Arrays.asList(0, 1), Collections.emptyList(), 0, 1)); + + DescribeTabletServersResponse resp = CoordinatorService.describeTabletServers(ctx); + + assertLoad(resp, 0, 1, 1, 1, 1); + assertLoad(resp, 1, 1, 1, 0, 0); + assertLoad(resp, 2, 1, 0, 0, 0); + } + + @Test + void testInactiveLeader() { + TableBucket tb = new TableBucket(1L, 0); + ctx.updateBucketReplicaAssignment(tb, Arrays.asList(0, 1)); + ctx.putBucketLeaderAndIsr( + tb, new LeaderAndIsr(0, 1, Arrays.asList(0, 1), Collections.emptyList(), 0, 1)); + ctx.addPendingLeaderActivation(tb); + + assertLoad(CoordinatorService.describeTabletServers(ctx), 0, 1, 1, 1, 0); + + ctx.clearPendingLeaderActivation(tb); + + assertLoad(CoordinatorService.describeTabletServers(ctx), 0, 1, 1, 1, 1); + } + + @Test + void testNoLeaderBucketAttributesLeadershipToNoServer() { + TableBucket tb = new TableBucket(1L, 0); + ctx.updateBucketReplicaAssignment(tb, Arrays.asList(0, 1)); + ctx.putBucketLeaderAndIsr( + tb, + new LeaderAndIsr( + LeaderAndIsr.NO_LEADER, + 1, + Arrays.asList(0, 1), + Collections.emptyList(), + 0, + 1)); + + DescribeTabletServersResponse resp = CoordinatorService.describeTabletServers(ctx); + + assertLoad(resp, 0, 1, 1, 0, 0); + assertLoad(resp, 1, 1, 1, 0, 0); + } + + @Test + void testEvacuatedServerReportsZeroReplicas() { + // only servers 0 and 1 host replicas; server 2 is evacuated and safe to remove + TableBucket tb = new TableBucket(1L, 0); + ctx.updateBucketReplicaAssignment(tb, Arrays.asList(0, 1)); + ctx.putBucketLeaderAndIsr( + tb, new LeaderAndIsr(0, 1, Arrays.asList(0, 1), Collections.emptyList(), 0, 1)); + + assertLoad(CoordinatorService.describeTabletServers(ctx), 2, 0, 0, 0, 0); + } + + @Test + void testDeadServerStillAssignedIsReported() { + // server 5 is not live but still has an assigned replica; it must be visible so + // that a scale-in gate does not treat the cluster as fully drained + TableBucket tb = new TableBucket(1L, 0); + ctx.updateBucketReplicaAssignment(tb, Arrays.asList(0, 1, 5)); + ctx.putBucketLeaderAndIsr( + tb, new LeaderAndIsr(0, 1, Arrays.asList(0, 1), Collections.emptyList(), 0, 1)); + + assertLoad(CoordinatorService.describeTabletServers(ctx), 5, 1, 0, 0, 0); + } + + @Test + void testPartitionedTableBucket() { + TableBucket tb = new TableBucket(1L, 100L, 0); + ctx.updateBucketReplicaAssignment(tb, Arrays.asList(0, 1)); + ctx.putBucketLeaderAndIsr( + tb, + new LeaderAndIsr( + 0, 1, Collections.singletonList(0), Collections.emptyList(), 0, 1)); + + DescribeTabletServersResponse resp = CoordinatorService.describeTabletServers(ctx); + + assertLoad(resp, 0, 1, 1, 1, 1); + assertLoad(resp, 1, 1, 0, 0, 0); + } + + @Test + void testPerServerSumsReconcileWithClusterHealth() { + TableBucket tb1 = new TableBucket(1L, 0); + TableBucket tb2 = new TableBucket(1L, 1); + TableBucket tb3 = new TableBucket(2L, 0); + + ctx.updateBucketReplicaAssignment(tb1, Arrays.asList(0, 1)); + ctx.updateBucketReplicaAssignment(tb2, Arrays.asList(1, 2)); + ctx.updateBucketReplicaAssignment(tb3, Arrays.asList(0, 2)); + + ctx.putBucketLeaderAndIsr( + tb1, new LeaderAndIsr(0, 1, Arrays.asList(0, 1), Collections.emptyList(), 0, 1)); + ctx.putBucketLeaderAndIsr( + tb2, + new LeaderAndIsr( + 1, 1, Collections.singletonList(1), Collections.emptyList(), 0, 1)); + ctx.putBucketLeaderAndIsr( + tb3, new LeaderAndIsr(0, 1, Arrays.asList(0, 2), Collections.emptyList(), 0, 1)); + ctx.addPendingLeaderActivation(tb3); + + DescribeTabletServersResponse resp = CoordinatorService.describeTabletServers(ctx); + GetClusterHealthResponse health = CoordinatorService.computeClusterHealth(ctx); + + int numReplicas = 0; + int inSyncReplicas = 0; + int numLeaderReplicas = 0; + int activeLeaderReplicas = 0; + for (PbTabletServerLoad load : resp.getTabletServersList()) { + numReplicas += load.getNumReplicas(); + inSyncReplicas += load.getInSyncReplicas(); + numLeaderReplicas += load.getNumLeaderReplicas(); + activeLeaderReplicas += load.getActiveLeaderReplicas(); + } + + assertThat(numReplicas).isEqualTo(health.getNumReplicas()); + assertThat(inSyncReplicas).isEqualTo(health.getInSyncReplicas()); + assertThat(numLeaderReplicas).isEqualTo(health.getNumLeaderReplicas()); + assertThat(activeLeaderReplicas).isEqualTo(health.getActiveLeaderReplicas()); + } + + private static void assertLoad( + DescribeTabletServersResponse resp, + int serverId, + int numReplicas, + int inSyncReplicas, + int numLeaderReplicas, + int activeLeaderReplicas) { + PbTabletServerLoad load = + resp.getTabletServersList().stream() + .filter(l -> l.getServerId() == serverId) + .findFirst() + .orElseThrow( + () -> + new AssertionError( + "no load reported for tablet server " + serverId)); + assertThat(load.getNumReplicas()) + .as("numReplicas of server %s", serverId) + .isEqualTo(numReplicas); + assertThat(load.getInSyncReplicas()) + .as("inSyncReplicas of server %s", serverId) + .isEqualTo(inSyncReplicas); + assertThat(load.getNumLeaderReplicas()) + .as("numLeaderReplicas of server %s", serverId) + .isEqualTo(numLeaderReplicas); + assertThat(load.getActiveLeaderReplicas()) + .as("activeLeaderReplicas of server %s", serverId) + .isEqualTo(activeLeaderReplicas); + } + + private static ServerInfo makeServerInfo(int id) { + return new ServerInfo( + id, + "RACK" + id, + Endpoint.fromListenersString("CLIENT://host" + id + ":9124"), + ServerType.TABLET_SERVER); + } +} diff --git a/fluss-server/src/test/java/org/apache/fluss/server/coordinator/TestCoordinatorGateway.java b/fluss-server/src/test/java/org/apache/fluss/server/coordinator/TestCoordinatorGateway.java index 7f3bc32e8c4..f9a3eea370d 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/coordinator/TestCoordinatorGateway.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/coordinator/TestCoordinatorGateway.java @@ -60,6 +60,8 @@ import org.apache.fluss.rpc.messages.DeleteProducerOffsetsResponse; import org.apache.fluss.rpc.messages.DescribeClusterConfigsRequest; import org.apache.fluss.rpc.messages.DescribeClusterConfigsResponse; +import org.apache.fluss.rpc.messages.DescribeTabletServersRequest; +import org.apache.fluss.rpc.messages.DescribeTabletServersResponse; import org.apache.fluss.rpc.messages.DropAclsRequest; import org.apache.fluss.rpc.messages.DropAclsResponse; import org.apache.fluss.rpc.messages.DropDatabaseRequest; @@ -485,6 +487,12 @@ public CompletableFuture getClusterHealth( throw new UnsupportedOperationException(); } + @Override + public CompletableFuture describeTabletServers( + DescribeTabletServersRequest request) { + throw new UnsupportedOperationException(); + } + @Override public CompletableFuture registerProducerOffsets( RegisterProducerOffsetsRequest request) { diff --git a/fluss-server/src/test/java/org/apache/fluss/server/tablet/TestTabletServerGateway.java b/fluss-server/src/test/java/org/apache/fluss/server/tablet/TestTabletServerGateway.java index c78270ea5ca..571936ce151 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/tablet/TestTabletServerGateway.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/tablet/TestTabletServerGateway.java @@ -28,6 +28,8 @@ import org.apache.fluss.rpc.messages.DatabaseExistsResponse; import org.apache.fluss.rpc.messages.DescribeClusterConfigsRequest; import org.apache.fluss.rpc.messages.DescribeClusterConfigsResponse; +import org.apache.fluss.rpc.messages.DescribeTabletServersRequest; +import org.apache.fluss.rpc.messages.DescribeTabletServersResponse; import org.apache.fluss.rpc.messages.FetchLogRequest; import org.apache.fluss.rpc.messages.FetchLogResponse; import org.apache.fluss.rpc.messages.GetClusterHealthRequest; @@ -242,6 +244,12 @@ public CompletableFuture getClusterHealth( throw new UnsupportedOperationException(); } + @Override + public CompletableFuture describeTabletServers( + DescribeTabletServersRequest request) { + throw new UnsupportedOperationException(); + } + @Override public CompletableFuture databaseExists(DatabaseExistsRequest request) { throw new UnsupportedOperationException(); From eb91bc6856292357db9bd1ece1deb559bc5f11a3 Mon Sep 17 00:00:00 2001 From: morazow Date: Fri, 24 Jul 2026 14:16:54 +0200 Subject: [PATCH 2/3] minor refactors --- .../org/apache/fluss/client/admin/Admin.java | 4 ++- .../client/admin/TabletServerDescription.java | 8 ++++-- .../fluss/client/admin/FlussAdminITCase.java | 11 +++++++- .../coordinator/CoordinatorService.java | 4 +-- .../DescribeTabletServersTest.java | 26 +++++++++++-------- 5 files changed, 36 insertions(+), 17 deletions(-) diff --git a/fluss-client/src/main/java/org/apache/fluss/client/admin/Admin.java b/fluss-client/src/main/java/org/apache/fluss/client/admin/Admin.java index dcc368bfefb..65e8fcc13f2 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/admin/Admin.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/admin/Admin.java @@ -813,7 +813,9 @@ CompletableFuture registerProducerOffsets( * unavailability; *

  • rolling-upgrade gate - only proceed to the next server when the previous one is green * again ({@code inSyncReplicas == numReplicas && activeLeaderReplicas == - * numLeaderReplicas}); + * numLeaderReplicas}) and {@link #getClusterHealth()} is not {@link + * ClusterHealthStatus#RED}, since a leaderless bucket may not surface in the per-server + * counters (see {@link TabletServerDescription}); *
  • reporting per-server tablet load in cluster status. * * diff --git a/fluss-client/src/main/java/org/apache/fluss/client/admin/TabletServerDescription.java b/fluss-client/src/main/java/org/apache/fluss/client/admin/TabletServerDescription.java index a4605e38302..8f3a0a93663 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/admin/TabletServerDescription.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/admin/TabletServerDescription.java @@ -39,8 +39,12 @@ * *

    Note: a bucket without an elected leader is attributed to no server's {@code * numLeaderReplicas}, so the per-server values sum to the cluster-wide {@link - * ClusterHealth#getNumLeaderReplicas()} only when every bucket has an elected leader. Such buckets - * still surface through {@code inSyncReplicas < numReplicas} on the assigned servers. + * ClusterHealth#getNumLeaderReplicas()} only when every bucket has an elected leader. Such a bucket + * does not always surface through {@code inSyncReplicas < numReplicas}: when the last ISR member + * goes offline, it is kept in the ISR so that a leader can be re-elected later, and the assigned + * servers still look green by the predicate above. A rolling-upgrade gate should therefore + * additionally require {@link Admin#getClusterHealth()} to not report {@link + * ClusterHealthStatus#RED}, which counts leaderless buckets unconditionally. * * @since 1.0 */ diff --git a/fluss-client/src/test/java/org/apache/fluss/client/admin/FlussAdminITCase.java b/fluss-client/src/test/java/org/apache/fluss/client/admin/FlussAdminITCase.java index f784706c37f..c79e3feacdb 100644 --- a/fluss-client/src/test/java/org/apache/fluss/client/admin/FlussAdminITCase.java +++ b/fluss-client/src/test/java/org/apache/fluss/client/admin/FlussAdminITCase.java @@ -2870,7 +2870,16 @@ void testDescribeTabletServersDuringRollingUpgrade() throws Exception { waitAllReplicasReady(tableId, 3); // Phase 1: Cluster is healthy - every live server is reported, hosts replicas of the - // created table (replication factor 3 on 3 servers) and is green. + // created table (replication factor 3 on 3 servers) and is green. The cluster is shared + // with other tests, so wait until residue from them (e.g. a recovering ISR) has settled + // before taking the snapshot asserted below. + waitUntil( + () -> + admin.describeTabletServers().get().stream() + .allMatch(FlussAdminITCase::isServerGreen), + Duration.ofMinutes(1), + "All tablet servers should be green before the rolling upgrade starts"); + List servers = admin.describeTabletServers().get(); assertThat(servers).extracting(TabletServerDescription::getServerId).contains(0, 1, 2); for (TabletServerDescription server : servers) { diff --git a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorService.java b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorService.java index 5b99dc3b4ee..15d3e8d4c36 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorService.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorService.java @@ -1508,13 +1508,13 @@ public CompletableFuture describeTabletServers( } AccessContextEvent event = - new AccessContextEvent<>(CoordinatorService::describeTabletServers); + new AccessContextEvent<>(CoordinatorService::computeTabletServerLoads); eventManagerSupplier.get().put(event); return event.getResultFuture(); } @VisibleForTesting - static DescribeTabletServersResponse describeTabletServers(CoordinatorContext ctx) { + static DescribeTabletServersResponse computeTabletServerLoads(CoordinatorContext ctx) { // sorted by server id for a deterministic response order Map loads = new TreeMap<>(); // report live and shutting-down servers even if they host no replicas, so that diff --git a/fluss-server/src/test/java/org/apache/fluss/server/coordinator/DescribeTabletServersTest.java b/fluss-server/src/test/java/org/apache/fluss/server/coordinator/DescribeTabletServersTest.java index a83cb7ebc1f..cdc1ae0cb25 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/coordinator/DescribeTabletServersTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/coordinator/DescribeTabletServersTest.java @@ -35,7 +35,7 @@ import static org.assertj.core.api.Assertions.assertThat; -/** Tests for {@link CoordinatorService#describeTabletServers(CoordinatorContext)}. */ +/** Tests for {@link CoordinatorService#computeTabletServerLoads(CoordinatorContext)}. */ class DescribeTabletServersTest { private CoordinatorContext ctx; @@ -49,7 +49,7 @@ void setUp() { @Test void testEmptyClusterReportsLiveServersWithZeroCounters() { - DescribeTabletServersResponse resp = CoordinatorService.describeTabletServers(ctx); + DescribeTabletServersResponse resp = CoordinatorService.computeTabletServerLoads(ctx); assertThat(resp.getTabletServersList()) .extracting(PbTabletServerLoad::getServerId) @@ -66,7 +66,7 @@ void testAllInSyncAllLeadersActive() { ctx.putBucketLeaderAndIsr( tb, new LeaderAndIsr(0, 1, Arrays.asList(0, 1, 2), Collections.emptyList(), 0, 1)); - DescribeTabletServersResponse resp = CoordinatorService.describeTabletServers(ctx); + DescribeTabletServersResponse resp = CoordinatorService.computeTabletServerLoads(ctx); assertLoad(resp, 0, 1, 1, 1, 1); assertLoad(resp, 1, 1, 1, 0, 0); @@ -80,7 +80,7 @@ void testReplicaOutOfIsrOnlyAffectsThatServer() { ctx.putBucketLeaderAndIsr( tb, new LeaderAndIsr(0, 1, Arrays.asList(0, 1), Collections.emptyList(), 0, 1)); - DescribeTabletServersResponse resp = CoordinatorService.describeTabletServers(ctx); + DescribeTabletServersResponse resp = CoordinatorService.computeTabletServerLoads(ctx); assertLoad(resp, 0, 1, 1, 1, 1); assertLoad(resp, 1, 1, 1, 0, 0); @@ -95,15 +95,19 @@ void testInactiveLeader() { tb, new LeaderAndIsr(0, 1, Arrays.asList(0, 1), Collections.emptyList(), 0, 1)); ctx.addPendingLeaderActivation(tb); - assertLoad(CoordinatorService.describeTabletServers(ctx), 0, 1, 1, 1, 0); + assertLoad(CoordinatorService.computeTabletServerLoads(ctx), 0, 1, 1, 1, 0); ctx.clearPendingLeaderActivation(tb); - assertLoad(CoordinatorService.describeTabletServers(ctx), 0, 1, 1, 1, 1); + assertLoad(CoordinatorService.computeTabletServerLoads(ctx), 0, 1, 1, 1, 1); } @Test void testNoLeaderBucketAttributesLeadershipToNoServer() { + // this is a real state: when the last ISR member goes offline, ReplicaStateMachine keeps + // it in the ISR (so a leader can be re-elected later) and sets the leader to NO_LEADER. + // Both servers then still report inSync == num, so the per-server green predicate alone + // cannot detect the offline bucket - see the TabletServerDescription javadoc. TableBucket tb = new TableBucket(1L, 0); ctx.updateBucketReplicaAssignment(tb, Arrays.asList(0, 1)); ctx.putBucketLeaderAndIsr( @@ -116,7 +120,7 @@ void testNoLeaderBucketAttributesLeadershipToNoServer() { 0, 1)); - DescribeTabletServersResponse resp = CoordinatorService.describeTabletServers(ctx); + DescribeTabletServersResponse resp = CoordinatorService.computeTabletServerLoads(ctx); assertLoad(resp, 0, 1, 1, 0, 0); assertLoad(resp, 1, 1, 1, 0, 0); @@ -130,7 +134,7 @@ void testEvacuatedServerReportsZeroReplicas() { ctx.putBucketLeaderAndIsr( tb, new LeaderAndIsr(0, 1, Arrays.asList(0, 1), Collections.emptyList(), 0, 1)); - assertLoad(CoordinatorService.describeTabletServers(ctx), 2, 0, 0, 0, 0); + assertLoad(CoordinatorService.computeTabletServerLoads(ctx), 2, 0, 0, 0, 0); } @Test @@ -142,7 +146,7 @@ void testDeadServerStillAssignedIsReported() { ctx.putBucketLeaderAndIsr( tb, new LeaderAndIsr(0, 1, Arrays.asList(0, 1), Collections.emptyList(), 0, 1)); - assertLoad(CoordinatorService.describeTabletServers(ctx), 5, 1, 0, 0, 0); + assertLoad(CoordinatorService.computeTabletServerLoads(ctx), 5, 1, 0, 0, 0); } @Test @@ -154,7 +158,7 @@ void testPartitionedTableBucket() { new LeaderAndIsr( 0, 1, Collections.singletonList(0), Collections.emptyList(), 0, 1)); - DescribeTabletServersResponse resp = CoordinatorService.describeTabletServers(ctx); + DescribeTabletServersResponse resp = CoordinatorService.computeTabletServerLoads(ctx); assertLoad(resp, 0, 1, 1, 1, 1); assertLoad(resp, 1, 1, 0, 0, 0); @@ -180,7 +184,7 @@ void testPerServerSumsReconcileWithClusterHealth() { tb3, new LeaderAndIsr(0, 1, Arrays.asList(0, 2), Collections.emptyList(), 0, 1)); ctx.addPendingLeaderActivation(tb3); - DescribeTabletServersResponse resp = CoordinatorService.describeTabletServers(ctx); + DescribeTabletServersResponse resp = CoordinatorService.computeTabletServerLoads(ctx); GetClusterHealthResponse health = CoordinatorService.computeClusterHealth(ctx); int numReplicas = 0; From 482d2baec26016a5d269b802d4fe7d5065a081fc Mon Sep 17 00:00:00 2001 From: morazow Date: Fri, 31 Jul 2026 23:18:25 +0200 Subject: [PATCH 3/3] [server][rpc] Report server tags in DescribeTabletServers response Expose the ServerTag (PERMANENT_OFFLINE / TEMPORARY_OFFLINE) of each tablet server in the DescribeTabletServers response, since there was no read API for tags set via addServerTag. A tagged server is reported even when it is dead and hosts no replicas, so its tag stays visible. --- .../client/admin/TabletServerDescription.java | 30 +++++++++++++++-- .../client/utils/ClientRpcMessageUtils.java | 6 +++- .../fluss/client/admin/FlussAdminITCase.java | 8 +++++ fluss-rpc/src/main/proto/FlussApi.proto | 1 + .../coordinator/CoordinatorService.java | 6 ++++ .../DescribeTabletServersTest.java | 32 ++++++++++++++----- 6 files changed, 71 insertions(+), 12 deletions(-) diff --git a/fluss-client/src/main/java/org/apache/fluss/client/admin/TabletServerDescription.java b/fluss-client/src/main/java/org/apache/fluss/client/admin/TabletServerDescription.java index 8f3a0a93663..d9847b8b3ec 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/admin/TabletServerDescription.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/admin/TabletServerDescription.java @@ -18,8 +18,12 @@ package org.apache.fluss.client.admin; import org.apache.fluss.annotation.PublicEvolving; +import org.apache.fluss.cluster.rebalance.ServerTag; + +import javax.annotation.Nullable; import java.util.Objects; +import java.util.Optional; /** * Per tablet server replica load returned by {@link Admin#describeTabletServers()}. @@ -56,18 +60,21 @@ public final class TabletServerDescription { private final int inSyncReplicas; private final int numLeaderReplicas; private final int activeLeaderReplicas; + private final @Nullable ServerTag serverTag; public TabletServerDescription( int serverId, int numReplicas, int inSyncReplicas, int numLeaderReplicas, - int activeLeaderReplicas) { + int activeLeaderReplicas, + @Nullable ServerTag serverTag) { this.serverId = serverId; this.numReplicas = numReplicas; this.inSyncReplicas = inSyncReplicas; this.numLeaderReplicas = numLeaderReplicas; this.activeLeaderReplicas = activeLeaderReplicas; + this.serverTag = serverTag; } public int getServerId() { @@ -94,6 +101,15 @@ public int getActiveLeaderReplicas() { return activeLeaderReplicas; } + /** + * The {@link ServerTag} set on this server via {@link Admin#addServerTag}, or empty if the + * server is untagged. A tagged server is being drained, so it is reported even when it is dead + * and hosts no replicas. + */ + public Optional getServerTag() { + return Optional.ofNullable(serverTag); + } + @Override public boolean equals(Object o) { if (this == o) { @@ -107,13 +123,19 @@ public boolean equals(Object o) { && numReplicas == that.numReplicas && inSyncReplicas == that.inSyncReplicas && numLeaderReplicas == that.numLeaderReplicas - && activeLeaderReplicas == that.activeLeaderReplicas; + && activeLeaderReplicas == that.activeLeaderReplicas + && serverTag == that.serverTag; } @Override public int hashCode() { return Objects.hash( - serverId, numReplicas, inSyncReplicas, numLeaderReplicas, activeLeaderReplicas); + serverId, + numReplicas, + inSyncReplicas, + numLeaderReplicas, + activeLeaderReplicas, + serverTag); } @Override @@ -129,6 +151,8 @@ public String toString() { + numLeaderReplicas + ", activeLeaderReplicas=" + activeLeaderReplicas + + ", serverTag=" + + serverTag + '}'; } } diff --git a/fluss-client/src/main/java/org/apache/fluss/client/utils/ClientRpcMessageUtils.java b/fluss-client/src/main/java/org/apache/fluss/client/utils/ClientRpcMessageUtils.java index 24ac21aca8f..7baba104f8f 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/utils/ClientRpcMessageUtils.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/utils/ClientRpcMessageUtils.java @@ -36,6 +36,7 @@ import org.apache.fluss.cluster.rebalance.RebalanceProgress; import org.apache.fluss.cluster.rebalance.RebalanceResultForBucket; import org.apache.fluss.cluster.rebalance.RebalanceStatus; +import org.apache.fluss.cluster.rebalance.ServerTag; import org.apache.fluss.config.cluster.AlterConfigOpType; import org.apache.fluss.config.cluster.ColumnPositionType; import org.apache.fluss.config.cluster.ConfigEntry; @@ -911,7 +912,10 @@ public static List toTabletServerDescriptions( load.getNumReplicas(), load.getInSyncReplicas(), load.getNumLeaderReplicas(), - load.getActiveLeaderReplicas())) + load.getActiveLeaderReplicas(), + load.hasServerTag() + ? ServerTag.valueOf(load.getServerTag()) + : null)) .collect(Collectors.toList()); } } diff --git a/fluss-client/src/test/java/org/apache/fluss/client/admin/FlussAdminITCase.java b/fluss-client/src/test/java/org/apache/fluss/client/admin/FlussAdminITCase.java index c79e3feacdb..9d0594a3308 100644 --- a/fluss-client/src/test/java/org/apache/fluss/client/admin/FlussAdminITCase.java +++ b/fluss-client/src/test/java/org/apache/fluss/client/admin/FlussAdminITCase.java @@ -2305,6 +2305,14 @@ public void testAddAndRemoveServerTags() throws Exception { .containsEntry(0, ServerTag.PERMANENT_OFFLINE) .containsEntry(1, ServerTag.PERMANENT_OFFLINE); + // the tags are also visible via describeTabletServers. + List described = admin.describeTabletServers().get(); + assertThat(getTabletServerDescription(described, 0).getServerTag()) + .contains(ServerTag.PERMANENT_OFFLINE); + assertThat(getTabletServerDescription(described, 1).getServerTag()) + .contains(ServerTag.PERMANENT_OFFLINE); + assertThat(getTabletServerDescription(described, 2).getServerTag()).isNotPresent(); + // 3.add different server tag for server 0,2. error will be thrown and tag for 2 will not be // added. assertThatThrownBy( diff --git a/fluss-rpc/src/main/proto/FlussApi.proto b/fluss-rpc/src/main/proto/FlussApi.proto index 258adc65246..3dc5094dac9 100644 --- a/fluss-rpc/src/main/proto/FlussApi.proto +++ b/fluss-rpc/src/main/proto/FlussApi.proto @@ -822,6 +822,7 @@ message PbTabletServerLoad { required int32 in_sync_replicas = 3; required int32 num_leader_replicas = 4; required int32 active_leader_replicas = 5; + optional int32 server_tag = 6; // ServerTag: PERMANENT_OFFLINE=0, TEMPORARY_OFFLINE=1 } message PbApiVersion { diff --git a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorService.java b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorService.java index 15d3e8d4c36..cb302e6ac0b 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorService.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorService.java @@ -1527,6 +1527,12 @@ static DescribeTabletServersResponse computeTabletServerLoads(CoordinatorContext countBucketReplicas(ctx, tb, loads); } + // tagged servers are being drained, so report them even when dead and empty + ctx.getServerTags() + .forEach( + (serverId, tag) -> + getOrCreateLoad(loads, serverId).setServerTag(tag.value)); + DescribeTabletServersResponse response = new DescribeTabletServersResponse(); response.addAllTabletServers(loads.values()); return response; diff --git a/fluss-server/src/test/java/org/apache/fluss/server/coordinator/DescribeTabletServersTest.java b/fluss-server/src/test/java/org/apache/fluss/server/coordinator/DescribeTabletServersTest.java index cdc1ae0cb25..ae62f02c42c 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/coordinator/DescribeTabletServersTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/coordinator/DescribeTabletServersTest.java @@ -19,6 +19,7 @@ import org.apache.fluss.cluster.Endpoint; import org.apache.fluss.cluster.ServerType; +import org.apache.fluss.cluster.rebalance.ServerTag; import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.rpc.messages.DescribeTabletServersResponse; import org.apache.fluss.rpc.messages.GetClusterHealthResponse; @@ -164,6 +165,20 @@ void testPartitionedTableBucket() { assertLoad(resp, 1, 1, 0, 0, 0); } + @Test + void testServerTagsAreReported() { + ctx.putServerTag(1, ServerTag.TEMPORARY_OFFLINE); + // server 5 is dead and hosts no replicas, but its tag must remain visible + ctx.putServerTag(5, ServerTag.PERMANENT_OFFLINE); + + DescribeTabletServersResponse resp = CoordinatorService.computeTabletServerLoads(ctx); + + assertThat(findLoad(resp, 0).hasServerTag()).isFalse(); + assertThat(findLoad(resp, 1).getServerTag()).isEqualTo(ServerTag.TEMPORARY_OFFLINE.value); + assertThat(findLoad(resp, 5).getServerTag()).isEqualTo(ServerTag.PERMANENT_OFFLINE.value); + assertLoad(resp, 5, 0, 0, 0, 0); + } + @Test void testPerServerSumsReconcileWithClusterHealth() { TableBucket tb1 = new TableBucket(1L, 0); @@ -211,14 +226,7 @@ private static void assertLoad( int inSyncReplicas, int numLeaderReplicas, int activeLeaderReplicas) { - PbTabletServerLoad load = - resp.getTabletServersList().stream() - .filter(l -> l.getServerId() == serverId) - .findFirst() - .orElseThrow( - () -> - new AssertionError( - "no load reported for tablet server " + serverId)); + PbTabletServerLoad load = findLoad(resp, serverId); assertThat(load.getNumReplicas()) .as("numReplicas of server %s", serverId) .isEqualTo(numReplicas); @@ -233,6 +241,14 @@ private static void assertLoad( .isEqualTo(activeLeaderReplicas); } + private static PbTabletServerLoad findLoad(DescribeTabletServersResponse resp, int serverId) { + return resp.getTabletServersList().stream() + .filter(l -> l.getServerId() == serverId) + .findFirst() + .orElseThrow( + () -> new AssertionError("no load reported for tablet server " + serverId)); + } + private static ServerInfo makeServerInfo(int id) { return new ServerInfo( id,