diff --git a/fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java b/fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java index 5685dcb06f74d2..f0f9c1d20f04a6 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java +++ b/fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java @@ -33,6 +33,7 @@ import org.apache.doris.mtmv.MTMVCache; import org.apache.doris.mtmv.MTMVJobInfo; import org.apache.doris.mtmv.MTMVJobManager; +import org.apache.doris.mtmv.MTMVPartitionExpander; import org.apache.doris.mtmv.MTMVPartitionInfo; import org.apache.doris.mtmv.MTMVPartitionInfo.MTMVPartitionType; import org.apache.doris.mtmv.MTMVPartitionUtil; @@ -381,7 +382,8 @@ public Map generateMvPartitionDescs() throws AnalysisE * @return mvPartitionName ==> relationPartitionNames and relationPartitionName ==> mvPartitionName * @throws AnalysisException */ - public Pair>, Map> calculateDoublyPartitionMappings() + public Pair>, Map> calculateDoublyPartitionMappings( + Set queryUsedBaseTablePartitions) throws AnalysisException { if (mvPartitionInfo.getPartitionType() == MTMVPartitionType.SELF_MANAGE) { return Pair.of(Maps.newHashMap(), Maps.newHashMap()); @@ -389,9 +391,11 @@ public Pair>, Map> calculateDoublyPartit long start = System.currentTimeMillis(); Map> mvToBase = Maps.newHashMap(); Map baseToMv = Maps.newHashMap(); - Map> relatedPartitionDescs = MTMVPartitionUtil - .generateRelatedPartitionDescs(mvPartitionInfo, mvProperties); Map mvPartitionItems = getAndCopyPartitionItems(); + Set effectiveFilter = getEffectiveQueryUsedBaseTablePartitions( + queryUsedBaseTablePartitions, mvPartitionItems); + Map> relatedPartitionDescs = MTMVPartitionUtil + .generateRelatedPartitionDescs(mvPartitionInfo, mvProperties, effectiveFilter); for (Entry entry : mvPartitionItems.entrySet()) { Set basePartitionNames = relatedPartitionDescs.getOrDefault(entry.getValue().toPartitionKeyDesc(), Sets.newHashSet()); @@ -417,14 +421,21 @@ public Pair>, Map> calculateDoublyPartit * @throws AnalysisException */ public Map> calculatePartitionMappings() throws AnalysisException { + return calculatePartitionMappings(null); + } + + public Map> calculatePartitionMappings(Set queryUsedBaseTablePartitions) + throws AnalysisException { if (mvPartitionInfo.getPartitionType() == MTMVPartitionType.SELF_MANAGE) { return Maps.newHashMap(); } long start = System.currentTimeMillis(); + Map mvPartitionItems = getAndCopyPartitionItems(); + Set effectiveFilter = getEffectiveQueryUsedBaseTablePartitions( + queryUsedBaseTablePartitions, mvPartitionItems); Map> res = Maps.newHashMap(); Map> relatedPartitionDescs = MTMVPartitionUtil - .generateRelatedPartitionDescs(mvPartitionInfo, mvProperties); - Map mvPartitionItems = getAndCopyPartitionItems(); + .generateRelatedPartitionDescs(mvPartitionInfo, mvProperties, effectiveFilter); for (Entry entry : mvPartitionItems.entrySet()) { res.put(entry.getKey(), relatedPartitionDescs.getOrDefault(entry.getValue().toPartitionKeyDesc(), Sets.newHashSet())); @@ -436,6 +447,17 @@ public Map> calculatePartitionMappings() throws AnalysisExce return res; } + private Set getEffectiveQueryUsedBaseTablePartitions( + Set queryUsedBaseTablePartitions, Map mvPartitionItems) + throws AnalysisException { + if (queryUsedBaseTablePartitions == null + || mvPartitionInfo.getPartitionType() != MTMVPartitionType.EXPR) { + return queryUsedBaseTablePartitions; + } + return MTMVPartitionExpander.expandToMvPartitionGranularity(queryUsedBaseTablePartitions, + mvPartitionItems, mvPartitionInfo.getRelatedTable()); + } + public ConcurrentLinkedQueue getHistoryTasks() { return jobInfo.getHistoryTasks(); } diff --git a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVPartitionExpander.java b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVPartitionExpander.java new file mode 100644 index 00000000000000..2978843f2b2b00 --- /dev/null +++ b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVPartitionExpander.java @@ -0,0 +1,91 @@ +// 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.mtmv; + +import org.apache.doris.catalog.PartitionItem; +import org.apache.doris.catalog.PartitionKey; +import org.apache.doris.catalog.PartitionType; +import org.apache.doris.catalog.RangePartitionItem; +import org.apache.doris.common.AnalysisException; +import org.apache.doris.datasource.mvcc.MvccSnapshot; +import org.apache.doris.datasource.mvcc.MvccUtil; + +import com.google.common.collect.Range; +import com.google.common.collect.Sets; + +import java.util.Map; +import java.util.Map.Entry; +import java.util.NavigableMap; +import java.util.Optional; +import java.util.Set; +import java.util.TreeMap; + +/** + * Expands query-used partitions to MV partition granularity without running + * partition expressions for every base table partition. + */ +public class MTMVPartitionExpander { + + public static Set expandToMvPartitionGranularity( + Set queryUsedBaseTablePartitions, + Map mvPartitionItems, + MTMVRelatedTableIf relatedTable) throws AnalysisException { + Optional snapshot = MvccUtil.getSnapshotFromContext(relatedTable); + if (relatedTable.getPartitionType(snapshot) != PartitionType.RANGE) { + return queryUsedBaseTablePartitions; + } + + NavigableMap> mvRanges = new TreeMap<>(); + for (PartitionItem item : mvPartitionItems.values()) { + Range range = ((RangePartitionItem) item).getItems(); + mvRanges.put(range.lowerEndpoint(), range); + } + + Map basePartitionItems = relatedTable.getAndCopyPartitionItems(snapshot); + NavigableMap> relevantMvRanges = new TreeMap<>(); + for (String queriedBasePartition : queryUsedBaseTablePartitions) { + PartitionItem baseItem = basePartitionItems.get(queriedBasePartition); + if (baseItem == null) { + continue; + } + Range baseRange = ((RangePartitionItem) baseItem).getItems(); + Range mvRange = findEnclosingRange(mvRanges, baseRange); + if (mvRange != null) { + relevantMvRanges.put(mvRange.lowerEndpoint(), mvRange); + } + } + + Set expandedPartitions = Sets.newHashSet(); + for (Entry baseEntry : basePartitionItems.entrySet()) { + Range baseRange = ((RangePartitionItem) baseEntry.getValue()).getItems(); + if (findEnclosingRange(relevantMvRanges, baseRange) != null) { + expandedPartitions.add(baseEntry.getKey()); + } + } + return expandedPartitions; + } + + private static Range findEnclosingRange( + NavigableMap> ranges, Range baseRange) { + Entry> candidate = ranges.floorEntry(baseRange.lowerEndpoint()); + return candidate != null && candidate.getValue().encloses(baseRange) ? candidate.getValue() : null; + } + + private MTMVPartitionExpander() { + } +} diff --git a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVPartitionUtil.java b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVPartitionUtil.java index ffbe55af08ceff..1d7811241b2308 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVPartitionUtil.java +++ b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVPartitionUtil.java @@ -166,10 +166,15 @@ public static List getPartitionDescsByRelatedTable( public static Map> generateRelatedPartitionDescs(MTMVPartitionInfo mvPartitionInfo, Map mvProperties) throws AnalysisException { + return generateRelatedPartitionDescs(mvPartitionInfo, mvProperties, null); + } + + public static Map> generateRelatedPartitionDescs(MTMVPartitionInfo mvPartitionInfo, + Map mvProperties, Set queryUsedPartitions) throws AnalysisException { long start = System.currentTimeMillis(); RelatedPartitionDescResult result = new RelatedPartitionDescResult(); for (MTMVRelatedPartitionDescGeneratorService service : partitionDescGenerators) { - service.apply(mvPartitionInfo, mvProperties, result); + service.apply(mvPartitionInfo, mvProperties, result, queryUsedPartitions); } if (LOG.isDebugEnabled()) { LOG.debug("generateRelatedPartitionDescs use [{}] mills, mvPartitionInfo is [{}]", @@ -581,11 +586,13 @@ public static Type getPartitionColumnType(MTMVRelatedTableIf relatedTable, Strin throw new AnalysisException("can not getPartitionColumnType by:" + col); } - public static MTMVBaseVersions getBaseVersions(MTMV mtmv) throws AnalysisException { - return new MTMVBaseVersions(getTableVersions(mtmv), getPartitionVersions(mtmv)); + public static MTMVBaseVersions getBaseVersions(MTMV mtmv, Map> partitionMappings) + throws AnalysisException { + return new MTMVBaseVersions(getTableVersions(mtmv), getPartitionVersions(mtmv, partitionMappings)); } - private static Map getPartitionVersions(MTMV mtmv) throws AnalysisException { + private static Map getPartitionVersions(MTMV mtmv, Map> partitionMappings) + throws AnalysisException { Map res = Maps.newHashMap(); if (mtmv.getMvPartitionInfo().getPartitionType().equals(MTMVPartitionType.SELF_MANAGE)) { return res; @@ -594,7 +601,14 @@ private static Map getPartitionVersions(MTMV mtmv) throws Analysis if (!(relatedTable instanceof OlapTable)) { return res; } - List partitions = Lists.newArrayList(((OlapTable) relatedTable).getPartitions()); + Set mappedPartitionNames = Sets.newHashSet(); + for (Set partitionNames : partitionMappings.values()) { + mappedPartitionNames.addAll(partitionNames); + } + List partitions = Lists.newArrayList(); + for (String partitionName : mappedPartitionNames) { + partitions.add(((OlapTable) relatedTable).getPartitionOrAnalysisException(partitionName)); + } List versions = null; try { versions = Partition.getVisibleVersions(partitions); diff --git a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRefreshContext.java b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRefreshContext.java index c59abd9ebdcda7..12be3dd3e04952 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRefreshContext.java +++ b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRefreshContext.java @@ -55,9 +55,14 @@ public Map getBaseTableSnapshotCache() { } public static MTMVRefreshContext buildContext(MTMV mtmv) throws AnalysisException { + return buildContext(mtmv, null); + } + + public static MTMVRefreshContext buildContext(MTMV mtmv, Set queryUsedPartitions) + throws AnalysisException { MTMVRefreshContext context = new MTMVRefreshContext(mtmv); - context.partitionMappings = mtmv.calculatePartitionMappings(); - context.baseVersions = MTMVPartitionUtil.getBaseVersions(mtmv); + context.partitionMappings = mtmv.calculatePartitionMappings(queryUsedPartitions); + context.baseVersions = MTMVPartitionUtil.getBaseVersions(mtmv, context.partitionMappings); return context; } diff --git a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescGeneratorService.java b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescGeneratorService.java index 09d85576b5cba4..c1635708a7c749 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescGeneratorService.java +++ b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescGeneratorService.java @@ -20,6 +20,7 @@ import org.apache.doris.common.AnalysisException; import java.util.Map; +import java.util.Set; /** * Interface for a series of processes to generate PartitionDesc @@ -34,5 +35,5 @@ public interface MTMVRelatedPartitionDescGeneratorService { * @throws AnalysisException */ void apply(MTMVPartitionInfo mvPartitionInfo, Map mvProperties, - RelatedPartitionDescResult lastResult) throws AnalysisException; + RelatedPartitionDescResult lastResult, Set queryUsedPartitions) throws AnalysisException; } diff --git a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescInitGenerator.java b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescInitGenerator.java index 28900f38e59f62..2e25e8fa39f6b8 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescInitGenerator.java +++ b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescInitGenerator.java @@ -21,6 +21,7 @@ import org.apache.doris.datasource.mvcc.MvccUtil; import java.util.Map; +import java.util.Set; /** * get all related partition descs @@ -29,7 +30,7 @@ public class MTMVRelatedPartitionDescInitGenerator implements MTMVRelatedPartiti @Override public void apply(MTMVPartitionInfo mvPartitionInfo, Map mvProperties, - RelatedPartitionDescResult lastResult) throws AnalysisException { + RelatedPartitionDescResult lastResult, Set queryUsedPartitions) throws AnalysisException { MTMVRelatedTableIf relatedTable = mvPartitionInfo.getRelatedTable(); lastResult.setItems(relatedTable.getAndCopyPartitionItems(MvccUtil.getSnapshotFromContext(relatedTable))); } diff --git a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescOnePartitionColGenerator.java b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescOnePartitionColGenerator.java index c5dad9bdb41891..c58fd87702b519 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescOnePartitionColGenerator.java +++ b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescOnePartitionColGenerator.java @@ -47,7 +47,7 @@ public class MTMVRelatedPartitionDescOnePartitionColGenerator implements MTMVRel @Override public void apply(MTMVPartitionInfo mvPartitionInfo, Map mvProperties, - RelatedPartitionDescResult lastResult) throws AnalysisException { + RelatedPartitionDescResult lastResult, Set queryUsedPartitions) throws AnalysisException { if (mvPartitionInfo.getPartitionType() == MTMVPartitionType.SELF_MANAGE) { return; } @@ -55,6 +55,9 @@ public void apply(MTMVPartitionInfo mvPartitionInfo, Map mvPrope Map relatedPartitionItems = lastResult.getItems(); int relatedColPos = mvPartitionInfo.getRelatedColPos(); for (Entry entry : relatedPartitionItems.entrySet()) { + if (queryUsedPartitions != null && !queryUsedPartitions.contains(entry.getKey())) { + continue; + } PartitionKeyDesc partitionKeyDesc = entry.getValue().toPartitionKeyDesc(relatedColPos); if (res.containsKey(partitionKeyDesc)) { res.get(partitionKeyDesc).add(entry.getKey()); diff --git a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescRollUpGenerator.java b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescRollUpGenerator.java index 71f7fc358f5975..ec1e888cd771d9 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescRollUpGenerator.java +++ b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescRollUpGenerator.java @@ -41,7 +41,7 @@ public class MTMVRelatedPartitionDescRollUpGenerator implements MTMVRelatedParti @Override public void apply(MTMVPartitionInfo mvPartitionInfo, Map mvProperties, - RelatedPartitionDescResult lastResult) throws AnalysisException { + RelatedPartitionDescResult lastResult, Set queryUsedPartitions) throws AnalysisException { if (mvPartitionInfo.getPartitionType() != MTMVPartitionType.EXPR) { return; } diff --git a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescSyncLimitGenerator.java b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescSyncLimitGenerator.java index 3bdcc1c57b630f..a029cbe2fe4e5c 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescSyncLimitGenerator.java +++ b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescSyncLimitGenerator.java @@ -34,6 +34,7 @@ import java.util.Map; import java.util.Map.Entry; import java.util.Optional; +import java.util.Set; /** * Only focus on partial partitions of related tables @@ -42,7 +43,7 @@ public class MTMVRelatedPartitionDescSyncLimitGenerator implements MTMVRelatedPa @Override public void apply(MTMVPartitionInfo mvPartitionInfo, Map mvProperties, - RelatedPartitionDescResult lastResult) throws AnalysisException { + RelatedPartitionDescResult lastResult, Set queryUsedPartitions) throws AnalysisException { Map partitionItems = lastResult.getItems(); MTMVPartitionSyncConfig config = generateMTMVPartitionSyncConfigByProperties(mvProperties); if (config.getSyncLimit() <= 0) { diff --git a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRewriteUtil.java b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRewriteUtil.java index 58b2a37d504810..bb962d519bf541 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRewriteUtil.java +++ b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRewriteUtil.java @@ -73,7 +73,7 @@ public static Collection getMTMVCanRewritePartitions(MTMV mtmv, Conne } if (refreshContext == null) { try { - refreshContext = MTMVRefreshContext.buildContext(mtmv); + refreshContext = MTMVRefreshContext.buildContext(mtmv, relatedPartitions); } catch (AnalysisException e) { LOG.warn("buildContext failed", e); // After failure, one should quickly return to avoid repeated failures diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/exploration/mv/PartitionCompensator.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/exploration/mv/PartitionCompensator.java index f17d5364ff0a45..e2f3ece6c8d406 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/exploration/mv/PartitionCompensator.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/exploration/mv/PartitionCompensator.java @@ -96,7 +96,8 @@ public static Pair>, Map mvValidPartitionNameSet = new HashSet<>(); Set mvValidBaseTablePartitionNameSet = new HashSet<>(); Set mvValidHasDataRelatedBaseTableNameSet = new HashSet<>(); - Pair>, Map> partitionMapping = mtmv.calculateDoublyPartitionMappings(); + Pair>, Map> partitionMapping = mtmv + .calculateDoublyPartitionMappings(queryUsedBaseTablePartitionNameSet); for (Partition mvValidPartition : mvValidPartitions) { mvValidPartitionNameSet.add(mvValidPartition.getName()); Set relatedBaseTablePartitions = partitionMapping.key().get(mvValidPartition.getName()); diff --git a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVExpandPartitionTest.java b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVExpandPartitionTest.java new file mode 100644 index 00000000000000..36eec9ad888079 --- /dev/null +++ b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVExpandPartitionTest.java @@ -0,0 +1,130 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.mtmv; + +import org.apache.doris.analysis.PartitionValue; +import org.apache.doris.catalog.Column; +import org.apache.doris.catalog.PartitionItem; +import org.apache.doris.catalog.PartitionKey; +import org.apache.doris.catalog.PartitionType; +import org.apache.doris.catalog.PrimitiveType; +import org.apache.doris.catalog.RangePartitionItem; +import org.apache.doris.catalog.ScalarType; +import org.apache.doris.common.AnalysisException; + +import com.google.common.collect.Lists; +import com.google.common.collect.Maps; +import com.google.common.collect.Range; +import com.google.common.collect.Sets; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; + +import java.lang.reflect.InvocationHandler; +import java.lang.reflect.Proxy; +import java.util.Map; +import java.util.Set; + +public class MTMVExpandPartitionTest { + + private static final Column DATE_COL = new Column("c1", ScalarType.createType(PrimitiveType.DATE), + true, null, "", ""); + + private MTMVRelatedTableIf rangeTable; + private MTMVRelatedTableIf listTable; + private Map monthlyMvPartitions; + + @Before + public void setUp() throws Exception { + Map dailyBasePartitions = Maps.newHashMap(); + dailyBasePartitions.put("p20210101", buildRange("2021-01-01", "2021-01-02")); + dailyBasePartitions.put("p20210102", buildRange("2021-01-02", "2021-01-03")); + dailyBasePartitions.put("p20210103", buildRange("2021-01-03", "2021-01-04")); + dailyBasePartitions.put("p20210201", buildRange("2021-02-01", "2021-02-02")); + dailyBasePartitions.put("p20210202", buildRange("2021-02-02", "2021-02-03")); + + monthlyMvPartitions = Maps.newHashMap(); + monthlyMvPartitions.put("mv_202101", buildRange("2021-01-01", "2021-02-01")); + monthlyMvPartitions.put("mv_202102", buildRange("2021-02-01", "2021-03-01")); + + rangeTable = createMockTable(PartitionType.RANGE, dailyBasePartitions); + listTable = createMockTable(PartitionType.LIST, Maps.newHashMap()); + } + + @Test + public void testExpandSinglePartitionToMonth() throws Exception { + Set result = MTMVPartitionExpander.expandToMvPartitionGranularity( + Sets.newHashSet("p20210102"), monthlyMvPartitions, rangeTable); + + Assert.assertEquals(Sets.newHashSet("p20210101", "p20210102", "p20210103"), result); + } + + @Test + public void testExpandMultipleMonths() throws Exception { + Set result = MTMVPartitionExpander.expandToMvPartitionGranularity( + Sets.newHashSet("p20210102", "p20210202"), monthlyMvPartitions, rangeTable); + + Assert.assertEquals(Sets.newHashSet( + "p20210101", "p20210102", "p20210103", "p20210201", "p20210202"), result); + } + + @Test + public void testListPartitionPassthrough() throws Exception { + Set queryUsedPartitions = Sets.newHashSet("p1"); + + Assert.assertSame(queryUsedPartitions, MTMVPartitionExpander.expandToMvPartitionGranularity( + queryUsedPartitions, monthlyMvPartitions, listTable)); + } + + @Test + public void testNonExistentPartition() throws Exception { + Set result = MTMVPartitionExpander.expandToMvPartitionGranularity( + Sets.newHashSet("p_nonexistent"), monthlyMvPartitions, rangeTable); + + Assert.assertTrue(result.isEmpty()); + } + + private static MTMVRelatedTableIf createMockTable( + PartitionType partitionType, Map partitionItems) { + InvocationHandler handler = (proxy, method, args) -> { + switch (method.getName()) { + case "getPartitionType": + return partitionType; + case "getAndCopyPartitionItems": + return partitionItems; + case "hashCode": + return System.identityHashCode(proxy); + case "equals": + return proxy == args[0]; + default: + throw new UnsupportedOperationException( + "MTMVExpandPartitionTest mock does not support: " + method.getName()); + } + }; + return (MTMVRelatedTableIf) Proxy.newProxyInstance( + MTMVRelatedTableIf.class.getClassLoader(), new Class[] {MTMVRelatedTableIf.class}, handler); + } + + private static RangePartitionItem buildRange(String lower, String upper) throws AnalysisException { + PartitionKey lowerKey = PartitionKey.createPartitionKey( + Lists.newArrayList(new PartitionValue(lower)), Lists.newArrayList(DATE_COL)); + PartitionKey upperKey = PartitionKey.createPartitionKey( + Lists.newArrayList(new PartitionValue(upper)), Lists.newArrayList(DATE_COL)); + return new RangePartitionItem(Range.closedOpen(lowerKey, upperKey)); + } +} diff --git a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVPartitionUtilTest.java b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVPartitionUtilTest.java index d90b43482eba1a..e709f03e4c1622 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVPartitionUtilTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVPartitionUtilTest.java @@ -32,11 +32,13 @@ import com.google.common.collect.Maps; import com.google.common.collect.Sets; import mockit.Expectations; +import mockit.Injectable; import mockit.Mocked; import org.junit.Assert; import org.junit.Before; import org.junit.Test; +import java.util.Collection; import java.util.List; import java.util.Map; import java.util.Optional; @@ -45,7 +47,7 @@ public class MTMVPartitionUtilTest { @Mocked private MTMV mtmv; - @Mocked + @Injectable private Partition p1; @Mocked private MTMVRelation relation; @@ -53,7 +55,7 @@ public class MTMVPartitionUtilTest { private BaseTableInfo baseTableInfo; @Mocked private MTMVPartitionInfo mtmvPartitionInfo; - @Mocked + @Injectable private OlapTable baseOlapTable; @Mocked private DatabaseIf databaseIf; @@ -67,7 +69,7 @@ public class MTMVPartitionUtilTest { private MTMVUtil mtmvUtil; @Mocked private MTMVRefreshContext context; - @Mocked + @Injectable private MTMVBaseVersions versions; private Set baseTables = Sets.newHashSet(); @@ -117,6 +119,10 @@ public void setUp() throws NoSuchMethodException, SecurityException, AnalysisExc minTimes = 0; result = MTMVPartitionType.SELF_MANAGE; + mtmvPartitionInfo.getRelatedTable(); + minTimes = 0; + result = baseOlapTable; + mtmvUtil.getTable(baseTableInfo); minTimes = 0; result = baseOlapTable; @@ -302,6 +308,25 @@ public void testIsTableNamelike() { Assert.assertFalse(MTMVPartitionUtil.isTableNamelike(new TableName("ctl1"), tableNameToCheck)); } + @Test + public void testGetBaseVersionsUsesMappedPartitions() throws AnalysisException { + Map> partitionMappings = Maps.newHashMap(); + partitionMappings.put("mv1", Sets.newHashSet("p1")); + + assertFetchedPartitionNames(partitionMappings, Sets.newHashSet("p1", "p2", "p3"), + Sets.newHashSet("p1")); + } + + @Test + public void testGetBaseVersionsDeduplicatesMappedPartitions() throws AnalysisException { + Map> partitionMappings = Maps.newHashMap(); + partitionMappings.put("mv1", Sets.newHashSet("p1", "p2")); + partitionMappings.put("mv2", Sets.newHashSet("p2", "p3")); + + assertFetchedPartitionNames(partitionMappings, Sets.newHashSet("p1", "p2", "p3", "p4"), + Sets.newHashSet("p1", "p2", "p3")); + } + @Test public void testGetTableSnapshotFromContext() throws AnalysisException { Map cache = Maps.newHashMap(); @@ -317,4 +342,47 @@ public void testGetTableSnapshotFromContext() throws AnalysisException { Assert.assertEquals(1, cache.size()); Assert.assertEquals(baseSnapshotIf, cache.values().iterator().next()); } + + private void assertFetchedPartitionNames(Map> partitionMappings, + Set allPartitionNames, Set expectedPartitionNames) throws AnalysisException { + Map partitions = Maps.newHashMap(); + long partitionId = 1; + for (String partitionName : allPartitionNames) { + partitions.put(partitionName, new Partition(partitionId++, partitionName, null, null)); + } + List fetchedPartitionNames = Lists.newArrayList(); + OlapTable relatedTable = new OlapTable() { + @Override + public Partition getPartitionOrAnalysisException(String partitionName) { + fetchedPartitionNames.add(partitionName); + return partitions.get(partitionName); + } + + @Override + public Collection getPartitions() { + Assert.fail("getPartitions should not be called"); + return null; + } + }; + new Expectations() { + { + mtmv.getRelation(); + minTimes = 0; + result = null; + + mtmvPartitionInfo.getPartitionType(); + minTimes = 0; + result = MTMVPartitionType.FOLLOW_BASE_TABLE; + + mtmvPartitionInfo.getRelatedTable(); + minTimes = 0; + result = relatedTable; + } + }; + + Assert.assertEquals(expectedPartitionNames, + MTMVPartitionUtil.getBaseVersions(mtmv, partitionMappings).getPartitionVersions().keySet()); + Assert.assertEquals(expectedPartitionNames, Sets.newHashSet(fetchedPartitionNames)); + Assert.assertEquals(expectedPartitionNames.size(), fetchedPartitionNames.size()); + } } diff --git a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescOnePartitionColGeneratorTest.java b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescOnePartitionColGeneratorTest.java new file mode 100644 index 00000000000000..a6f26d5cac655e --- /dev/null +++ b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVRelatedPartitionDescOnePartitionColGeneratorTest.java @@ -0,0 +1,82 @@ +// 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.mtmv; + +import org.apache.doris.analysis.PartitionValue; +import org.apache.doris.catalog.Column; +import org.apache.doris.catalog.PartitionItem; +import org.apache.doris.catalog.PartitionKey; +import org.apache.doris.catalog.PrimitiveType; +import org.apache.doris.catalog.RangePartitionItem; +import org.apache.doris.catalog.ScalarType; +import org.apache.doris.common.AnalysisException; +import org.apache.doris.mtmv.MTMVPartitionInfo.MTMVPartitionType; + +import com.google.common.collect.Lists; +import com.google.common.collect.Maps; +import com.google.common.collect.Range; +import com.google.common.collect.Sets; +import mockit.Expectations; +import mockit.Mocked; +import org.junit.Assert; +import org.junit.Test; + +import java.util.Map; +import java.util.Set; + +public class MTMVRelatedPartitionDescOnePartitionColGeneratorTest { + + private static final Column DATE_COL = new Column("c1", ScalarType.createType(PrimitiveType.DATE), + true, null, "", ""); + + @Mocked + private MTMVPartitionInfo partitionInfo; + + @Test + public void testQueryUsedPartitionsFilter() throws Exception { + new Expectations() { + { + partitionInfo.getPartitionType(); + result = MTMVPartitionType.FOLLOW_BASE_TABLE; + + partitionInfo.getRelatedColPos(); + result = 0; + } + }; + RelatedPartitionDescResult result = new RelatedPartitionDescResult(); + Map partitionItems = Maps.newHashMap(); + partitionItems.put("p1", buildRange("2021-01-01", "2021-01-02")); + partitionItems.put("p2", buildRange("2021-01-02", "2021-01-03")); + result.setItems(partitionItems); + + new MTMVRelatedPartitionDescOnePartitionColGenerator().apply( + partitionInfo, Maps.newHashMap(), result, Sets.newHashSet("p2")); + + Set mappedPartitions = Sets.newHashSet(); + result.getDescs().values().forEach(mappedPartitions::addAll); + Assert.assertEquals(Sets.newHashSet("p2"), mappedPartitions); + } + + private static RangePartitionItem buildRange(String lower, String upper) throws AnalysisException { + PartitionKey lowerKey = PartitionKey.createPartitionKey( + Lists.newArrayList(new PartitionValue(lower)), Lists.newArrayList(DATE_COL)); + PartitionKey upperKey = PartitionKey.createPartitionKey( + Lists.newArrayList(new PartitionValue(upper)), Lists.newArrayList(DATE_COL)); + return new RangePartitionItem(Range.closedOpen(lowerKey, upperKey)); + } +}