From 27fce5701e2f1c56351b1e76119ab5b7fffc58ff Mon Sep 17 00:00:00 2001 From: meiyi Date: Mon, 12 Jan 2026 16:23:53 +0800 Subject: [PATCH] fix tablet and replica --- .../org/apache/doris/catalog/LocalReplica.java | 8 ++++---- .../org/apache/doris/catalog/LocalTablet.java | 13 +++++++++++++ .../java/org/apache/doris/catalog/Replica.java | 8 ++++---- .../java/org/apache/doris/catalog/Tablet.java | 10 +++++++++- .../apache/doris/cloud/catalog/CloudReplica.java | 16 ++++++++++++++-- .../doris/consistency/ConsistencyChecker.java | 10 ++++++---- 6 files changed, 50 insertions(+), 15 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/catalog/LocalReplica.java b/fe/fe-core/src/main/java/org/apache/doris/catalog/LocalReplica.java index 57b980152dba12..3cf0a35b74a123 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/catalog/LocalReplica.java +++ b/fe/fe-core/src/main/java/org/apache/doris/catalog/LocalReplica.java @@ -34,9 +34,9 @@ public class LocalReplica extends Replica { @SerializedName(value = "rds", alternate = {"remoteDataSize"}) private volatile long remoteDataSize = 0; @SerializedName(value = "ris", alternate = {"remoteInvertedIndexSize"}) - private Long remoteInvertedIndexSize = 0L; + private long remoteInvertedIndexSize = 0L; @SerializedName(value = "rss", alternate = {"remoteSegmentSize"}) - private Long remoteSegmentSize = 0L; + private long remoteSegmentSize = 0L; // the last load failed version @SerializedName(value = "lfv", alternate = {"lastFailedVersion"}) @@ -207,7 +207,7 @@ public void setRemoteDataSize(long remoteDataSize) { } @Override - public Long getRemoteInvertedIndexSize() { + public long getRemoteInvertedIndexSize() { return remoteInvertedIndexSize; } @@ -217,7 +217,7 @@ public void setRemoteInvertedIndexSize(long remoteInvertedIndexSize) { } @Override - public Long getRemoteSegmentSize() { + public long getRemoteSegmentSize() { return remoteSegmentSize; } diff --git a/fe/fe-core/src/main/java/org/apache/doris/catalog/LocalTablet.java b/fe/fe-core/src/main/java/org/apache/doris/catalog/LocalTablet.java index 0531d499c56209..924c425aa07274 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/catalog/LocalTablet.java +++ b/fe/fe-core/src/main/java/org/apache/doris/catalog/LocalTablet.java @@ -28,6 +28,9 @@ public class LocalTablet extends Tablet { private static final Logger LOG = LogManager.getLogger(LocalTablet.class); + @SerializedName(value = "lastCheckTime") + private long lastCheckTime; + // cooldown conf @SerializedName(value = "cri", alternate = {"cooldownReplicaId"}) private long cooldownReplicaId = -1; @@ -148,4 +151,14 @@ protected long getLastTimeNoPathForNewReplica() { public void setLastTimeNoPathForNewReplica(long lastTimeNoPathForNewReplica) { this.lastTimeNoPathForNewReplica = lastTimeNoPathForNewReplica; } + + @Override + public long getLastCheckTime() { + return lastCheckTime; + } + + @Override + public void setLastCheckTime(long lastCheckTime) { + this.lastCheckTime = lastCheckTime; + } } diff --git a/fe/fe-core/src/main/java/org/apache/doris/catalog/Replica.java b/fe/fe-core/src/main/java/org/apache/doris/catalog/Replica.java index e1aa322b2832ef..b6ac45f04840b4 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/catalog/Replica.java +++ b/fe/fe-core/src/main/java/org/apache/doris/catalog/Replica.java @@ -99,11 +99,11 @@ public static class ReplicaContext { @Setter @Getter @SerializedName(value = "lis", alternate = {"localInvertedIndexSize"}) - private Long localInvertedIndexSize = 0L; + private long localInvertedIndexSize = 0L; @Setter @Getter @SerializedName(value = "lss", alternate = {"localSegmentSize"}) - private Long localSegmentSize = 0L; + private long localSegmentSize = 0L; public Replica() { } @@ -199,7 +199,7 @@ public void setRemoteDataSize(long remoteDataSize) { } } - public Long getRemoteInvertedIndexSize() { + public long getRemoteInvertedIndexSize() { return 0L; } @@ -209,7 +209,7 @@ public void setRemoteInvertedIndexSize(long remoteInvertedIndexSize) { } } - public Long getRemoteSegmentSize() { + public long getRemoteSegmentSize() { return 0L; } diff --git a/fe/fe-core/src/main/java/org/apache/doris/catalog/Tablet.java b/fe/fe-core/src/main/java/org/apache/doris/catalog/Tablet.java index 7386cae8a27097..96ae1f9b692763 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/catalog/Tablet.java +++ b/fe/fe-core/src/main/java/org/apache/doris/catalog/Tablet.java @@ -51,7 +51,7 @@ /** * This class represents the olap tablet related metadata. */ -public class Tablet extends MetaObject { +public class Tablet { private static final Logger LOG = LogManager.getLogger(Tablet.class); // if current version count of replica is mor than // QUERYABLE_TIMES_OF_MIN_VERSION_COUNT times the minimum version count, @@ -889,4 +889,12 @@ public void setLastTimeNoPathForNewReplica(long lastTimeNoPathForNewReplica) { throw new UnsupportedOperationException("setLastTimeNoPathForNewReplica is not supported in Tablet"); } } + + public long getLastCheckTime() { + return -1; + } + + public void setLastCheckTime(long lastCheckTime) { + throw new UnsupportedOperationException("setLastCheckTime is not supported in Tablet"); + } } diff --git a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudReplica.java b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudReplica.java index c703d4af62a15c..8f2c6127e25f2c 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudReplica.java +++ b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudReplica.java @@ -72,7 +72,7 @@ public class CloudReplica extends Replica { private static final Random rand = new Random(); - private Map> memClusterToBackends = new ConcurrentHashMap>(); + private Map> memClusterToBackends = null; // clusterId, secondaryBe, changeTimestamp private Map> secondaryClusterToBackends @@ -311,6 +311,7 @@ private long getBackendIdImpl(String clusterId) throws ComputeGroupException { int indexRand = rand.nextInt(Config.cloud_replica_num); int coldReadRand = rand.nextInt(100); boolean allowColdRead = coldReadRand < Config.cloud_cold_read_percent; + initMemClusterToBackends(); boolean replicaEnough = memClusterToBackends.get(clusterId) != null && memClusterToBackends.get(clusterId).size() > indexRand; @@ -470,7 +471,18 @@ private long getIndexByBeNum(long hashValue, int beNum) { return (hashValue % beNum + beNum) % beNum; } - public List hashReplicaToBes(String clusterId, boolean isBackGround, int replicaNum) + private void initMemClusterToBackends() { + // the enable_cloud_multi_replica is not used now + if (memClusterToBackends == null) { + synchronized (this) { + if (memClusterToBackends == null) { + memClusterToBackends = new ConcurrentHashMap<>(); + } + } + } + } + + private List hashReplicaToBes(String clusterId, boolean isBackGround, int replicaNum) throws ComputeGroupException { // TODO(luwei) list should be sorted List clusterBes = ((CloudSystemInfoService) Env.getCurrentSystemInfo()) diff --git a/fe/fe-core/src/main/java/org/apache/doris/consistency/ConsistencyChecker.java b/fe/fe-core/src/main/java/org/apache/doris/consistency/ConsistencyChecker.java index 471e235684e390..4a19692154a513 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/consistency/ConsistencyChecker.java +++ b/fe/fe-core/src/main/java/org/apache/doris/consistency/ConsistencyChecker.java @@ -56,6 +56,8 @@ public class ConsistencyChecker extends MasterDaemon { private static final Comparator COMPARATOR = (first, second) -> Long.signum(first.getLastCheckTime() - second.getLastCheckTime()); + private static final Comparator TABLET_COMPARATOR = + (first, second) -> Long.signum(first.getLastCheckTime() - second.getLastCheckTime()); // tabletId -> job private Map jobs; @@ -315,12 +317,12 @@ private List chooseTablets() { MaterializedIndex index = (MaterializedIndex) chosenOne; // sort tablets - Queue tabletQueue - = new PriorityQueue<>(Math.max(index.getTablets().size(), 1), COMPARATOR); + Queue tabletQueue = new PriorityQueue<>(Math.max(index.getTablets().size(), 1), + TABLET_COMPARATOR); tabletQueue.addAll(index.getTablets()); + Tablet tablet = null; - while ((chosenOne = tabletQueue.poll()) != null) { - Tablet tablet = (Tablet) chosenOne; + while ((tablet = tabletQueue.poll()) != null) { long chosenTabletId = tablet.getId(); if (this.jobs.containsKey(chosenTabletId)) {