Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
32 changes: 27 additions & 5 deletions fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -381,17 +382,20 @@ public Map<String, PartitionKeyDesc> generateMvPartitionDescs() throws AnalysisE
* @return mvPartitionName ==> relationPartitionNames and relationPartitionName ==> mvPartitionName
* @throws AnalysisException
*/
public Pair<Map<String, Set<String>>, Map<String, String>> calculateDoublyPartitionMappings()
public Pair<Map<String, Set<String>>, Map<String, String>> calculateDoublyPartitionMappings(
Set<String> queryUsedBaseTablePartitions)
throws AnalysisException {
if (mvPartitionInfo.getPartitionType() == MTMVPartitionType.SELF_MANAGE) {
return Pair.of(Maps.newHashMap(), Maps.newHashMap());
}
long start = System.currentTimeMillis();
Map<String, Set<String>> mvToBase = Maps.newHashMap();
Map<String, String> baseToMv = Maps.newHashMap();
Map<PartitionKeyDesc, Set<String>> relatedPartitionDescs = MTMVPartitionUtil
.generateRelatedPartitionDescs(mvPartitionInfo, mvProperties);
Map<String, PartitionItem> mvPartitionItems = getAndCopyPartitionItems();
Set<String> effectiveFilter = getEffectiveQueryUsedBaseTablePartitions(
queryUsedBaseTablePartitions, mvPartitionItems);
Map<PartitionKeyDesc, Set<String>> relatedPartitionDescs = MTMVPartitionUtil
.generateRelatedPartitionDescs(mvPartitionInfo, mvProperties, effectiveFilter);
for (Entry<String, PartitionItem> entry : mvPartitionItems.entrySet()) {
Set<String> basePartitionNames = relatedPartitionDescs.getOrDefault(entry.getValue().toPartitionKeyDesc(),
Sets.newHashSet());
Expand All @@ -417,14 +421,21 @@ public Pair<Map<String, Set<String>>, Map<String, String>> calculateDoublyPartit
* @throws AnalysisException
*/
public Map<String, Set<String>> calculatePartitionMappings() throws AnalysisException {
return calculatePartitionMappings(null);
}

public Map<String, Set<String>> calculatePartitionMappings(Set<String> queryUsedBaseTablePartitions)
throws AnalysisException {
if (mvPartitionInfo.getPartitionType() == MTMVPartitionType.SELF_MANAGE) {
return Maps.newHashMap();
}
long start = System.currentTimeMillis();
Map<String, PartitionItem> mvPartitionItems = getAndCopyPartitionItems();
Set<String> effectiveFilter = getEffectiveQueryUsedBaseTablePartitions(
queryUsedBaseTablePartitions, mvPartitionItems);
Map<String, Set<String>> res = Maps.newHashMap();
Map<PartitionKeyDesc, Set<String>> relatedPartitionDescs = MTMVPartitionUtil
.generateRelatedPartitionDescs(mvPartitionInfo, mvProperties);
Map<String, PartitionItem> mvPartitionItems = getAndCopyPartitionItems();
.generateRelatedPartitionDescs(mvPartitionInfo, mvProperties, effectiveFilter);
for (Entry<String, PartitionItem> entry : mvPartitionItems.entrySet()) {
res.put(entry.getKey(),
relatedPartitionDescs.getOrDefault(entry.getValue().toPartitionKeyDesc(), Sets.newHashSet()));
Expand All @@ -436,6 +447,17 @@ public Map<String, Set<String>> calculatePartitionMappings() throws AnalysisExce
return res;
}

private Set<String> getEffectiveQueryUsedBaseTablePartitions(
Set<String> queryUsedBaseTablePartitions, Map<String, PartitionItem> mvPartitionItems)
throws AnalysisException {
if (queryUsedBaseTablePartitions == null
|| mvPartitionInfo.getPartitionType() != MTMVPartitionType.EXPR) {
return queryUsedBaseTablePartitions;
}
return MTMVPartitionExpander.expandToMvPartitionGranularity(queryUsedBaseTablePartitions,
mvPartitionItems, mvPartitionInfo.getRelatedTable());
}

public ConcurrentLinkedQueue<MTMVTask> getHistoryTasks() {
return jobInfo.getHistoryTasks();
}
Expand Down
Original file line number Diff line number Diff line change
@@ -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<String> expandToMvPartitionGranularity(
Set<String> queryUsedBaseTablePartitions,
Map<String, PartitionItem> mvPartitionItems,
MTMVRelatedTableIf relatedTable) throws AnalysisException {
Optional<MvccSnapshot> snapshot = MvccUtil.getSnapshotFromContext(relatedTable);
if (relatedTable.getPartitionType(snapshot) != PartitionType.RANGE) {
return queryUsedBaseTablePartitions;
}

NavigableMap<PartitionKey, Range<PartitionKey>> mvRanges = new TreeMap<>();
for (PartitionItem item : mvPartitionItems.values()) {
Range<PartitionKey> range = ((RangePartitionItem) item).getItems();
mvRanges.put(range.lowerEndpoint(), range);
}

Map<String, PartitionItem> basePartitionItems = relatedTable.getAndCopyPartitionItems(snapshot);
NavigableMap<PartitionKey, Range<PartitionKey>> relevantMvRanges = new TreeMap<>();
for (String queriedBasePartition : queryUsedBaseTablePartitions) {
PartitionItem baseItem = basePartitionItems.get(queriedBasePartition);
if (baseItem == null) {
continue;
}
Range<PartitionKey> baseRange = ((RangePartitionItem) baseItem).getItems();
Range<PartitionKey> mvRange = findEnclosingRange(mvRanges, baseRange);
if (mvRange != null) {
relevantMvRanges.put(mvRange.lowerEndpoint(), mvRange);
}
}

Set<String> expandedPartitions = Sets.newHashSet();
for (Entry<String, PartitionItem> baseEntry : basePartitionItems.entrySet()) {
Range<PartitionKey> baseRange = ((RangePartitionItem) baseEntry.getValue()).getItems();
if (findEnclosingRange(relevantMvRanges, baseRange) != null) {
expandedPartitions.add(baseEntry.getKey());
}
}
return expandedPartitions;
}

private static Range<PartitionKey> findEnclosingRange(
NavigableMap<PartitionKey, Range<PartitionKey>> ranges, Range<PartitionKey> baseRange) {
Entry<PartitionKey, Range<PartitionKey>> candidate = ranges.floorEntry(baseRange.lowerEndpoint());
return candidate != null && candidate.getValue().encloses(baseRange) ? candidate.getValue() : null;
}

private MTMVPartitionExpander() {
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -166,10 +166,15 @@ public static List<AllPartitionDesc> getPartitionDescsByRelatedTable(

public static Map<PartitionKeyDesc, Set<String>> generateRelatedPartitionDescs(MTMVPartitionInfo mvPartitionInfo,
Map<String, String> mvProperties) throws AnalysisException {
return generateRelatedPartitionDescs(mvPartitionInfo, mvProperties, null);
}

public static Map<PartitionKeyDesc, Set<String>> generateRelatedPartitionDescs(MTMVPartitionInfo mvPartitionInfo,
Map<String, String> mvProperties, Set<String> 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 [{}]",
Expand Down Expand Up @@ -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<String, Set<String>> partitionMappings)
throws AnalysisException {
return new MTMVBaseVersions(getTableVersions(mtmv), getPartitionVersions(mtmv, partitionMappings));
}

private static Map<String, Long> getPartitionVersions(MTMV mtmv) throws AnalysisException {
private static Map<String, Long> getPartitionVersions(MTMV mtmv, Map<String, Set<String>> partitionMappings)
throws AnalysisException {
Map<String, Long> res = Maps.newHashMap();
if (mtmv.getMvPartitionInfo().getPartitionType().equals(MTMVPartitionType.SELF_MANAGE)) {
return res;
Expand All @@ -594,7 +601,14 @@ private static Map<String, Long> getPartitionVersions(MTMV mtmv) throws Analysis
if (!(relatedTable instanceof OlapTable)) {
return res;
}
List<Partition> partitions = Lists.newArrayList(((OlapTable) relatedTable).getPartitions());
Set<String> mappedPartitionNames = Sets.newHashSet();
for (Set<String> partitionNames : partitionMappings.values()) {
mappedPartitionNames.addAll(partitionNames);
}
List<Partition> partitions = Lists.newArrayList();
for (String partitionName : mappedPartitionNames) {
partitions.add(((OlapTable) relatedTable).getPartitionOrAnalysisException(partitionName));
}
List<Long> versions = null;
try {
versions = Partition.getVisibleVersions(partitions);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -55,9 +55,14 @@ public Map<BaseTableInfo, MTMVSnapshotIf> getBaseTableSnapshotCache() {
}

public static MTMVRefreshContext buildContext(MTMV mtmv) throws AnalysisException {
return buildContext(mtmv, null);
}

public static MTMVRefreshContext buildContext(MTMV mtmv, Set<String> 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;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -34,5 +35,5 @@ public interface MTMVRelatedPartitionDescGeneratorService {
* @throws AnalysisException
*/
void apply(MTMVPartitionInfo mvPartitionInfo, Map<String, String> mvProperties,
RelatedPartitionDescResult lastResult) throws AnalysisException;
RelatedPartitionDescResult lastResult, Set<String> queryUsedPartitions) throws AnalysisException;
}
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
import org.apache.doris.datasource.mvcc.MvccUtil;

import java.util.Map;
import java.util.Set;

/**
* get all related partition descs
Expand All @@ -29,7 +30,7 @@ public class MTMVRelatedPartitionDescInitGenerator implements MTMVRelatedPartiti

@Override
public void apply(MTMVPartitionInfo mvPartitionInfo, Map<String, String> mvProperties,
RelatedPartitionDescResult lastResult) throws AnalysisException {
RelatedPartitionDescResult lastResult, Set<String> queryUsedPartitions) throws AnalysisException {
MTMVRelatedTableIf relatedTable = mvPartitionInfo.getRelatedTable();
lastResult.setItems(relatedTable.getAndCopyPartitionItems(MvccUtil.getSnapshotFromContext(relatedTable)));
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -47,14 +47,17 @@ public class MTMVRelatedPartitionDescOnePartitionColGenerator implements MTMVRel

@Override
public void apply(MTMVPartitionInfo mvPartitionInfo, Map<String, String> mvProperties,
RelatedPartitionDescResult lastResult) throws AnalysisException {
RelatedPartitionDescResult lastResult, Set<String> queryUsedPartitions) throws AnalysisException {
if (mvPartitionInfo.getPartitionType() == MTMVPartitionType.SELF_MANAGE) {
return;
}
Map<PartitionKeyDesc, Set<String>> res = Maps.newHashMap();
Map<String, PartitionItem> relatedPartitionItems = lastResult.getItems();
int relatedColPos = mvPartitionInfo.getRelatedColPos();
for (Entry<String, PartitionItem> 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());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ public class MTMVRelatedPartitionDescRollUpGenerator implements MTMVRelatedParti

@Override
public void apply(MTMVPartitionInfo mvPartitionInfo, Map<String, String> mvProperties,
RelatedPartitionDescResult lastResult) throws AnalysisException {
RelatedPartitionDescResult lastResult, Set<String> queryUsedPartitions) throws AnalysisException {
if (mvPartitionInfo.getPartitionType() != MTMVPartitionType.EXPR) {
return;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -42,7 +43,7 @@ public class MTMVRelatedPartitionDescSyncLimitGenerator implements MTMVRelatedPa

@Override
public void apply(MTMVPartitionInfo mvPartitionInfo, Map<String, String> mvProperties,
RelatedPartitionDescResult lastResult) throws AnalysisException {
RelatedPartitionDescResult lastResult, Set<String> queryUsedPartitions) throws AnalysisException {
Map<String, PartitionItem> partitionItems = lastResult.getItems();
MTMVPartitionSyncConfig config = generateMTMVPartitionSyncConfigByProperties(mvProperties);
if (config.getSyncLimit() <= 0) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,7 @@ public static Collection<Partition> 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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -96,7 +96,8 @@ public static Pair<Map<BaseTableInfo, Set<String>>, Map<BaseTableInfo, Set<Strin
Set<String> mvValidPartitionNameSet = new HashSet<>();
Set<String> mvValidBaseTablePartitionNameSet = new HashSet<>();
Set<String> mvValidHasDataRelatedBaseTableNameSet = new HashSet<>();
Pair<Map<String, Set<String>>, Map<String, String>> partitionMapping = mtmv.calculateDoublyPartitionMappings();
Pair<Map<String, Set<String>>, Map<String, String>> partitionMapping = mtmv
.calculateDoublyPartitionMappings(queryUsedBaseTablePartitionNameSet);
for (Partition mvValidPartition : mvValidPartitions) {
mvValidPartitionNameSet.add(mvValidPartition.getName());
Set<String> relatedBaseTablePartitions = partitionMapping.key().get(mvValidPartition.getName());
Expand Down
Loading
Loading