From 91215a725d03fc288e96b21ac0d3a4fbc804295f Mon Sep 17 00:00:00 2001 From: deardeng Date: Mon, 10 Aug 2026 10:47:06 +0800 Subject: [PATCH] [improvement](fe) Add cloud tablet rebalancer metrics pick from https://github.com/apache/doris/pull/66576 Record cloud tablet rebalancer round duration, allocation, and scan metrics. (cherry picked from commit 472f5c9f317702527a5a85f356126fa6089b39d3) --- .../cloud/catalog/CloudTabletRebalancer.java | 78 +++++++++++-------- .../catalog/CloudTabletRebalancerMetrics.java | 78 +++++++++++++++++++ .../org/apache/doris/metric/CloudMetrics.java | 37 +++++++++ .../org/apache/doris/metric/MetricRepo.java | 15 ++++ 4 files changed, 177 insertions(+), 31 deletions(-) create mode 100644 fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTabletRebalancerMetrics.java diff --git a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTabletRebalancer.java b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTabletRebalancer.java index 82d29f3d701f68..d1758eae96ebf7 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTabletRebalancer.java +++ b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTabletRebalancer.java @@ -82,6 +82,9 @@ public class CloudTabletRebalancer extends MasterDaemon { private static final Logger LOG = LogManager.getLogger(CloudTabletRebalancer.class); + private final CloudTabletRebalancerMetrics rebalancerMetrics; + private long currentRoundTabletScanCount; + private volatile ConcurrentHashMap> beToTabletsGlobal = new ConcurrentHashMap>(); @@ -243,8 +246,14 @@ private boolean isComputeGroupBalanceChanged(String clusterId) { } public CloudTabletRebalancer(CloudSystemInfoService cloudSystemInfoService) { + this(cloudSystemInfoService, CloudTabletRebalancerMetrics.create()); + } + + CloudTabletRebalancer(CloudSystemInfoService cloudSystemInfoService, + CloudTabletRebalancerMetrics rebalancerMetrics) { super("cloud tablet rebalancer", Config.cloud_tablet_rebalancer_interval_second * 1000); this.cloudSystemInfoService = cloudSystemInfoService; + this.rebalancerMetrics = rebalancerMetrics; } private void initializeWarmupExecutorsIfNeeded() { @@ -503,43 +512,49 @@ protected void runAfterCatalogReady() { } LOG.info("cloud tablet rebalance begin"); - long start = System.currentTimeMillis(); - activeTabletIds = getActiveTabletIds(); - globalBalanceTypeEnum = BalanceTypeEnum.getCloudWarmUpForRebalanceTypeEnum(); + CloudTabletRebalancerMetrics.Round metricRound = rebalancerMetrics.startRound(); + currentRoundTabletScanCount = 0L; + try { + long start = System.currentTimeMillis(); + activeTabletIds = getActiveTabletIds(); + globalBalanceTypeEnum = BalanceTypeEnum.getCloudWarmUpForRebalanceTypeEnum(); - buildClusterToBackendMap(); - if (!completeRouteInfo()) { - return; - } + buildClusterToBackendMap(); + if (!completeRouteInfo()) { + return; + } - statRouteInfo(); - migrateTabletsForSmoothUpgrade(); - statRouteInfo(); + statRouteInfo(); + migrateTabletsForSmoothUpgrade(); + statRouteInfo(); - indexBalanced = true; - tableBalanced = true; + indexBalanced = true; + tableBalanced = true; - performBalancing(); + performBalancing(); - checkDecommissionState(clusterToBes); - inited = true; - long sleepSeconds = Config.cloud_tablet_rebalancer_interval_second; - if (sleepSeconds < 0L) { - LOG.warn("cloud tablet rebalance interval second is negative, change it to default 1s"); - sleepSeconds = 1L; - } - long balanceEnd = System.currentTimeMillis(); - if (DebugPointUtil.isEnable("CloudTabletRebalancer.balanceEnd.tooLong")) { - LOG.info("debug pointCloudTabletRebalancer.balanceEnd.tooLong"); - // slower the balance end time to trigger next balance immediately - balanceEnd += (Config.cloud_tablet_rebalancer_interval_second + 10L) * 1000L; - } - if (balanceEnd - start > Config.cloud_tablet_rebalancer_interval_second * 1000L) { - sleepSeconds = 1L; + checkDecommissionState(clusterToBes); + inited = true; + long sleepSeconds = Config.cloud_tablet_rebalancer_interval_second; + if (sleepSeconds < 0L) { + LOG.warn("cloud tablet rebalance interval second is negative, change it to default 1s"); + sleepSeconds = 1L; + } + long balanceEnd = System.currentTimeMillis(); + if (DebugPointUtil.isEnable("CloudTabletRebalancer.balanceEnd.tooLong")) { + LOG.info("debug pointCloudTabletRebalancer.balanceEnd.tooLong"); + // slower the balance end time to trigger next balance immediately + balanceEnd += (Config.cloud_tablet_rebalancer_interval_second + 10L) * 1000L; + } + if (balanceEnd - start > Config.cloud_tablet_rebalancer_interval_second * 1000L) { + sleepSeconds = 1L; + } + setInterval(sleepSeconds * 1000L); + LOG.info("finished to rebalancer. cost: {} ms, rebalancer sche interval {} s", + (System.currentTimeMillis() - start), sleepSeconds); + } finally { + rebalancerMetrics.finishRound(metricRound, currentRoundTabletScanCount); } - setInterval(sleepSeconds * 1000L); - LOG.info("finished to rebalancer. cost: {} ms, rebalancer sche interval {} s", - (System.currentTimeMillis() - start), sleepSeconds); } private void buildClusterToBackendMap() { @@ -1241,6 +1256,7 @@ public void loopCloudReplica(Operator operator) { for (MaterializedIndex index : partition.getMaterializedIndices(IndexExtState.VISIBLE)) { for (Map.Entry> entry : clusterToBes.entrySet()) { String cluster = entry.getKey(); + currentRoundTabletScanCount += index.getTablets().size(); operator.op(db, table, partition, index, cluster); } } // end for indices diff --git a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTabletRebalancerMetrics.java b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTabletRebalancerMetrics.java new file mode 100644 index 00000000000000..96d5b31ad1a0b7 --- /dev/null +++ b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTabletRebalancerMetrics.java @@ -0,0 +1,78 @@ +// 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.doris.cloud.catalog; + +import org.apache.doris.metric.MetricRepo; + +import java.lang.management.ManagementFactory; +import java.util.concurrent.TimeUnit; +import java.util.function.LongSupplier; + +final class CloudTabletRebalancerMetrics { + private static final long ALLOCATED_BYTES_UNAVAILABLE = -1L; + + private final LongSupplier nanoTimeSupplier; + private final LongSupplier allocatedBytesSupplier; + + CloudTabletRebalancerMetrics(LongSupplier nanoTimeSupplier, LongSupplier allocatedBytesSupplier) { + this.nanoTimeSupplier = nanoTimeSupplier; + this.allocatedBytesSupplier = allocatedBytesSupplier; + } + + static CloudTabletRebalancerMetrics create() { + com.sun.management.ThreadMXBean threadMxBean = + ManagementFactory.getPlatformMXBean(com.sun.management.ThreadMXBean.class); + return new CloudTabletRebalancerMetrics(System::nanoTime, createAllocatedBytesSupplier(threadMxBean)); + } + + Round startRound() { + return new Round(nanoTimeSupplier.getAsLong(), allocatedBytesSupplier.getAsLong()); + } + + void finishRound(Round round, long tabletScanCount) { + long durationMs = TimeUnit.NANOSECONDS.toMillis(nanoTimeSupplier.getAsLong() - round.startNanos); + long currentAllocatedBytes = allocatedBytesSupplier.getAsLong(); + long allocatedBytes = round.startAllocatedBytes < 0L || currentAllocatedBytes < 0L + ? ALLOCATED_BYTES_UNAVAILABLE : currentAllocatedBytes - round.startAllocatedBytes; + MetricRepo.updateCloudTabletRebalancerMetrics(durationMs, allocatedBytes, tabletScanCount); + } + + static LongSupplier createAllocatedBytesSupplier(com.sun.management.ThreadMXBean threadMxBean) { + if (threadMxBean == null || !threadMxBean.isThreadAllocatedMemorySupported()) { + return () -> ALLOCATED_BYTES_UNAVAILABLE; + } + if (!threadMxBean.isThreadAllocatedMemoryEnabled()) { + try { + threadMxBean.setThreadAllocatedMemoryEnabled(true); + } catch (SecurityException | UnsupportedOperationException e) { + return () -> ALLOCATED_BYTES_UNAVAILABLE; + } + } + return () -> threadMxBean.getThreadAllocatedBytes(Thread.currentThread().getId()); + } + + static final class Round { + private final long startNanos; + private final long startAllocatedBytes; + + private Round(long startNanos, long startAllocatedBytes) { + this.startNanos = startNanos; + this.startAllocatedBytes = startAllocatedBytes; + } + } +} diff --git a/fe/fe-core/src/main/java/org/apache/doris/metric/CloudMetrics.java b/fe/fe-core/src/main/java/org/apache/doris/metric/CloudMetrics.java index e74ed0b1bc31b8..af8c17bef6ed59 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/metric/CloudMetrics.java +++ b/fe/fe-core/src/main/java/org/apache/doris/metric/CloudMetrics.java @@ -71,6 +71,13 @@ public class CloudMetrics { protected static AutoMappedMetric CLUSTER_CLOUD_WARM_UP_CACHE_BALANCE_NUM; protected static AutoMappedMetric VIRTUAL_COMPUTE_GROUP_SWITCH_COUNTER; + protected static LongCounterMetric CLOUD_TABLET_REBALANCER_ROUND_TOTAL; + protected static LongCounterMetric CLOUD_TABLET_REBALANCER_ALLOCATED_BYTES_TOTAL; + protected static GaugeMetricImpl CLOUD_TABLET_REBALANCER_LAST_ROUND_ALLOCATED_BYTES; + protected static LongCounterMetric CLOUD_TABLET_REBALANCER_DURATION_MS_TOTAL; + protected static GaugeMetricImpl CLOUD_TABLET_REBALANCER_LAST_ROUND_DURATION_MS; + protected static LongCounterMetric CLOUD_TABLET_REBALANCER_TABLET_SCAN_TOTAL; + protected static void init() { if (Config.isNotCloudMode()) { return; @@ -201,5 +208,35 @@ protected static void init() { VIRTUAL_COMPUTE_GROUP_SWITCH_COUNTER = new AutoMappedMetric<>(name -> new LongCounterMetric( "virtual_compute_group_switch_total", MetricUnit.NOUNIT, "virtual compute group active standby switch count")); + + initCloudTabletRebalancerMetrics(); + } + + static void initCloudTabletRebalancerMetrics() { + CLOUD_TABLET_REBALANCER_ROUND_TOTAL = new LongCounterMetric( + "cloud_tablet_rebalancer_round_total", MetricUnit.OPERATIONS, + "total cloud tablet rebalancer rounds"); + CLOUD_TABLET_REBALANCER_ALLOCATED_BYTES_TOTAL = new LongCounterMetric( + "cloud_tablet_rebalancer_allocated_bytes_total", MetricUnit.BYTES, + "total bytes allocated by cloud tablet rebalancer rounds"); + CLOUD_TABLET_REBALANCER_LAST_ROUND_ALLOCATED_BYTES = new GaugeMetricImpl<>( + "cloud_tablet_rebalancer_last_round_allocated_bytes", MetricUnit.BYTES, + "bytes allocated by the last cloud tablet rebalancer round, or -1 when unavailable", -1L); + CLOUD_TABLET_REBALANCER_DURATION_MS_TOTAL = new LongCounterMetric( + "cloud_tablet_rebalancer_duration_ms_total", MetricUnit.MILLISECONDS, + "total cloud tablet rebalancer round duration in milliseconds"); + CLOUD_TABLET_REBALANCER_LAST_ROUND_DURATION_MS = new GaugeMetricImpl<>( + "cloud_tablet_rebalancer_last_round_duration_ms", MetricUnit.MILLISECONDS, + "duration of the last cloud tablet rebalancer round in milliseconds", 0L); + CLOUD_TABLET_REBALANCER_TABLET_SCAN_TOTAL = new LongCounterMetric( + "cloud_tablet_rebalancer_tablet_scan_total", MetricUnit.OPERATIONS, + "total tablet route entries scanned by cloud tablet rebalancer rounds"); + + MetricRepo.DORIS_METRIC_REGISTER.addMetrics(CLOUD_TABLET_REBALANCER_ROUND_TOTAL); + MetricRepo.DORIS_METRIC_REGISTER.addMetrics(CLOUD_TABLET_REBALANCER_ALLOCATED_BYTES_TOTAL); + MetricRepo.DORIS_METRIC_REGISTER.addMetrics(CLOUD_TABLET_REBALANCER_LAST_ROUND_ALLOCATED_BYTES); + MetricRepo.DORIS_METRIC_REGISTER.addMetrics(CLOUD_TABLET_REBALANCER_DURATION_MS_TOTAL); + MetricRepo.DORIS_METRIC_REGISTER.addMetrics(CLOUD_TABLET_REBALANCER_LAST_ROUND_DURATION_MS); + MetricRepo.DORIS_METRIC_REGISTER.addMetrics(CLOUD_TABLET_REBALANCER_TABLET_SCAN_TOTAL); } } diff --git a/fe/fe-core/src/main/java/org/apache/doris/metric/MetricRepo.java b/fe/fe-core/src/main/java/org/apache/doris/metric/MetricRepo.java index 1e7a979b01ff48..97ac8d2446095a 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/metric/MetricRepo.java +++ b/fe/fe-core/src/main/java/org/apache/doris/metric/MetricRepo.java @@ -2370,4 +2370,19 @@ public static void updateClusterCloudBalanceNum(String clusterName, String clust counter.setLabels(labels); MetricRepo.DORIS_METRIC_REGISTER.addMetrics(counter); } + + public static void updateCloudTabletRebalancerMetrics(long durationMs, long allocatedBytes, + long tabletScanCount) { + if (!MetricRepo.isInit || Config.isNotCloudMode()) { + return; + } + CloudMetrics.CLOUD_TABLET_REBALANCER_ROUND_TOTAL.increase(1L); + CloudMetrics.CLOUD_TABLET_REBALANCER_DURATION_MS_TOTAL.increase(durationMs); + CloudMetrics.CLOUD_TABLET_REBALANCER_LAST_ROUND_DURATION_MS.setValue(durationMs); + CloudMetrics.CLOUD_TABLET_REBALANCER_TABLET_SCAN_TOTAL.increase(tabletScanCount); + CloudMetrics.CLOUD_TABLET_REBALANCER_LAST_ROUND_ALLOCATED_BYTES.setValue(allocatedBytes); + if (allocatedBytes >= 0L) { + CloudMetrics.CLOUD_TABLET_REBALANCER_ALLOCATED_BYTES_TOTAL.increase(allocatedBytes); + } + } }