From 0b24b9d7ddcf7c9631d1ec4364be48b7af7560e5 Mon Sep 17 00:00:00 2001 From: Allen Cao <258265534+heyallencao@users.noreply.github.com> Date: Wed, 16 Sep 2026 18:20:23 +0800 Subject: [PATCH 1/4] [coordinator] Add rebalance progress and outcome metrics Expose leader-scoped rebalance gauges and round counters using existing task state. Avoid inferring historical bucket outcomes during recovery, and add regression tests and metric documentation. Refs apache/fluss#4055 --- .../org/apache/fluss/metrics/MetricNames.java | 10 + .../CoordinatorEventProcessor.java | 6 +- .../rebalance/RebalanceManager.java | 77 ++++- .../metrics/group/CoordinatorMetricGroup.java | 38 +++ .../rebalance/RebalanceManagerTest.java | 322 +++++++++++++++++- .../observability/monitor-metrics.md | 26 ++ 6 files changed, 469 insertions(+), 10 deletions(-) diff --git a/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java b/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java index ac60ee5f101..3649f7e23c0 100644 --- a/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java +++ b/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java @@ -47,6 +47,16 @@ public class MetricNames { public static final String KV_LEADER_REPLICA_CAPACITY = "kvLeaderReplicaCapacity"; public static final String REPLICAS_TO_DELETE_COUNT = "replicasToDeleteCount"; public static final String PENDING_LEADER_ACTIVATION_COUNT = "pendingLeaderActivationCount"; + public static final String REBALANCE_IN_PROGRESS = "rebalanceInProgress"; + public static final String REBALANCE_BUCKETS_PENDING = "rebalanceBucketsPending"; + public static final String REBALANCE_BUCKETS_COMPLETED = "rebalanceBucketsCompleted"; + public static final String REBALANCE_BUCKETS_FAILED = "rebalanceBucketsFailed"; + public static final String REBALANCE_BUCKETS_TIMED_OUT = "rebalanceBucketsTimedOut"; + public static final String REBALANCE_DURATION_MS = "rebalanceDurationMs"; + public static final String INFLIGHT_BUCKET_DURATION_MS = "inflightBucketDurationMs"; + public static final String REBALANCES_COMPLETED_TOTAL = "rebalancesCompletedTotal"; + public static final String REBALANCES_FAILED_TOTAL = "rebalancesFailedTotal"; + public static final String REBALANCES_CANCELED_TOTAL = "rebalancesCanceledTotal"; // for coordinator sender (per-tablet-server control request sender threads) public static final String SENDER_QUEUE_SIZE = "senderQueueSize"; public static final String SENDER_QUEUE_TIME_MS = "senderQueueTimeMs"; diff --git a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessor.java b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessor.java index aa04f8bec9e..c6b26a59d22 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessor.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessor.java @@ -277,7 +277,11 @@ public CoordinatorEventProcessor( this.internalListenerName = conf.getString(ConfigOptions.INTERNAL_LISTENER_NAME); this.rebalanceManager = new RebalanceManager( - this, zooKeeperClient, coordinatorEventManager, SystemClock.getInstance()); + this, + zooKeeperClient, + coordinatorEventManager, + SystemClock.getInstance(), + coordinatorMetricGroup); this.offlineLeaderRetryDelayMs = conf.get(ConfigOptions.COORDINATOR_OFFLINE_LEADER_RETRY_DELAY).toMillis(); if (offlineLeaderRetryDelayMs <= 0) { diff --git a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java index ce95dd862a0..14c04a31bbd 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java @@ -25,6 +25,10 @@ import org.apache.fluss.cluster.rebalance.ServerTag; import org.apache.fluss.exception.NoRebalanceInProgressException; import org.apache.fluss.metadata.TableBucket; +import org.apache.fluss.metrics.Counter; +import org.apache.fluss.metrics.MetricNames; +import org.apache.fluss.metrics.ThreadSafeSimpleCounter; +import org.apache.fluss.metrics.groups.MetricGroup; import org.apache.fluss.server.coordinator.CoordinatorContext; import org.apache.fluss.server.coordinator.CoordinatorEventProcessor; import org.apache.fluss.server.coordinator.event.EventManager; @@ -36,6 +40,7 @@ import org.apache.fluss.server.coordinator.rebalance.model.RackModel; import org.apache.fluss.server.coordinator.rebalance.model.ServerModel; import org.apache.fluss.server.metadata.ServerInfo; +import org.apache.fluss.server.metrics.group.CoordinatorMetricGroup; import org.apache.fluss.server.zk.ZooKeeperClient; import org.apache.fluss.server.zk.data.LeaderAndIsr; import org.apache.fluss.server.zk.data.RebalanceTask; @@ -65,9 +70,11 @@ import static org.apache.fluss.cluster.rebalance.RebalanceStatus.CANCELED; import static org.apache.fluss.cluster.rebalance.RebalanceStatus.COMPLETED; +import static org.apache.fluss.cluster.rebalance.RebalanceStatus.FAILED; import static org.apache.fluss.cluster.rebalance.RebalanceStatus.FINAL_STATUSES; import static org.apache.fluss.cluster.rebalance.RebalanceStatus.NOT_STARTED; import static org.apache.fluss.cluster.rebalance.RebalanceStatus.REBALANCING; +import static org.apache.fluss.cluster.rebalance.RebalanceStatus.TIMEOUT; import static org.apache.fluss.utils.Preconditions.checkArgument; import static org.apache.fluss.utils.Preconditions.checkNotNull; @@ -90,6 +97,10 @@ public class RebalanceManager { private final EventManager eventManager; private final Clock clock; private final ScheduledExecutorService timeoutChecker; + private final MetricGroup metricGroup; + private final Counter rebalancesCompleted; + private final Counter rebalancesFailed; + private final Counter rebalancesCanceled; /** A queue of in progress table bucket to rebalance. */ private final Queue inProgressRebalanceTasksQueue = new ArrayDeque<>(); @@ -104,6 +115,7 @@ public class RebalanceManager { private final GoalOptimizer goalOptimizer; private volatile long registerTime; + private volatile boolean bucketResultsAvailable; private volatile @Nullable RebalanceStatus rebalanceStatus; private volatile @Nullable String currentRebalanceId; private volatile boolean isClosed = false; @@ -126,12 +138,14 @@ public RebalanceManager( CoordinatorEventProcessor eventProcessor, ZooKeeperClient zkClient, EventManager eventManager, - Clock clock) { + Clock clock, + CoordinatorMetricGroup coordinatorMetricGroup) { this( eventProcessor, zkClient, eventManager, clock, + coordinatorMetricGroup, // TODO: Reuse the CoordinatorServer shared scheduler for this lightweight // coordinator timeout checker instead of creating a component-owned scheduler. Executors.newScheduledThreadPool( @@ -144,6 +158,7 @@ public RebalanceManager( ZooKeeperClient zkClient, EventManager eventManager, Clock clock, + CoordinatorMetricGroup coordinatorMetricGroup, ScheduledExecutorService timeoutChecker) { this.eventProcessor = eventProcessor; this.zkClient = zkClient; @@ -151,6 +166,50 @@ public RebalanceManager( this.clock = clock == null ? SystemClock.getInstance() : clock; this.timeoutChecker = timeoutChecker; this.goalOptimizer = new GoalOptimizer(); + this.metricGroup = coordinatorMetricGroup.getOrAddRebalanceMetricGroup(); + this.rebalancesCompleted = + metricGroup.counter( + MetricNames.REBALANCES_COMPLETED_TOTAL, new ThreadSafeSimpleCounter()); + this.rebalancesFailed = + metricGroup.counter( + MetricNames.REBALANCES_FAILED_TOTAL, new ThreadSafeSimpleCounter()); + this.rebalancesCanceled = + metricGroup.counter( + MetricNames.REBALANCES_CANCELED_TOTAL, new ThreadSafeSimpleCounter()); + registerMetrics(); + } + + private void registerMetrics() { + metricGroup.gauge( + MetricNames.REBALANCE_IN_PROGRESS, () -> rebalanceStatus == REBALANCING ? 1 : 0); + metricGroup.gauge(MetricNames.REBALANCE_BUCKETS_PENDING, inProgressRebalanceTasks::size); + metricGroup.gauge( + MetricNames.REBALANCE_BUCKETS_COMPLETED, () -> finishedBucketCount(COMPLETED)); + metricGroup.gauge(MetricNames.REBALANCE_BUCKETS_FAILED, () -> finishedBucketCount(FAILED)); + metricGroup.gauge( + MetricNames.REBALANCE_BUCKETS_TIMED_OUT, () -> finishedBucketCount(TIMEOUT)); + metricGroup.gauge( + MetricNames.REBALANCE_DURATION_MS, + () -> + rebalanceStatus == REBALANCING + ? Math.max(0L, clock.milliseconds() - registerTime) + : 0L); + metricGroup.gauge(MetricNames.INFLIGHT_BUCKET_DURATION_MS, this::inflightBucketDurationMs); + } + + private long finishedBucketCount(RebalanceStatus status) { + if (!bucketResultsAvailable) { + return 0L; + } + return finishedRebalanceTasks.values().stream() + .filter(result -> result.status() == status) + .count(); + } + + private long inflightBucketDurationMs() { + TableBucket bucket = inflightTaskBucket; + long startMs = inflightTaskStartMs; + return bucket == null || startMs < 0 ? 0L : Math.max(0L, clock.milliseconds() - startMs); } public void startup() { @@ -194,11 +253,12 @@ public void registerRebalance( Map rebalancePlan, RebalanceStatus newStatus) { checkNotClosed(); - registerTime = System.currentTimeMillis(); + registerTime = clock.milliseconds(); // first clear all exists tasks. inProgressRebalanceTasks.clear(); inProgressRebalanceTasksQueue.clear(); finishedRebalanceTasks.clear(); + bucketResultsAvailable = !FINAL_STATUSES.contains(newStatus); // Clear gate (bucket) first, then data (startMs). inflightTaskBucket = null; inflightTaskStartMs = -1; @@ -320,6 +380,9 @@ public void cancelRebalance(@Nullable String rebalanceId) { LOG.error("Error when delete rebalance plan from zookeeper.", e); } + if (rebalanceStatus == REBALANCING) { + rebalancesCanceled.inc(); + } rebalanceStatus = CANCELED; inProgressRebalanceTasksQueue.clear(); inProgressRebalanceTasks.clear(); @@ -407,6 +470,13 @@ private void completeRebalance() { LOG.error("Error when update rebalance plan from zookeeper.", e); } + if (bucketResultsAvailable) { + if (finishedBucketCount(FAILED) > 0 || finishedBucketCount(TIMEOUT) > 0) { + rebalancesFailed.inc(); + } else { + rebalancesCompleted.inc(); + } + } rebalanceStatus = COMPLETED; inProgressRebalanceTasks.clear(); inProgressRebalanceTasksQueue.clear(); @@ -414,7 +484,7 @@ private void completeRebalance() { // Here, it will not clear finishedRebalanceTasks, because it will be used by // listRebalanceProgress. It will be cleared when next register. - LOG.info("Rebalance complete with {} ms.", System.currentTimeMillis() - registerTime); + LOG.info("Rebalance complete with {} ms.", clock.milliseconds() - registerTime); } private ClusterModel buildClusterModel(CoordinatorContext coordinatorContext) { @@ -519,6 +589,7 @@ private void checkNotClosed() { public void close() { isClosed = true; timeoutChecker.shutdownNow(); + metricGroup.close(); } @VisibleForTesting diff --git a/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/CoordinatorMetricGroup.java b/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/CoordinatorMetricGroup.java index 31120962f6d..0d444373884 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/CoordinatorMetricGroup.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/CoordinatorMetricGroup.java @@ -54,6 +54,8 @@ public class CoordinatorMetricGroup extends AbstractMetricGroup { private final Map, CoordinatorEventMetricGroup> eventMetricGroups = new ConcurrentHashMap<>(); + private @Nullable AbstractMetricGroup rebalanceMetricGroup; + public CoordinatorMetricGroup( MetricRegistry registry, String clusterId, String hostname, String serverId) { super(registry, new String[] {clusterId, hostname, NAME}, null); @@ -80,6 +82,29 @@ public CoordinatorEventMetricGroup getOrAddEventTypeMetricGroup( eventClass, e -> new CoordinatorEventMetricGroup(registry, eventClass, this)); } + /** + * Returns the leader-scoped rebalance metric group, retaining the coordinator metric scope. The + * rebalance manager closes this group on leadership loss so a new leader term can register + * fresh gauges and counters without replacing the server metric group. + */ + public synchronized AbstractMetricGroup getOrAddRebalanceMetricGroup() { + if (rebalanceMetricGroup == null || rebalanceMetricGroup.isClosed()) { + rebalanceMetricGroup = new RebalanceMetricGroup(registry, this); + if (isClosed()) { + rebalanceMetricGroup.close(); + } + } + return rebalanceMetricGroup; + } + + @Override + public synchronized void close() { + if (rebalanceMetricGroup != null) { + rebalanceMetricGroup.close(); + } + super.close(); + } + // ------------------------------------------------------------------------ // table buckets groups // ------------------------------------------------------------------------ @@ -125,6 +150,19 @@ public void removeTablePartitionMetricsGroup( } } + /** Rebalance metrics share the coordinator scope but have a leader-specific lifetime. */ + private static class RebalanceMetricGroup extends AbstractMetricGroup { + + private RebalanceMetricGroup(MetricRegistry registry, CoordinatorMetricGroup parent) { + super(registry, parent.getScopeComponents(), parent); + } + + @Override + protected String getGroupName(CharacterFilter filter) { + return ""; + } + } + /** The metric group for table. */ private static class SimpleTableMetricGroup extends AbstractMetricGroup { diff --git a/fluss-server/src/test/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManagerTest.java b/fluss-server/src/test/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManagerTest.java index 3bf3680fda6..e3dd9aec4b5 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManagerTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManagerTest.java @@ -23,6 +23,12 @@ import org.apache.fluss.config.ConfigOptions; import org.apache.fluss.config.Configuration; import org.apache.fluss.metadata.TableBucket; +import org.apache.fluss.metrics.Counter; +import org.apache.fluss.metrics.Gauge; +import org.apache.fluss.metrics.MetricNames; +import org.apache.fluss.metrics.groups.AbstractMetricGroup; +import org.apache.fluss.metrics.registry.MetricRegistry; +import org.apache.fluss.metrics.registry.NOPMetricRegistry; import org.apache.fluss.server.coordinator.AutoPartitionManager; import org.apache.fluss.server.coordinator.CoordinatorContext; import org.apache.fluss.server.coordinator.CoordinatorEventProcessor; @@ -38,6 +44,7 @@ import org.apache.fluss.server.coordinator.lease.KvSnapshotLeaseManager; import org.apache.fluss.server.coordinator.remote.RemoteDirDynamicLoader; import org.apache.fluss.server.metadata.CoordinatorMetadataCache; +import org.apache.fluss.server.metrics.group.CoordinatorMetricGroup; import org.apache.fluss.server.metrics.group.TestingMetricGroups; import org.apache.fluss.server.zk.NOPErrorHandler; import org.apache.fluss.server.zk.ZkEpoch; @@ -56,10 +63,13 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.RegisterExtension; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; import java.time.Duration; import java.util.ArrayList; import java.util.Arrays; +import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -67,9 +77,16 @@ import java.util.concurrent.ScheduledThreadPoolExecutor; import static org.apache.fluss.cluster.rebalance.RebalanceStatus.COMPLETED; +import static org.apache.fluss.cluster.rebalance.RebalanceStatus.FAILED; import static org.apache.fluss.cluster.rebalance.RebalanceStatus.NOT_STARTED; import static org.apache.fluss.cluster.rebalance.RebalanceStatus.TIMEOUT; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; /** Test for {@link RebalanceManager}. */ public class RebalanceManagerTest { @@ -88,6 +105,10 @@ public class RebalanceManagerTest { private ReplicaCapacityController replicaCapacityController; private LakeTableTieringManager lakeTableTieringManager; private RebalanceManager rebalanceManager; + private CoordinatorMetricGroup coordinatorMetricGroup; + private AbstractMetricGroup rebalanceMetricGroup; + private MetricRegistry metricRegistry; + private ManualClock metricClock; private KvSnapshotLeaseManager kvSnapshotLeaseManager; private Scheduler scheduler; @@ -134,18 +155,25 @@ void beforeEach() { new LakeTableTieringManager(TestingMetricGroups.LAKE_TIERING_METRICS); CoordinatorEventProcessor eventProcessor = buildCoordinatorEventProcessor(conf); RecordingEventManager recordingEventManager = new RecordingEventManager(); + metricRegistry = mock(MetricRegistry.class); + coordinatorMetricGroup = + new CoordinatorMetricGroup(metricRegistry, "cluster", "host", "coordinator"); + metricClock = new ManualClock(0L); rebalanceManager = new RebalanceManager( eventProcessor, zookeeperClient, recordingEventManager, - SystemClock.getInstance()); + metricClock, + coordinatorMetricGroup); + rebalanceMetricGroup = coordinatorMetricGroup.getOrAddRebalanceMetricGroup(); rebalanceManager.startup(); } @AfterEach void afterEach() throws Exception { rebalanceManager.close(); + coordinatorMetricGroup.close(); if (scheduler != null) { scheduler.shutdown(); } @@ -192,7 +220,12 @@ void testStartupQueuesRecoverRebalanceEvent() throws Exception { RebalanceManager manager = new RebalanceManager( - eventProcessor, zookeeperClient, eventManager, clock, executor); + eventProcessor, + zookeeperClient, + eventManager, + clock, + createCoordinatorMetricGroup(), + executor); // If startup() finds a pending rebalance task in ZooKeeper, it should enqueue a // RecoverRebalanceEvent to be processed by the coordinator event thread, instead of // calling registerRebalance() directly on the startup thread. @@ -229,7 +262,12 @@ void testTimeoutEnqueuesEvent() throws Exception { RebalanceManager manager = new RebalanceManager( - eventProcessor, zookeeperClient, eventManager, clock, executor); + eventProcessor, + zookeeperClient, + eventManager, + clock, + createCoordinatorMetricGroup(), + executor); manager.startup(); TableBucket tb1 = new TableBucket(1L, 0); @@ -281,7 +319,12 @@ void testTimeoutAfterCompletionIsNoOp() throws Exception { RebalanceManager manager = new RebalanceManager( - eventProcessor, zookeeperClient, eventManager, clock, executor); + eventProcessor, + zookeeperClient, + eventManager, + clock, + createCoordinatorMetricGroup(), + executor); manager.startup(); TableBucket tb1 = new TableBucket(1L, 0); @@ -318,7 +361,12 @@ void testTimeoutTreatsTaskAsCompleted() throws Exception { RebalanceManager manager = new RebalanceManager( - eventProcessor, zookeeperClient, eventManager, clock, executor); + eventProcessor, + zookeeperClient, + eventManager, + clock, + createCoordinatorMetricGroup(), + executor); manager.startup(); TableBucket tb1 = new TableBucket(1L, 0); @@ -356,6 +404,268 @@ void testTimeoutTreatsTaskAsCompleted() throws Exception { manager.close(); } + @Test + void testRebalanceMetricsLifecycle() { + assertThat(rebalanceMetricGroup.getLogicalScope(value -> value, '_')) + .isEqualTo("coordinator"); + assertThat(rebalanceMetricGroup.getScopeComponents()) + .containsExactly(coordinatorMetricGroup.getScopeComponents()); + assertThat(rebalanceMetricGroup.getAllVariables()) + .isEqualTo(coordinatorMetricGroup.getAllVariables()); + assertThat(rebalanceMetricGroup.getMetrics()).hasSize(10); + rebalanceMetricGroup + .getMetrics() + .values() + .forEach( + metric -> { + if (metric instanceof Gauge) { + assertThat(((Number) ((Gauge) metric).getValue()).longValue()) + .isZero(); + } else { + assertThat(((Counter) metric).getCount()).isZero(); + } + }); + + Map plan = createRebalancePlan(2); + List buckets = new ArrayList<>(plan.keySet()); + rebalanceManager.registerRebalance("success", plan, NOT_STARTED); + assertThat(gaugeValue(MetricNames.REBALANCE_IN_PROGRESS)).isEqualTo(1); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_PENDING)).isEqualTo(2); + + metricClock.advanceTime(Duration.ofMillis(500)); + assertThat(gaugeValue(MetricNames.REBALANCE_DURATION_MS)).isEqualTo(500); + assertThat(gaugeValue(MetricNames.INFLIGHT_BUCKET_DURATION_MS)).isEqualTo(500); + + rebalanceManager.finishRebalanceTask(buckets.get(0), COMPLETED); + metricClock.advanceTime(Duration.ofMillis(200)); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_PENDING)).isEqualTo(1); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_COMPLETED)).isEqualTo(1); + assertThat(gaugeValue(MetricNames.REBALANCE_DURATION_MS)).isEqualTo(700); + assertThat(gaugeValue(MetricNames.INFLIGHT_BUCKET_DURATION_MS)).isEqualTo(200); + assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isZero(); + + rebalanceManager.finishRebalanceTask(buckets.get(1), COMPLETED); + rebalanceManager.finishRebalanceTask(buckets.get(1), COMPLETED); + assertThat(gaugeValue(MetricNames.REBALANCE_IN_PROGRESS)).isZero(); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_PENDING)).isZero(); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_COMPLETED)).isEqualTo(2); + assertThat(gaugeValue(MetricNames.REBALANCE_DURATION_MS)).isZero(); + assertThat(gaugeValue(MetricNames.INFLIGHT_BUCKET_DURATION_MS)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isEqualTo(1); + + rebalanceManager.registerRebalance("empty", Collections.emptyMap(), NOT_STARTED); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_COMPLETED)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isEqualTo(2); + assertThat(counterValue(MetricNames.REBALANCES_FAILED_TOTAL)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_CANCELED_TOTAL)).isZero(); + } + + @Test + void testRebalanceFailureMetrics() { + Map plan = createRebalancePlan(3); + List buckets = new ArrayList<>(plan.keySet()); + rebalanceManager.registerRebalance("mixed-results", plan, NOT_STARTED); + rebalanceManager.finishRebalanceTask(buckets.get(0), FAILED); + rebalanceManager.finishRebalanceTask(buckets.get(1), TIMEOUT); + + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_PENDING)).isEqualTo(1); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_FAILED)).isEqualTo(1); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_TIMED_OUT)).isEqualTo(1); + assertThat(counterValue(MetricNames.REBALANCES_FAILED_TOTAL)).isZero(); + + rebalanceManager.finishRebalanceTask(buckets.get(2), COMPLETED); + rebalanceManager.finishRebalanceTask(buckets.get(0), FAILED); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_COMPLETED)).isEqualTo(1); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_FAILED)).isEqualTo(1); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_TIMED_OUT)).isEqualTo(1); + assertThat(gaugeValue(MetricNames.REBALANCE_IN_PROGRESS)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_FAILED_TOTAL)).isEqualTo(1); + assertThat(rebalanceManager.getRebalanceStatus()).isEqualTo(COMPLETED); + + rebalanceManager.registerRebalance("next", createRebalancePlan(1), NOT_STARTED); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_COMPLETED)).isZero(); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_FAILED)).isZero(); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_TIMED_OUT)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_FAILED_TOTAL)).isEqualTo(1); + } + + @Test + void testRebalanceTimeoutMetrics() { + Map plan = createRebalancePlan(1); + TableBucket bucket = plan.keySet().iterator().next(); + rebalanceManager.registerRebalance("timeout", plan, NOT_STARTED); + metricClock.advanceTime(Duration.ofMinutes(3)); + rebalanceManager.checkTimeout(); + + assertThat(gaugeValue(MetricNames.INFLIGHT_BUCKET_DURATION_MS)).isZero(); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_PENDING)).isEqualTo(1); + assertThat(gaugeValue(MetricNames.REBALANCE_DURATION_MS)).isEqualTo(180_000); + rebalanceManager.finishRebalanceTask(bucket, TIMEOUT); + + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_TIMED_OUT)).isEqualTo(1); + assertThat(gaugeValue(MetricNames.REBALANCE_DURATION_MS)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_FAILED_TOTAL)).isEqualTo(1); + assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isZero(); + } + + @Test + void testRebalanceCancellationMetrics() { + rebalanceManager.cancelRebalance(null); + assertThat(counterValue(MetricNames.REBALANCES_CANCELED_TOTAL)).isZero(); + + Map plan = createRebalancePlan(2); + rebalanceManager.registerRebalance("cancel", plan, NOT_STARTED); + rebalanceManager.finishRebalanceTask(plan.keySet().iterator().next(), COMPLETED); + rebalanceManager.cancelRebalance("cancel"); + rebalanceManager.cancelRebalance("cancel"); + assertThat(counterValue(MetricNames.REBALANCES_CANCELED_TOTAL)).isEqualTo(1); + assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_FAILED_TOTAL)).isZero(); + assertThat(gaugeValue(MetricNames.REBALANCE_IN_PROGRESS)).isZero(); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_PENDING)).isZero(); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_COMPLETED)).isEqualTo(1); + assertThat(gaugeValue(MetricNames.REBALANCE_DURATION_MS)).isZero(); + assertThat(gaugeValue(MetricNames.INFLIGHT_BUCKET_DURATION_MS)).isZero(); + } + + @ParameterizedTest + @EnumSource( + value = RebalanceStatus.class, + names = {"COMPLETED", "CANCELED", "FAILED", "TIMEOUT"}) + void testRecoveringFinishedRebalanceDoesNotRestoreOutcomeMetrics(RebalanceStatus status) { + rebalanceManager.registerRebalance("recovered", createRebalancePlan(1), status); + assertThat(gaugeValue(MetricNames.REBALANCE_IN_PROGRESS)).isZero(); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_PENDING)).isZero(); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_COMPLETED)).isZero(); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_FAILED)).isZero(); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_TIMED_OUT)).isZero(); + assertThat(gaugeValue(MetricNames.REBALANCE_DURATION_MS)).isZero(); + assertThat(gaugeValue(MetricNames.INFLIGHT_BUCKET_DURATION_MS)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_FAILED_TOTAL)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_CANCELED_TOTAL)).isZero(); + + rebalanceManager.registerRebalance("recovered-empty", Collections.emptyMap(), status); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_COMPLETED)).isZero(); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_FAILED)).isZero(); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_TIMED_OUT)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_FAILED_TOTAL)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_CANCELED_TOTAL)).isZero(); + } + + @ParameterizedTest + @EnumSource( + value = RebalanceStatus.class, + names = {"FAILED", "TIMEOUT"}) + void testRecoveringMixedResultsDoesNotReportSuccessfulBuckets(RebalanceStatus failedStatus) + throws Exception { + Map plan = createRebalancePlan(2); + List buckets = new ArrayList<>(plan.keySet()); + String failureMetric = + failedStatus == FAILED + ? MetricNames.REBALANCE_BUCKETS_FAILED + : MetricNames.REBALANCE_BUCKETS_TIMED_OUT; + zookeeperClient.registerRebalanceTask( + new RebalanceTask("mixed-results", NOT_STARTED, plan)); + rebalanceManager.registerRebalance("mixed-results", plan, NOT_STARTED); + rebalanceManager.finishRebalanceTask(buckets.get(0), failedStatus); + rebalanceManager.finishRebalanceTask(buckets.get(1), COMPLETED); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_COMPLETED)).isEqualTo(1); + assertThat(gaugeValue(failureMetric)).isEqualTo(1); + assertThat(counterValue(MetricNames.REBALANCES_FAILED_TOTAL)).isEqualTo(1); + + rebalanceManager.close(); + RecordingEventManager eventManager = new RecordingEventManager(); + rebalanceManager = + new RebalanceManager( + mock(CoordinatorEventProcessor.class), + zookeeperClient, + eventManager, + metricClock, + coordinatorMetricGroup, + new NoOpScheduledExecutor()); + rebalanceMetricGroup = coordinatorMetricGroup.getOrAddRebalanceMetricGroup(); + rebalanceManager.startup(); + assertThat(eventManager.events).hasSize(1); + assertThat(eventManager.events.get(0)).isInstanceOf(RecoverRebalanceEvent.class); + RebalanceTask recoveredTask = + ((RecoverRebalanceEvent) eventManager.events.get(0)).getRebalanceTask(); + assertThat(recoveredTask).isEqualTo(new RebalanceTask("mixed-results", COMPLETED, plan)); + rebalanceManager.registerRebalance( + recoveredTask.getRebalanceId(), + recoveredTask.getExecutePlan(), + recoveredTask.getRebalanceStatus()); + + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_COMPLETED)).isZero(); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_FAILED)).isZero(); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_TIMED_OUT)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_FAILED_TOTAL)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_CANCELED_TOTAL)).isZero(); + assertThat(rebalanceManager.listRebalanceProgress(null).progressForBucketMap().values()) + .hasSize(2) + .extracting(RebalanceResultForBucket::status) + .containsOnly(COMPLETED); + + zookeeperClient.registerRebalanceTask(new RebalanceTask("next", NOT_STARTED, plan)); + rebalanceManager.registerRebalance("next", plan, NOT_STARTED); + rebalanceManager.finishRebalanceTask(buckets.get(0), failedStatus); + rebalanceManager.finishRebalanceTask(buckets.get(1), COMPLETED); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_COMPLETED)).isEqualTo(1); + assertThat(gaugeValue(failureMetric)).isEqualTo(1); + assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_FAILED_TOTAL)).isEqualTo(1); + } + + @Test + void testRebalanceMetricsResetOnLeadershipChange() { + rebalanceManager.registerRebalance("first-term", Collections.emptyMap(), NOT_STARTED); + assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isEqualTo(1); + AbstractMetricGroup previousMetricGroup = rebalanceMetricGroup; + rebalanceManager.close(); + assertThat(previousMetricGroup.isClosed()).isTrue(); + assertThat(coordinatorMetricGroup.isClosed()).isFalse(); + verify(metricRegistry, times(10)).unregister(any(), anyString(), eq(previousMetricGroup)); + + rebalanceManager = + new RebalanceManager( + mock(CoordinatorEventProcessor.class), + zookeeperClient, + new RecordingEventManager(), + metricClock, + coordinatorMetricGroup, + new NoOpScheduledExecutor()); + rebalanceMetricGroup = coordinatorMetricGroup.getOrAddRebalanceMetricGroup(); + assertThat(rebalanceMetricGroup).isNotSameAs(previousMetricGroup); + verify(metricRegistry, times(10)).register(any(), anyString(), eq(rebalanceMetricGroup)); + assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isZero(); + + Map plan = createRebalancePlan(1); + rebalanceManager.registerRebalance("recovered-active", plan, NOT_STARTED); + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_PENDING)).isEqualTo(1); + rebalanceManager.finishRebalanceTask(plan.keySet().iterator().next(), COMPLETED); + assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isEqualTo(1); + + coordinatorMetricGroup.close(); + assertThat(rebalanceMetricGroup.isClosed()).isTrue(); + assertThat(coordinatorMetricGroup.getOrAddRebalanceMetricGroup().isClosed()).isTrue(); + } + + private long gaugeValue(String metricName) { + return ((Number) ((Gauge) rebalanceMetricGroup.getMetrics().get(metricName)).getValue()) + .longValue(); + } + + private long counterValue(String metricName) { + return ((Counter) rebalanceMetricGroup.getMetrics().get(metricName)).getCount(); + } + + private static CoordinatorMetricGroup createCoordinatorMetricGroup() { + return new CoordinatorMetricGroup(NOPMetricRegistry.INSTANCE, "cluster", "host", "0"); + } + private CoordinatorEventProcessor buildCoordinatorEventProcessor(Configuration conf) { return new CoordinatorEventProcessor( zookeeperClient, @@ -365,7 +675,7 @@ private CoordinatorEventProcessor buildCoordinatorEventProcessor(Configuration c replicaCapacityController, autoPartitionManager, lakeTableTieringManager, - TestingMetricGroups.COORDINATOR_METRICS, + createCoordinatorMetricGroup(), conf, Executors.newFixedThreadPool(1, new ExecutorThreadFactory("test-coordinator-io")), metadataManager, diff --git a/website/docs/maintenance/observability/monitor-metrics.md b/website/docs/maintenance/observability/monitor-metrics.md index 0d31c09dd3b..942420bd337 100644 --- a/website/docs/maintenance/observability/monitor-metrics.md +++ b/website/docs/maintenance/observability/monitor-metrics.md @@ -449,6 +449,32 @@ Some metrics might not be exposed when using other JVM implementations (e.g. IBM +#### Rebalance Metrics + +Rebalance metrics use the `coordinator` scope with no additional infix or per-table/per-bucket labels. +They are registered only on the active coordinator and removed when it loses leadership. +Gauges describe the most recently registered rebalance. Finished bucket counts include only outcomes +observed in the current coordinator leader term, remain available after completion or cancellation, +and are cleared when the next rebalance is registered. Recovering a previously finished rebalance +leaves these counts at 0 because per-bucket outcomes are not persisted. These zeros indicate no +outcomes observed by this leader, not that the historical rebalance had no failures or timeouts. +Counters accumulate across rebalances within a coordinator leader term and reset on failover. +Recovering a previously finished rebalance does not increment the counters; a recovered unfinished +rebalance is counted when it finishes in the new leader term. + +| Metrics | Description | Type | +| --- | --- | --- | +| `rebalanceInProgress` | 1 while a rebalance is running; otherwise 0. | Gauge | +| `rebalanceBucketsPending` | Number of unfinished buckets, including the currently running bucket. | Gauge | +| `rebalanceBucketsCompleted` | Number of buckets observed to complete successfully by this leader in the most recently registered rebalance. | Gauge | +| `rebalanceBucketsFailed` | Number of buckets observed to fail by this leader in the most recently registered rebalance, excluding timeouts. | Gauge | +| `rebalanceBucketsTimedOut` | Number of buckets observed to time out by this leader in the most recently registered rebalance. | Gauge | +| `rebalanceDurationMs` | Elapsed milliseconds since the running rebalance was registered on this leader; 0 when idle. | Gauge | +| `inflightBucketDurationMs` | Elapsed milliseconds for the currently running bucket; 0 when no bucket is running. | Gauge | +| `rebalancesCompletedTotal` | Number of rebalances that finished without failed or timed-out buckets, including empty plans. | Counter | +| `rebalancesFailedTotal` | Number of rebalances that finished with at least one failed or timed-out bucket. This is an outcome metric; the existing Admin API can still report the rebalance status as `COMPLETED` because execution has finished. | Counter | +| `rebalancesCanceledTotal` | Number of running rebalances canceled on this leader. Repeated cancellation and cancellation when idle do not increment it. | Counter | + ### Tablet Server From 6db197a45722a5f759c5781e50462fa78e13319b Mon Sep 17 00:00:00 2001 From: Allen Cao <258265534+heyallencao@users.noreply.github.com> Date: Thu, 17 Sep 2026 20:05:35 +0800 Subject: [PATCH 2/4] [server] Address rebalance metrics review comments --- .../rebalance/RebalanceManager.java | 15 ++-- .../observability/monitor-metrics.md | 73 ++++++++++++++++--- 2 files changed, 71 insertions(+), 17 deletions(-) diff --git a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java index 14c04a31bbd..2bae2fec8ba 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java @@ -206,6 +206,11 @@ private long finishedBucketCount(RebalanceStatus status) { .count(); } + private boolean hasFinishedBucket(RebalanceStatus status) { + return finishedRebalanceTasks.values().stream() + .anyMatch(result -> result.status() == status); + } + private long inflightBucketDurationMs() { TableBucket bucket = inflightTaskBucket; long startMs = inflightTaskStartMs; @@ -406,20 +411,20 @@ public RebalanceTask generateRebalanceTask(List goalsByPriority) { String rebalanceId = UUID.randomUUID().toString(); try { // Generate the latest cluster model. - long startTime = System.currentTimeMillis(); + long startTime = clock.milliseconds(); ClusterModel clusterModel = buildClusterModel(eventProcessor.getCoordinatorContext()); LOG.info( "Build cluster model for rebalance id {} with {} ms.", rebalanceId, - System.currentTimeMillis() - startTime); + clock.milliseconds() - startTime); // do optimize. - startTime = System.currentTimeMillis(); + startTime = clock.milliseconds(); rebalancePlanForBuckets = goalOptimizer.doOptimizeOnce(clusterModel, goalsByPriority); LOG.info( "Do optimize for rebalance id {} with {} ms.", rebalanceId, - System.currentTimeMillis() - startTime); + clock.milliseconds() - startTime); } catch (Exception e) { LOG.error("Failed to generate rebalance plan.", e); throw e; @@ -471,7 +476,7 @@ private void completeRebalance() { } if (bucketResultsAvailable) { - if (finishedBucketCount(FAILED) > 0 || finishedBucketCount(TIMEOUT) > 0) { + if (hasFinishedBucket(FAILED) || hasFinishedBucket(TIMEOUT)) { rebalancesFailed.inc(); } else { rebalancesCompleted.inc(); diff --git a/website/docs/maintenance/observability/monitor-metrics.md b/website/docs/maintenance/observability/monitor-metrics.md index 942420bd337..2e540214b2a 100644 --- a/website/docs/maintenance/observability/monitor-metrics.md +++ b/website/docs/maintenance/observability/monitor-metrics.md @@ -462,18 +462,67 @@ Counters accumulate across rebalances within a coordinator leader term and reset Recovering a previously finished rebalance does not increment the counters; a recovered unfinished rebalance is counted when it finishes in the new leader term. -| Metrics | Description | Type | -| --- | --- | --- | -| `rebalanceInProgress` | 1 while a rebalance is running; otherwise 0. | Gauge | -| `rebalanceBucketsPending` | Number of unfinished buckets, including the currently running bucket. | Gauge | -| `rebalanceBucketsCompleted` | Number of buckets observed to complete successfully by this leader in the most recently registered rebalance. | Gauge | -| `rebalanceBucketsFailed` | Number of buckets observed to fail by this leader in the most recently registered rebalance, excluding timeouts. | Gauge | -| `rebalanceBucketsTimedOut` | Number of buckets observed to time out by this leader in the most recently registered rebalance. | Gauge | -| `rebalanceDurationMs` | Elapsed milliseconds since the running rebalance was registered on this leader; 0 when idle. | Gauge | -| `inflightBucketDurationMs` | Elapsed milliseconds for the currently running bucket; 0 when no bucket is running. | Gauge | -| `rebalancesCompletedTotal` | Number of rebalances that finished without failed or timed-out buckets, including empty plans. | Counter | -| `rebalancesFailedTotal` | Number of rebalances that finished with at least one failed or timed-out bucket. This is an outcome metric; the existing Admin API can still report the rebalance status as `COMPLETED` because execution has finished. | Counter | -| `rebalancesCanceledTotal` | Number of running rebalances canceled on this leader. Repeated cancellation and cancellation when idle do not increment it. | Counter | +
+ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + +
MetricsDescriptionType
rebalanceInProgress1 while a rebalance is running; otherwise 0.Gauge
rebalanceBucketsPendingNumber of unfinished buckets, including the currently running bucket.Gauge
rebalanceBucketsCompletedNumber of buckets observed to complete successfully by this leader in the most recently registered rebalance.Gauge
rebalanceBucketsFailedNumber of buckets observed to fail by this leader in the most recently registered rebalance, excluding timeouts.Gauge
rebalanceBucketsTimedOutNumber of buckets observed to time out by this leader in the most recently registered rebalance.Gauge
rebalanceDurationMsElapsed milliseconds since the running rebalance was registered on this leader; 0 when idle.Gauge
inflightBucketDurationMsElapsed milliseconds for the currently running bucket; 0 when no bucket is running.Gauge
rebalancesCompletedTotalNumber of rebalances that finished without failed or timed-out buckets, including empty plans.Counter
rebalancesFailedTotalNumber of rebalances that finished with at least one failed or timed-out bucket. This is an outcome metric; the existing Admin API can still report the rebalance status as COMPLETED because execution has finished.Counter
rebalancesCanceledTotalNumber of running rebalances canceled on this leader. Repeated cancellation and cancellation when idle do not increment it.Counter
### Tablet Server From 595af653bf2654b6e44ccb825b49896ca5881664 Mon Sep 17 00:00:00 2001 From: "caoanfu.1" Date: Fri, 18 Sep 2026 16:52:03 +0800 Subject: [PATCH 3/4] [server] Keep rebalance metrics across leader terms Own RebalanceMetrics in CoordinatorMetricGroup and register metrics once for the coordinator process lifetime. Preserve counters across leader terms and return zero from gauges on standby. Bind managers on startup and unbind on close. Update lifecycle and recovery tests and document the process-scoped semantics. --- .../CoordinatorEventProcessor.java | 2 +- .../rebalance/RebalanceManager.java | 69 +++---- .../rebalance/RebalanceMetrics.java | 96 +++++++++ .../metrics/group/CoordinatorMetricGroup.java | 41 +--- .../rebalance/RebalanceManagerTest.java | 184 ++++++++++++++---- .../observability/monitor-metrics.md | 23 ++- 6 files changed, 287 insertions(+), 128 deletions(-) create mode 100644 fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceMetrics.java diff --git a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessor.java b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessor.java index c6b26a59d22..6eb44ab4328 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessor.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessor.java @@ -281,7 +281,7 @@ public CoordinatorEventProcessor( zooKeeperClient, coordinatorEventManager, SystemClock.getInstance(), - coordinatorMetricGroup); + coordinatorMetricGroup.getRebalanceMetrics()); this.offlineLeaderRetryDelayMs = conf.get(ConfigOptions.COORDINATOR_OFFLINE_LEADER_RETRY_DELAY).toMillis(); if (offlineLeaderRetryDelayMs <= 0) { diff --git a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java index 2bae2fec8ba..94325d2e78d 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java @@ -25,10 +25,6 @@ import org.apache.fluss.cluster.rebalance.ServerTag; import org.apache.fluss.exception.NoRebalanceInProgressException; import org.apache.fluss.metadata.TableBucket; -import org.apache.fluss.metrics.Counter; -import org.apache.fluss.metrics.MetricNames; -import org.apache.fluss.metrics.ThreadSafeSimpleCounter; -import org.apache.fluss.metrics.groups.MetricGroup; import org.apache.fluss.server.coordinator.CoordinatorContext; import org.apache.fluss.server.coordinator.CoordinatorEventProcessor; import org.apache.fluss.server.coordinator.event.EventManager; @@ -40,7 +36,6 @@ import org.apache.fluss.server.coordinator.rebalance.model.RackModel; import org.apache.fluss.server.coordinator.rebalance.model.ServerModel; import org.apache.fluss.server.metadata.ServerInfo; -import org.apache.fluss.server.metrics.group.CoordinatorMetricGroup; import org.apache.fluss.server.zk.ZooKeeperClient; import org.apache.fluss.server.zk.data.LeaderAndIsr; import org.apache.fluss.server.zk.data.RebalanceTask; @@ -97,10 +92,7 @@ public class RebalanceManager { private final EventManager eventManager; private final Clock clock; private final ScheduledExecutorService timeoutChecker; - private final MetricGroup metricGroup; - private final Counter rebalancesCompleted; - private final Counter rebalancesFailed; - private final Counter rebalancesCanceled; + private final RebalanceMetrics metrics; /** A queue of in progress table bucket to rebalance. */ private final Queue inProgressRebalanceTasksQueue = new ArrayDeque<>(); @@ -139,13 +131,13 @@ public RebalanceManager( ZooKeeperClient zkClient, EventManager eventManager, Clock clock, - CoordinatorMetricGroup coordinatorMetricGroup) { + RebalanceMetrics metrics) { this( eventProcessor, zkClient, eventManager, clock, - coordinatorMetricGroup, + metrics, // TODO: Reuse the CoordinatorServer shared scheduler for this lightweight // coordinator timeout checker instead of creating a component-owned scheduler. Executors.newScheduledThreadPool( @@ -158,7 +150,7 @@ public RebalanceManager( ZooKeeperClient zkClient, EventManager eventManager, Clock clock, - CoordinatorMetricGroup coordinatorMetricGroup, + RebalanceMetrics metrics, ScheduledExecutorService timeoutChecker) { this.eventProcessor = eventProcessor; this.zkClient = zkClient; @@ -166,38 +158,24 @@ public RebalanceManager( this.clock = clock == null ? SystemClock.getInstance() : clock; this.timeoutChecker = timeoutChecker; this.goalOptimizer = new GoalOptimizer(); - this.metricGroup = coordinatorMetricGroup.getOrAddRebalanceMetricGroup(); - this.rebalancesCompleted = - metricGroup.counter( - MetricNames.REBALANCES_COMPLETED_TOTAL, new ThreadSafeSimpleCounter()); - this.rebalancesFailed = - metricGroup.counter( - MetricNames.REBALANCES_FAILED_TOTAL, new ThreadSafeSimpleCounter()); - this.rebalancesCanceled = - metricGroup.counter( - MetricNames.REBALANCES_CANCELED_TOTAL, new ThreadSafeSimpleCounter()); - registerMetrics(); + this.metrics = metrics; } - private void registerMetrics() { - metricGroup.gauge( - MetricNames.REBALANCE_IN_PROGRESS, () -> rebalanceStatus == REBALANCING ? 1 : 0); - metricGroup.gauge(MetricNames.REBALANCE_BUCKETS_PENDING, inProgressRebalanceTasks::size); - metricGroup.gauge( - MetricNames.REBALANCE_BUCKETS_COMPLETED, () -> finishedBucketCount(COMPLETED)); - metricGroup.gauge(MetricNames.REBALANCE_BUCKETS_FAILED, () -> finishedBucketCount(FAILED)); - metricGroup.gauge( - MetricNames.REBALANCE_BUCKETS_TIMED_OUT, () -> finishedBucketCount(TIMEOUT)); - metricGroup.gauge( - MetricNames.REBALANCE_DURATION_MS, - () -> - rebalanceStatus == REBALANCING - ? Math.max(0L, clock.milliseconds() - registerTime) - : 0L); - metricGroup.gauge(MetricNames.INFLIGHT_BUCKET_DURATION_MS, this::inflightBucketDurationMs); + long rebalanceInProgress() { + return rebalanceStatus == REBALANCING ? 1L : 0L; } - private long finishedBucketCount(RebalanceStatus status) { + long pendingBucketCount() { + return inProgressRebalanceTasks.size(); + } + + long rebalanceDurationMs() { + return rebalanceStatus == REBALANCING + ? Math.max(0L, clock.milliseconds() - registerTime) + : 0L; + } + + long finishedBucketCount(RebalanceStatus status) { if (!bucketResultsAvailable) { return 0L; } @@ -211,7 +189,7 @@ private boolean hasFinishedBucket(RebalanceStatus status) { .anyMatch(result -> result.status() == status); } - private long inflightBucketDurationMs() { + long inflightBucketDurationMs() { TableBucket bucket = inflightTaskBucket; long startMs = inflightTaskStartMs; return bucket == null || startMs < 0 ? 0L : Math.max(0L, clock.milliseconds() - startMs); @@ -219,6 +197,7 @@ private long inflightBucketDurationMs() { public void startup() { LOG.info("Start up rebalance manager."); + metrics.bind(this); initialize(); } @@ -386,7 +365,7 @@ public void cancelRebalance(@Nullable String rebalanceId) { } if (rebalanceStatus == REBALANCING) { - rebalancesCanceled.inc(); + metrics.incRebalancesCanceled(); } rebalanceStatus = CANCELED; inProgressRebalanceTasksQueue.clear(); @@ -477,9 +456,9 @@ private void completeRebalance() { if (bucketResultsAvailable) { if (hasFinishedBucket(FAILED) || hasFinishedBucket(TIMEOUT)) { - rebalancesFailed.inc(); + metrics.incRebalancesFailed(); } else { - rebalancesCompleted.inc(); + metrics.incRebalancesCompleted(); } } rebalanceStatus = COMPLETED; @@ -593,8 +572,8 @@ private void checkNotClosed() { public void close() { isClosed = true; + metrics.unbind(this); timeoutChecker.shutdownNow(); - metricGroup.close(); } @VisibleForTesting diff --git a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceMetrics.java b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceMetrics.java new file mode 100644 index 00000000000..395a08c5bb3 --- /dev/null +++ b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceMetrics.java @@ -0,0 +1,96 @@ +/* + * 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.rebalance; + +import org.apache.fluss.metrics.Counter; +import org.apache.fluss.metrics.MetricNames; +import org.apache.fluss.metrics.ThreadSafeSimpleCounter; +import org.apache.fluss.metrics.groups.MetricGroup; + +import javax.annotation.Nullable; + +import java.util.function.ToLongFunction; + +import static org.apache.fluss.cluster.rebalance.RebalanceStatus.COMPLETED; +import static org.apache.fluss.cluster.rebalance.RebalanceStatus.FAILED; +import static org.apache.fluss.cluster.rebalance.RebalanceStatus.TIMEOUT; + +/** Rebalance metrics registered for the lifetime of a coordinator server. */ +public class RebalanceMetrics { + + private final Counter rebalancesCompleted = new ThreadSafeSimpleCounter(); + private final Counter rebalancesFailed = new ThreadSafeSimpleCounter(); + private final Counter rebalancesCanceled = new ThreadSafeSimpleCounter(); + + private volatile @Nullable RebalanceManager current; + + /** Registers the rebalance metrics on the coordinator metric group. */ + public RebalanceMetrics(MetricGroup metricGroup) { + metricGroup.counter(MetricNames.REBALANCES_COMPLETED_TOTAL, rebalancesCompleted); + metricGroup.counter(MetricNames.REBALANCES_FAILED_TOTAL, rebalancesFailed); + metricGroup.counter(MetricNames.REBALANCES_CANCELED_TOTAL, rebalancesCanceled); + metricGroup.gauge( + MetricNames.REBALANCE_IN_PROGRESS, + () -> readCurrent(RebalanceManager::rebalanceInProgress)); + metricGroup.gauge( + MetricNames.REBALANCE_BUCKETS_PENDING, + () -> readCurrent(RebalanceManager::pendingBucketCount)); + metricGroup.gauge( + MetricNames.REBALANCE_BUCKETS_COMPLETED, + () -> readCurrent(manager -> manager.finishedBucketCount(COMPLETED))); + metricGroup.gauge( + MetricNames.REBALANCE_BUCKETS_FAILED, + () -> readCurrent(manager -> manager.finishedBucketCount(FAILED))); + metricGroup.gauge( + MetricNames.REBALANCE_BUCKETS_TIMED_OUT, + () -> readCurrent(manager -> manager.finishedBucketCount(TIMEOUT))); + metricGroup.gauge( + MetricNames.REBALANCE_DURATION_MS, + () -> readCurrent(RebalanceManager::rebalanceDurationMs)); + metricGroup.gauge( + MetricNames.INFLIGHT_BUCKET_DURATION_MS, + () -> readCurrent(RebalanceManager::inflightBucketDurationMs)); + } + + void bind(RebalanceManager manager) { + current = manager; + } + + void unbind(RebalanceManager manager) { + if (current == manager) { + current = null; + } + } + + void incRebalancesCompleted() { + rebalancesCompleted.inc(); + } + + void incRebalancesFailed() { + rebalancesFailed.inc(); + } + + void incRebalancesCanceled() { + rebalancesCanceled.inc(); + } + + private long readCurrent(ToLongFunction read) { + RebalanceManager manager = current; + return manager == null ? 0L : read.applyAsLong(manager); + } +} diff --git a/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/CoordinatorMetricGroup.java b/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/CoordinatorMetricGroup.java index 0d444373884..9ade6cc31c4 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/CoordinatorMetricGroup.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/CoordinatorMetricGroup.java @@ -26,6 +26,7 @@ import org.apache.fluss.metrics.groups.MetricGroup; import org.apache.fluss.metrics.registry.MetricRegistry; import org.apache.fluss.server.coordinator.event.CoordinatorEvent; +import org.apache.fluss.server.coordinator.rebalance.RebalanceMetrics; import javax.annotation.Nullable; @@ -54,7 +55,7 @@ public class CoordinatorMetricGroup extends AbstractMetricGroup { private final Map, CoordinatorEventMetricGroup> eventMetricGroups = new ConcurrentHashMap<>(); - private @Nullable AbstractMetricGroup rebalanceMetricGroup; + private final RebalanceMetrics rebalanceMetrics; public CoordinatorMetricGroup( MetricRegistry registry, String clusterId, String hostname, String serverId) { @@ -62,6 +63,7 @@ public CoordinatorMetricGroup( this.clusterId = clusterId; this.hostname = hostname; this.serverId = serverId; + this.rebalanceMetrics = new RebalanceMetrics(this); } @Override @@ -82,27 +84,9 @@ public CoordinatorEventMetricGroup getOrAddEventTypeMetricGroup( eventClass, e -> new CoordinatorEventMetricGroup(registry, eventClass, this)); } - /** - * Returns the leader-scoped rebalance metric group, retaining the coordinator metric scope. The - * rebalance manager closes this group on leadership loss so a new leader term can register - * fresh gauges and counters without replacing the server metric group. - */ - public synchronized AbstractMetricGroup getOrAddRebalanceMetricGroup() { - if (rebalanceMetricGroup == null || rebalanceMetricGroup.isClosed()) { - rebalanceMetricGroup = new RebalanceMetricGroup(registry, this); - if (isClosed()) { - rebalanceMetricGroup.close(); - } - } - return rebalanceMetricGroup; - } - - @Override - public synchronized void close() { - if (rebalanceMetricGroup != null) { - rebalanceMetricGroup.close(); - } - super.close(); + /** Returns the rebalance metrics for this coordinator server. */ + public RebalanceMetrics getRebalanceMetrics() { + return rebalanceMetrics; } // ------------------------------------------------------------------------ @@ -150,19 +134,6 @@ public void removeTablePartitionMetricsGroup( } } - /** Rebalance metrics share the coordinator scope but have a leader-specific lifetime. */ - private static class RebalanceMetricGroup extends AbstractMetricGroup { - - private RebalanceMetricGroup(MetricRegistry registry, CoordinatorMetricGroup parent) { - super(registry, parent.getScopeComponents(), parent); - } - - @Override - protected String getGroupName(CharacterFilter filter) { - return ""; - } - } - /** The metric group for table. */ private static class SimpleTableMetricGroup extends AbstractMetricGroup { diff --git a/fluss-server/src/test/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManagerTest.java b/fluss-server/src/test/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManagerTest.java index e3dd9aec4b5..2b5ebe49f1d 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManagerTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManagerTest.java @@ -25,8 +25,8 @@ import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.metrics.Counter; import org.apache.fluss.metrics.Gauge; +import org.apache.fluss.metrics.Metric; import org.apache.fluss.metrics.MetricNames; -import org.apache.fluss.metrics.groups.AbstractMetricGroup; import org.apache.fluss.metrics.registry.MetricRegistry; import org.apache.fluss.metrics.registry.NOPMetricRegistry; import org.apache.fluss.server.coordinator.AutoPartitionManager; @@ -81,10 +81,12 @@ import static org.apache.fluss.cluster.rebalance.RebalanceStatus.NOT_STARTED; import static org.apache.fluss.cluster.rebalance.RebalanceStatus.TIMEOUT; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; @@ -106,7 +108,6 @@ public class RebalanceManagerTest { private LakeTableTieringManager lakeTableTieringManager; private RebalanceManager rebalanceManager; private CoordinatorMetricGroup coordinatorMetricGroup; - private AbstractMetricGroup rebalanceMetricGroup; private MetricRegistry metricRegistry; private ManualClock metricClock; private KvSnapshotLeaseManager kvSnapshotLeaseManager; @@ -165,8 +166,7 @@ void beforeEach() { zookeeperClient, recordingEventManager, metricClock, - coordinatorMetricGroup); - rebalanceMetricGroup = coordinatorMetricGroup.getOrAddRebalanceMetricGroup(); + coordinatorMetricGroup.getRebalanceMetrics()); rebalanceManager.startup(); } @@ -224,7 +224,7 @@ void testStartupQueuesRecoverRebalanceEvent() throws Exception { zookeeperClient, eventManager, clock, - createCoordinatorMetricGroup(), + createCoordinatorMetricGroup().getRebalanceMetrics(), executor); // If startup() finds a pending rebalance task in ZooKeeper, it should enqueue a // RecoverRebalanceEvent to be processed by the coordinator event thread, instead of @@ -266,7 +266,7 @@ void testTimeoutEnqueuesEvent() throws Exception { zookeeperClient, eventManager, clock, - createCoordinatorMetricGroup(), + createCoordinatorMetricGroup().getRebalanceMetrics(), executor); manager.startup(); @@ -323,7 +323,7 @@ void testTimeoutAfterCompletionIsNoOp() throws Exception { zookeeperClient, eventManager, clock, - createCoordinatorMetricGroup(), + createCoordinatorMetricGroup().getRebalanceMetrics(), executor); manager.startup(); @@ -365,7 +365,7 @@ void testTimeoutTreatsTaskAsCompleted() throws Exception { zookeeperClient, eventManager, clock, - createCoordinatorMetricGroup(), + createCoordinatorMetricGroup().getRebalanceMetrics(), executor); manager.startup(); @@ -406,14 +406,17 @@ void testTimeoutTreatsTaskAsCompleted() throws Exception { @Test void testRebalanceMetricsLifecycle() { - assertThat(rebalanceMetricGroup.getLogicalScope(value -> value, '_')) + assertThat(coordinatorMetricGroup.getLogicalScope(value -> value, '_')) .isEqualTo("coordinator"); - assertThat(rebalanceMetricGroup.getScopeComponents()) - .containsExactly(coordinatorMetricGroup.getScopeComponents()); - assertThat(rebalanceMetricGroup.getAllVariables()) - .isEqualTo(coordinatorMetricGroup.getAllVariables()); - assertThat(rebalanceMetricGroup.getMetrics()).hasSize(10); - rebalanceMetricGroup + assertThat(coordinatorMetricGroup.getScopeComponents()) + .containsExactly("cluster", "host", "coordinator"); + assertThat(coordinatorMetricGroup.getAllVariables()) + .containsOnlyKeys("cluster_id", "host", "server_id") + .containsEntry("cluster_id", "cluster") + .containsEntry("host", "host") + .containsEntry("server_id", "coordinator"); + assertThat(coordinatorMetricGroup.getMetrics()).hasSize(10); + coordinatorMetricGroup .getMetrics() .values() .forEach( @@ -584,9 +587,8 @@ void testRecoveringMixedResultsDoesNotReportSuccessfulBuckets(RebalanceStatus fa zookeeperClient, eventManager, metricClock, - coordinatorMetricGroup, + coordinatorMetricGroup.getRebalanceMetrics(), new NoOpScheduledExecutor()); - rebalanceMetricGroup = coordinatorMetricGroup.getOrAddRebalanceMetricGroup(); rebalanceManager.startup(); assertThat(eventManager.events).hasSize(1); assertThat(eventManager.events.get(0)).isInstanceOf(RecoverRebalanceEvent.class); @@ -602,7 +604,7 @@ void testRecoveringMixedResultsDoesNotReportSuccessfulBuckets(RebalanceStatus fa assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_FAILED)).isZero(); assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_TIMED_OUT)).isZero(); assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isZero(); - assertThat(counterValue(MetricNames.REBALANCES_FAILED_TOTAL)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_FAILED_TOTAL)).isEqualTo(1); assertThat(counterValue(MetricNames.REBALANCES_CANCELED_TOTAL)).isZero(); assertThat(rebalanceManager.listRebalanceProgress(null).progressForBucketMap().values()) .hasSize(2) @@ -616,50 +618,153 @@ void testRecoveringMixedResultsDoesNotReportSuccessfulBuckets(RebalanceStatus fa assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_COMPLETED)).isEqualTo(1); assertThat(gaugeValue(failureMetric)).isEqualTo(1); assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isZero(); - assertThat(counterValue(MetricNames.REBALANCES_FAILED_TOTAL)).isEqualTo(1); + assertThat(counterValue(MetricNames.REBALANCES_FAILED_TOTAL)).isEqualTo(2); } @Test - void testRebalanceMetricsResetOnLeadershipChange() { + void testRebalanceMetricsSurviveLeadershipChange() throws Exception { + Map registeredMetrics = new HashMap<>(coordinatorMetricGroup.getMetrics()); rebalanceManager.registerRebalance("first-term", Collections.emptyMap(), NOT_STARTED); assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isEqualTo(1); - AbstractMetricGroup previousMetricGroup = rebalanceMetricGroup; + + Map plan = createRebalancePlan(1); + rebalanceManager.registerRebalance("failed", plan, NOT_STARTED); + rebalanceManager.finishRebalanceTask(plan.keySet().iterator().next(), FAILED); + rebalanceManager.registerRebalance("canceled", plan, NOT_STARTED); + rebalanceManager.cancelRebalance("canceled"); + rebalanceManager.registerRebalance("unfinished", plan, NOT_STARTED); + zookeeperClient.registerRebalanceTask(new RebalanceTask("unfinished", NOT_STARTED, plan)); + metricClock.advanceTime(Duration.ofMillis(100)); + + RebalanceManager previousManager = rebalanceManager; rebalanceManager.close(); - assertThat(previousMetricGroup.isClosed()).isTrue(); assertThat(coordinatorMetricGroup.isClosed()).isFalse(); - verify(metricRegistry, times(10)).unregister(any(), anyString(), eq(previousMetricGroup)); + assertGaugesAreZero(); + assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isEqualTo(1); + assertThat(counterValue(MetricNames.REBALANCES_FAILED_TOTAL)).isEqualTo(1); + assertThat(counterValue(MetricNames.REBALANCES_CANCELED_TOTAL)).isEqualTo(1); + RecordingEventManager eventManager = new RecordingEventManager(); rebalanceManager = new RebalanceManager( mock(CoordinatorEventProcessor.class), zookeeperClient, - new RecordingEventManager(), + eventManager, metricClock, - coordinatorMetricGroup, + coordinatorMetricGroup.getRebalanceMetrics(), new NoOpScheduledExecutor()); - rebalanceMetricGroup = coordinatorMetricGroup.getOrAddRebalanceMetricGroup(); - assertThat(rebalanceMetricGroup).isNotSameAs(previousMetricGroup); - verify(metricRegistry, times(10)).register(any(), anyString(), eq(rebalanceMetricGroup)); - assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isZero(); + rebalanceManager.startup(); - Map plan = createRebalancePlan(1); - rebalanceManager.registerRebalance("recovered-active", plan, NOT_STARTED); + assertThat(eventManager.events).hasSize(1); + RebalanceTask recoveredTask = + ((RecoverRebalanceEvent) eventManager.events.get(0)).getRebalanceTask(); + rebalanceManager.registerRebalance( + recoveredTask.getRebalanceId(), + recoveredTask.getExecutePlan(), + recoveredTask.getRebalanceStatus()); + previousManager.close(); assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_PENDING)).isEqualTo(1); rebalanceManager.finishRebalanceTask(plan.keySet().iterator().next(), COMPLETED); - assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isEqualTo(1); + assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isEqualTo(2); + assertThat(counterValue(MetricNames.REBALANCES_FAILED_TOTAL)).isEqualTo(1); + assertThat(counterValue(MetricNames.REBALANCES_CANCELED_TOTAL)).isEqualTo(1); + assertThat(coordinatorMetricGroup.getMetrics()).hasSize(10); + registeredMetrics.forEach( + (name, metric) -> + assertThat(coordinatorMetricGroup.getMetrics().get(name)).isSameAs(metric)); + verify(metricRegistry, times(10)).register(any(), anyString(), eq(coordinatorMetricGroup)); + verify(metricRegistry, never()).unregister(any(), anyString(), eq(coordinatorMetricGroup)); coordinatorMetricGroup.close(); - assertThat(rebalanceMetricGroup.isClosed()).isTrue(); - assertThat(coordinatorMetricGroup.getOrAddRebalanceMetricGroup().isClosed()).isTrue(); + assertThat(coordinatorMetricGroup.getMetrics()).isEmpty(); + verify(metricRegistry, times(10)) + .unregister(any(), anyString(), eq(coordinatorMetricGroup)); + + rebalanceManager.close(); + coordinatorMetricGroup = createCoordinatorMetricGroup(); + assertGaugesAreZero(); + assertThat(counterValue(MetricNames.REBALANCES_COMPLETED_TOTAL)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_FAILED_TOTAL)).isZero(); + assertThat(counterValue(MetricNames.REBALANCES_CANCELED_TOTAL)).isZero(); + } + + @Test + void testRebalanceMetricsRegisteredOnStandby() { + CoordinatorMetricGroup standbyGroup = + new CoordinatorMetricGroup(metricRegistry, "cluster", "standby", "standby"); + try { + assertThat(standbyGroup.getMetrics().values()) + .hasSize(10) + .allSatisfy( + metric -> { + if (metric instanceof Gauge) { + assertThat( + ((Number) ((Gauge) metric).getValue()) + .longValue()) + .isZero(); + } else { + assertThat(((Counter) metric).getCount()).isZero(); + } + }); + verify(metricRegistry, times(10)).register(any(), anyString(), eq(standbyGroup)); + } finally { + standbyGroup.close(); + } + } + + @Test + void testUnstartedManagerDoesNotReplaceCurrentManager() { + rebalanceManager.registerRebalance("current", createRebalancePlan(1), NOT_STARTED); + RebalanceManager unstartedManager = + new RebalanceManager( + mock(CoordinatorEventProcessor.class), + zookeeperClient, + new RecordingEventManager(), + metricClock, + coordinatorMetricGroup.getRebalanceMetrics(), + new NoOpScheduledExecutor()); + try { + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_PENDING)).isEqualTo(1); + } finally { + unstartedManager.close(); + } + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_PENDING)).isEqualTo(1); + } + + @Test + void testFailedProcessorConstructionDoesNotReplaceCurrentManager() { + rebalanceManager.registerRebalance("current", createRebalancePlan(1), NOT_STARTED); + Configuration conf = new Configuration(); + conf.set(ConfigOptions.COORDINATOR_OFFLINE_LEADER_RETRY_DELAY, Duration.ZERO); + + assertThatThrownBy(() -> buildCoordinatorEventProcessor(conf, coordinatorMetricGroup)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining(ConfigOptions.COORDINATOR_OFFLINE_LEADER_RETRY_DELAY.key()); + + assertThat(gaugeValue(MetricNames.REBALANCE_BUCKETS_PENDING)).isEqualTo(1); + } + + private void assertGaugesAreZero() { + coordinatorMetricGroup + .getMetrics() + .values() + .forEach( + metric -> { + if (metric instanceof Gauge) { + assertThat(((Number) ((Gauge) metric).getValue()).longValue()) + .isZero(); + } + }); } private long gaugeValue(String metricName) { - return ((Number) ((Gauge) rebalanceMetricGroup.getMetrics().get(metricName)).getValue()) + return ((Number) + ((Gauge) coordinatorMetricGroup.getMetrics().get(metricName)).getValue()) .longValue(); } private long counterValue(String metricName) { - return ((Counter) rebalanceMetricGroup.getMetrics().get(metricName)).getCount(); + return ((Counter) coordinatorMetricGroup.getMetrics().get(metricName)).getCount(); } private static CoordinatorMetricGroup createCoordinatorMetricGroup() { @@ -667,6 +772,11 @@ private static CoordinatorMetricGroup createCoordinatorMetricGroup() { } private CoordinatorEventProcessor buildCoordinatorEventProcessor(Configuration conf) { + return buildCoordinatorEventProcessor(conf, createCoordinatorMetricGroup()); + } + + private CoordinatorEventProcessor buildCoordinatorEventProcessor( + Configuration conf, CoordinatorMetricGroup metricGroup) { return new CoordinatorEventProcessor( zookeeperClient, serverMetadataCache, @@ -675,7 +785,7 @@ private CoordinatorEventProcessor buildCoordinatorEventProcessor(Configuration c replicaCapacityController, autoPartitionManager, lakeTableTieringManager, - createCoordinatorMetricGroup(), + metricGroup, conf, Executors.newFixedThreadPool(1, new ExecutorThreadFactory("test-coordinator-io")), metadataManager, diff --git a/website/docs/maintenance/observability/monitor-metrics.md b/website/docs/maintenance/observability/monitor-metrics.md index 2e540214b2a..a8a610745f4 100644 --- a/website/docs/maintenance/observability/monitor-metrics.md +++ b/website/docs/maintenance/observability/monitor-metrics.md @@ -452,15 +452,18 @@ Some metrics might not be exposed when using other JVM implementations (e.g. IBM #### Rebalance Metrics Rebalance metrics use the `coordinator` scope with no additional infix or per-table/per-bucket labels. -They are registered only on the active coordinator and removed when it loses leadership. -Gauges describe the most recently registered rebalance. Finished bucket counts include only outcomes -observed in the current coordinator leader term, remain available after completion or cancellation, -and are cleared when the next rebalance is registered. Recovering a previously finished rebalance -leaves these counts at 0 because per-bucket outcomes are not persisted. These zeros indicate no +They are registered when the coordinator server starts and remain registered across leadership changes. +Gauges report 0 on standby. On the active coordinator, they describe the most recently registered +rebalance. Finished bucket counts include only outcomes observed in the current leader term, remain +available after completion or cancellation within that term, and are cleared when the next rebalance +is registered. Recovering a previously finished rebalance leaves these counts at 0 because per-bucket +outcomes are not persisted. These zeros indicate no outcomes observed by this leader, not that the historical rebalance had no failures or timeouts. -Counters accumulate across rebalances within a coordinator leader term and reset on failover. +Counters accumulate across rebalances and leader terms within the same coordinator process. They +retain their values on standby and reset when the coordinator server restarts, not on leadership +changes. These are per-process counts, not a cluster-wide total transferred between coordinators. Recovering a previously finished rebalance does not increment the counters; a recovered unfinished -rebalance is counted when it finishes in the new leader term. +rebalance is counted by the coordinator that observes its completion. @@ -508,17 +511,17 @@ rebalance is counted when it finishes in the new leader term. - + - + - + From 5d23ef1e93f0e22637adc13f662ab0b555f50ac7 Mon Sep 17 00:00:00 2001 From: "allen.cao" Date: Sat, 19 Sep 2026 11:38:18 +0800 Subject: [PATCH 4/4] [server] Address follow-up rebalance metrics review comments --- .../CoordinatorEventProcessor.java | 11 +++++---- .../rebalance/RebalanceManager.java | 23 ++++++++++++++----- .../rebalance/RebalanceMetrics.java | 10 +++++++- 3 files changed, 33 insertions(+), 11 deletions(-) diff --git a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessor.java b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessor.java index 6eb44ab4328..fb0d3ca2124 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessor.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessor.java @@ -346,10 +346,13 @@ public void startup() { } public void shutdown() { - clearOfflineLeaderRetryTask(); - // close the event manager - coordinatorEventManager.close(); - rebalanceManager.close(); + try { + clearOfflineLeaderRetryTask(); + // close the event manager + coordinatorEventManager.close(); + } finally { + rebalanceManager.close(); + } onShutdown(); coordinatorContext.resetContext(); updateObservedKvLeaderReplicaCount(); diff --git a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java index 94325d2e78d..4c818de1d1f 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceManager.java @@ -49,6 +49,7 @@ import javax.annotation.Nullable; import java.util.ArrayDeque; +import java.util.EnumMap; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -62,6 +63,7 @@ import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; import static org.apache.fluss.cluster.rebalance.RebalanceStatus.CANCELED; import static org.apache.fluss.cluster.rebalance.RebalanceStatus.COMPLETED; @@ -105,6 +107,8 @@ public class RebalanceManager { private final Map finishedRebalanceTasks = new ConcurrentHashMap<>(); + private final Map finishedBucketCounts = newFinishedBucketCounts(); + private final GoalOptimizer goalOptimizer; private volatile long registerTime; private volatile boolean bucketResultsAvailable; @@ -158,7 +162,7 @@ public RebalanceManager( this.clock = clock == null ? SystemClock.getInstance() : clock; this.timeoutChecker = timeoutChecker; this.goalOptimizer = new GoalOptimizer(); - this.metrics = metrics; + this.metrics = checkNotNull(metrics, "metrics"); } long rebalanceInProgress() { @@ -179,14 +183,11 @@ long finishedBucketCount(RebalanceStatus status) { if (!bucketResultsAvailable) { return 0L; } - return finishedRebalanceTasks.values().stream() - .filter(result -> result.status() == status) - .count(); + return finishedBucketCounts.get(status).get(); } private boolean hasFinishedBucket(RebalanceStatus status) { - return finishedRebalanceTasks.values().stream() - .anyMatch(result -> result.status() == status); + return finishedBucketCounts.get(status).get() > 0; } long inflightBucketDurationMs() { @@ -242,6 +243,7 @@ public void registerRebalance( inProgressRebalanceTasks.clear(); inProgressRebalanceTasksQueue.clear(); finishedRebalanceTasks.clear(); + finishedBucketCounts.values().forEach(count -> count.set(0L)); bucketResultsAvailable = !FINAL_STATUSES.contains(newStatus); // Clear gate (bucket) first, then data (startMs). inflightTaskBucket = null; @@ -284,6 +286,7 @@ public void finishRebalanceTask(TableBucket tableBucket, RebalanceStatus statusF finishedRebalanceTasks.put( tableBucket, RebalanceResultForBucket.of(resultForBucket.plan(), statusForBucket)); + finishedBucketCounts.get(statusForBucket).incrementAndGet(); // Clear gate (bucket) first, then data (startMs). inflightTaskBucket = null; inflightTaskStartMs = -1; @@ -524,6 +527,14 @@ private RebalanceTask buildRebalanceTask( return new RebalanceTask(rebalanceId, NOT_STARTED, bucketPlan); } + private static Map newFinishedBucketCounts() { + Map counts = new EnumMap<>(RebalanceStatus.class); + for (RebalanceStatus status : RebalanceStatus.values()) { + counts.put(status, new AtomicLong()); + } + return counts; + } + private boolean isOfflineTagged(ServerTag serverTag) { return serverTag == ServerTag.PERMANENT_OFFLINE || serverTag == ServerTag.TEMPORARY_OFFLINE; } diff --git a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceMetrics.java b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceMetrics.java index 395a08c5bb3..97b1bed7799 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceMetrics.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/rebalance/RebalanceMetrics.java @@ -23,6 +23,7 @@ import org.apache.fluss.metrics.groups.MetricGroup; import javax.annotation.Nullable; +import javax.annotation.concurrent.ThreadSafe; import java.util.function.ToLongFunction; @@ -30,7 +31,14 @@ import static org.apache.fluss.cluster.rebalance.RebalanceStatus.FAILED; import static org.apache.fluss.cluster.rebalance.RebalanceStatus.TIMEOUT; -/** Rebalance metrics registered for the lifetime of a coordinator server. */ +/** + * Rebalance metrics registered for the lifetime of a coordinator server. + * + *

The coordinator event thread binds and unbinds the current {@link RebalanceManager} and + * updates counters, while metric reporter threads may concurrently read gauges via {@link + * #readCurrent(ToLongFunction)}. + */ +@ThreadSafe public class RebalanceMetrics { private final Counter rebalancesCompleted = new ThreadSafeSimpleCounter();

rebalancesCompletedTotalNumber of rebalances that finished without failed or timed-out buckets, including empty plans.Number of rebalances observed by this coordinator process to finish without failed or timed-out buckets, including empty plans, across all its leader terms. Counter
rebalancesFailedTotalNumber of rebalances that finished with at least one failed or timed-out bucket. This is an outcome metric; the existing Admin API can still report the rebalance status as COMPLETED because execution has finished.Number of rebalances observed by this coordinator process to finish with at least one failed or timed-out bucket, across all its leader terms. This is an outcome metric; the existing Admin API can still report the rebalance status as COMPLETED because execution has finished. Counter
rebalancesCanceledTotalNumber of running rebalances canceled on this leader. Repeated cancellation and cancellation when idle do not increment it.Number of running rebalances canceled by this coordinator process across all its leader terms. Repeated cancellation and cancellation when idle do not increment it. Counter