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 ab32596f8633bf..14d9171f3ff0e4 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 @@ -17,6 +17,7 @@ package org.apache.doris.catalog; +import org.apache.doris.catalog.Replica.ReplicaState; import org.apache.doris.clone.TabletSchedCtx; import org.apache.doris.clone.TabletSchedCtx.Priority; import org.apache.doris.common.Config; @@ -33,6 +34,7 @@ import java.util.Comparator; import java.util.Iterator; import java.util.List; +import java.util.stream.LongStream; public class LocalTablet extends Tablet { private static final Logger LOG = LogManager.getLogger(LocalTablet.class); @@ -116,6 +118,51 @@ public long getRemoteDataSize() { return replicas.stream().max(Comparator.comparing(Replica::getRemoteDataSize)).get().getRemoteDataSize(); } + @Override + public Replica getReplicaById(long replicaId) { + for (Replica replica : getReplicas()) { + if (replica.getId() == replicaId) { + return replica; + } + } + return null; + } + + // ATTN: Replica::getDataSize may zero in cloud and non-cloud + // due to dataSize not write to image + @Override + public long getDataSize(boolean singleReplica, boolean filterSizeZero) { + LongStream s = getReplicas().stream().filter(r -> r.getState() == ReplicaState.NORMAL) + .filter(r -> !filterSizeZero || r.getDataSize() > 0) + .mapToLong(Replica::getDataSize); + return singleReplica ? Double.valueOf(s.average().orElse(0)).longValue() : s.sum(); + } + + @Override + public long getRowCount(boolean singleReplica) { + LongStream s = getReplicas().stream().filter(r -> r.getState() == ReplicaState.NORMAL) + .mapToLong(Replica::getRowCount); + return singleReplica ? Double.valueOf(s.average().orElse(0)).longValue() : s.sum(); + } + + // Get the least row count among all valid replicas. + // The replica with the least row count is the most accurate one. Because it performs most compaction. + @Override + public long getMinReplicaRowCount(long version) { + long minRowCount = Long.MAX_VALUE; + long maxReplicaVersion = 0; + for (Replica r : getReplicas()) { + if (r.isAlive() + && r.checkVersionCatchUp(version, false) + && (r.getVersion() > maxReplicaVersion + || r.getVersion() == maxReplicaVersion && r.getRowCount() < minRowCount)) { + minRowCount = r.getRowCount(); + maxReplicaVersion = r.getVersion(); + } + } + return minRowCount == Long.MAX_VALUE ? 0 : minRowCount; + } + @Override public long getCheckedVersion() { return this.checkedVersion; 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 81080a1754ae18..50b16f379feaae 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 @@ -44,7 +44,6 @@ import java.util.Map; import java.util.Set; import java.util.stream.Collectors; -import java.util.stream.LongStream; /** * This class represents the olap tablet related metadata. @@ -302,14 +301,7 @@ public String getDetailsStatusForQuery(long visibleVersion) { return sb.toString(); } - public Replica getReplicaById(long replicaId) { - for (Replica replica : getReplicas()) { - if (replica.getId() == replicaId) { - return replica; - } - } - return null; - } + public abstract Replica getReplicaById(long replicaId); public abstract Replica getReplicaByBackendId(long backendId); @@ -333,41 +325,17 @@ public String toString() { @Override public abstract boolean equals(Object obj); - // ATTN: Replica::getDataSize may zero in cloud and non-cloud - // due to dataSize not write to image - public long getDataSize(boolean singleReplica, boolean filterSizeZero) { - LongStream s = getReplicas().stream().filter(r -> r.getState() == ReplicaState.NORMAL) - .filter(r -> !filterSizeZero || r.getDataSize() > 0) - .mapToLong(Replica::getDataSize); - return singleReplica ? Double.valueOf(s.average().orElse(0)).longValue() : s.sum(); - } + public abstract long getDataSize(boolean singleReplica, boolean filterSizeZero); public long getRemoteDataSize() { return 0; } - public long getRowCount(boolean singleReplica) { - LongStream s = getReplicas().stream().filter(r -> r.getState() == ReplicaState.NORMAL) - .mapToLong(Replica::getRowCount); - return singleReplica ? Double.valueOf(s.average().orElse(0)).longValue() : s.sum(); - } + public abstract long getRowCount(boolean singleReplica); // Get the least row count among all valid replicas. // The replica with the least row count is the most accurate one. Because it performs most compaction. - public long getMinReplicaRowCount(long version) { - long minRowCount = Long.MAX_VALUE; - long maxReplicaVersion = 0; - for (Replica r : getReplicas()) { - if (r.isAlive() - && r.checkVersionCatchUp(version, false) - && (r.getVersion() > maxReplicaVersion - || r.getVersion() == maxReplicaVersion && r.getRowCount() < minRowCount)) { - minRowCount = r.getRowCount(); - maxReplicaVersion = r.getVersion(); - } - } - return minRowCount == Long.MAX_VALUE ? 0 : minRowCount; - } + public abstract long getMinReplicaRowCount(long version); /** * A replica is healthy only if diff --git a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTablet.java b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTablet.java index 3aaa20bb372876..5fdc53a6c556c6 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTablet.java +++ b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTablet.java @@ -19,19 +19,20 @@ import org.apache.doris.catalog.Env; import org.apache.doris.catalog.Replica; +import org.apache.doris.catalog.Replica.ReplicaState; import org.apache.doris.catalog.Tablet; import org.apache.doris.common.InternalErrorCode; import org.apache.doris.common.UserException; import org.apache.doris.persist.gson.GsonPostProcessable; import org.apache.doris.system.SystemInfoService; -import com.google.common.collect.Lists; import com.google.common.collect.Multimap; import com.google.gson.annotations.SerializedName; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; import java.io.IOException; +import java.util.Collections; import java.util.List; import java.util.Objects; @@ -82,33 +83,59 @@ public Multimap getNormalReplicaBackendPathMap(String beEndpoint) th return backendPathMapReprocess(pathMap); } - private boolean isLatestReplicaAndDeleteOld(Replica newReplica) { + @Override + public void addReplica(Replica replica, boolean isRestore) { + this.replica = replica; + if (!isRestore) { + Env.getCurrentInvertedIndex().addReplica(id, replica); + } + } + + @Override + public List getReplicas() { if (replica == null) { - return true; + return Collections.emptyList(); } - if (replica.getVersion() <= newReplica.getVersion()) { - replica = null; - return true; + return Collections.singletonList(replica); + } + + @Override + public Replica getReplicaById(long replicaId) { + if (replica != null && replica.getId() == replicaId) { + return replica; } - return false; + return null; + } + + public CloudReplica getCloudReplica() { + if (replica == null) { + return null; + } + return (CloudReplica) replica; } @Override - public void addReplica(Replica replica, boolean isRestore) { - if (isLatestReplicaAndDeleteOld(replica)) { - this.replica = replica; - if (!isRestore) { - Env.getCurrentInvertedIndex().addReplica(id, replica); - } + public long getDataSize(boolean singleReplica, boolean filterSizeZero) { + if (replica != null && replica.getState() == ReplicaState.NORMAL) { + return replica.getDataSize(); } + return 0; } @Override - public List getReplicas() { - if (replica == null) { - return Lists.newArrayList(); + public long getRowCount(boolean singleReplica) { + if (replica != null && replica.getState() == ReplicaState.NORMAL) { + return replica.getRowCount(); + } + return 0; + } + + @Override + public long getMinReplicaRowCount(long version) { + if (replica != null && replica.isAlive() && replica.checkVersionCatchUp(version, false)) { + return replica.getRowCount(); } - return Lists.newArrayList(replica); + return 0; } @Override diff --git a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTabletInvertedIndex.java b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTabletInvertedIndex.java index c0ebb3c134d4e3..af6a368f3f9d57 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTabletInvertedIndex.java +++ b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTabletInvertedIndex.java @@ -21,11 +21,11 @@ import org.apache.doris.catalog.TabletInvertedIndex; import com.google.common.base.Preconditions; -import com.google.common.collect.Lists; import com.google.common.collect.Maps; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; +import java.util.Collections; import java.util.List; import java.util.Map; @@ -45,9 +45,9 @@ public List getReplicas(Long tabletId) { long stamp = readLock(); try { if (replicaMetaMap.containsKey(tabletId)) { - return Lists.newArrayList(replicaMetaMap.get(tabletId)); + return Collections.singletonList(replicaMetaMap.get(tabletId)); } - return Lists.newArrayList(); + return Collections.emptyList(); } finally { readUnlock(stamp); } @@ -118,9 +118,9 @@ public List getReplicasByTabletId(long tabletId) { long stamp = readLock(); try { if (replicaMetaMap.containsKey(tabletId)) { - return Lists.newArrayList(replicaMetaMap.get(tabletId)); + return Collections.singletonList(replicaMetaMap.get(tabletId)); } - return Lists.newArrayList(); + return Collections.emptyList(); } finally { readUnlock(stamp); } 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 20264eba7d4fd2..80d29e325ec378 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 @@ -1243,7 +1243,7 @@ private void updateBeToTablets(Tablet pickedTablet, long srcBe, long destBe, ConcurrentHashMap>> beToTabletsInTable, ConcurrentHashMap>>> partToTablets) { - CloudReplica replica = (CloudReplica) pickedTablet.getReplicas().get(0); + CloudReplica replica = ((CloudTablet) pickedTablet).getCloudReplica(); long tableId = replica.getTableId(); long partId = replica.getPartitionId(); long indexId = replica.getIndexId(); @@ -1258,7 +1258,7 @@ private void updateBeToTablets(Tablet pickedTablet, long srcBe, long destBe, private void updateClusterToBeMap(Tablet pickedTablet, long destBe, String clusterId, List infos) { - CloudReplica cloudReplica = (CloudReplica) pickedTablet.getReplicas().get(0); + CloudReplica cloudReplica = ((CloudTablet) pickedTablet).getCloudReplica(); Database db = Env.getCurrentInternalCatalog().getDbNullable(cloudReplica.getDbId()); if (db == null) { return; @@ -1491,7 +1491,7 @@ private void balanceImpl(List bes, String clusterId, Map continue; // No tablet to pick } - CloudReplica cloudReplica = (CloudReplica) pickedTablet.getReplicas().get(0); + CloudReplica cloudReplica = ((CloudTablet) pickedTablet).getCloudReplica(); Backend srcBackend = Env.getCurrentSystemInfo().getBackend(srcBe); if ((BalanceTypeEnum.WITHOUT_WARMUP.equals(currentBalanceType) @@ -1616,7 +1616,7 @@ private void migrateTablets(Long srcBe, Long dstBe) { List infos = new ArrayList<>(); for (Tablet tablet : tablets) { // get replica - CloudReplica cloudReplica = (CloudReplica) tablet.getReplicas().get(0); + CloudReplica cloudReplica = ((CloudTablet) tablet).getCloudReplica(); Backend be = cloudSystemInfoService.getBackend(srcBe); if (be == null) { LOG.info("src backend {} not found", srcBe); diff --git a/fe/fe-core/src/main/java/org/apache/doris/cloud/datasource/CloudInternalCatalog.java b/fe/fe-core/src/main/java/org/apache/doris/cloud/datasource/CloudInternalCatalog.java index 3ffde0f00c4ba3..043a10258ea670 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/cloud/datasource/CloudInternalCatalog.java +++ b/fe/fe-core/src/main/java/org/apache/doris/cloud/datasource/CloudInternalCatalog.java @@ -42,6 +42,7 @@ import org.apache.doris.cloud.catalog.CloudEnv; import org.apache.doris.cloud.catalog.CloudPartition; import org.apache.doris.cloud.catalog.CloudReplica; +import org.apache.doris.cloud.catalog.CloudTablet; import org.apache.doris.cloud.persist.UpdateCloudReplicaInfo; import org.apache.doris.cloud.proto.Cloud; import org.apache.doris.cloud.proto.Cloud.CopyJobPB; @@ -849,7 +850,7 @@ public void erasePartitionDropBackendReplicas(List partitions) { for (MaterializedIndex index : partition.getMaterializedIndices(IndexExtState.ALL)) { indexIds.add(index.getId()); if (tableId == -1) { - tableId = ((CloudReplica) index.getTablets().get(0).getReplicas().get(0)).getTableId(); + tableId = ((CloudTablet) index.getTablets().get(0)).getCloudReplica().getTableId(); } } partitionIds.add(partition.getId()); @@ -1135,7 +1136,7 @@ private void unprotectUpdateCloudReplica(OlapTable olapTable, UpdateCloudReplica Tablet tablet = materializedIndex.getTablet(tabletIds.get(i)); Replica replica; if (info.getReplicaIds().isEmpty()) { - replica = tablet.getReplicas().get(0); + replica = ((CloudTablet) tablet).getCloudReplica(); } else { replica = tablet.getReplicaById(info.getReplicaIds().get(i)); } diff --git a/fe/fe-core/src/test/java/org/apache/doris/statistics/OlapAnalysisTaskTest.java b/fe/fe-core/src/test/java/org/apache/doris/statistics/OlapAnalysisTaskTest.java index fa1879db5f3bfc..89a0a67d810f95 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/statistics/OlapAnalysisTaskTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/statistics/OlapAnalysisTaskTest.java @@ -23,6 +23,7 @@ import org.apache.doris.catalog.DataProperty; import org.apache.doris.catalog.DatabaseIf; import org.apache.doris.catalog.KeysType; +import org.apache.doris.catalog.LocalTablet; import org.apache.doris.catalog.MaterializedIndex; import org.apache.doris.catalog.OlapTable; import org.apache.doris.catalog.Partition; @@ -652,7 +653,7 @@ public int getPartitionSampleCount() { } @Test - public void testGetSampleTablets(@Mocked MaterializedIndex index, @Mocked Tablet t) { + public void testGetSampleTablets(@Mocked MaterializedIndex index, @Mocked LocalTablet t) { OlapAnalysisTask task = new OlapAnalysisTask(); task.tbl = new OlapTable(); task.col = new Column("col1", PrimitiveType.STRING); @@ -704,7 +705,7 @@ public Tablet getTablet(long tabletId) { return t; } }; - new MockUp() { + new MockUp() { @Mock public long getMinReplicaRowCount(long version) { return tabletsRowCount[i[0]++];