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
Original file line number Diff line number Diff line change
Expand Up @@ -375,17 +375,18 @@ protected List<Plan> doRewrite(StructInfo queryStructInfo, CascadesContext casca
// if mv can not offer any partition for query, query rewrite bail out to avoid cycle run
return rewriteResults;
}
boolean partitionNeedUnion = PartitionCompensator.needUnionRewrite(invalidPartitions, cascadesContext);
boolean canUnionRewrite = canUnionRewrite(queryPlan,
(AsyncMaterializationContext) materializationContext, cascadesContext);
if (partitionNeedUnion && !canUnionRewrite) {
boolean hasPartitionCompensation =
PartitionCompensator.hasPartitionCompensation(invalidPartitions);
boolean needBaseTableUnion = !invalidPartitions.value().isEmpty();
if (needBaseTableUnion && !canUnionRewrite(queryPlan,
(AsyncMaterializationContext) materializationContext, cascadesContext)) {
materializationContext.recordFailReason(queryStructInfo,
"need compensate union all, but can not, because the query structInfo",
() -> String.format("mv partition info is %s, and the query plan is %s",
mtmv.getMvPartitionInfo(), queryPlan.treeString()));
return rewriteResults;
}
if (partitionNeedUnion) {
if (hasPartitionCompensation) {
Pair<Plan, Boolean> planAndNeedAddFilterPair =
StructInfo.addFilterOnTableScan(queryPlan, invalidPartitions.value(), cascadesContext);
if (planAndNeedAddFilterPair == null) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -105,8 +105,8 @@ protected Pair<Map<BaseTableInfo, Set<String>>, Map<BaseColInfo, Set<String>>> c
Pair<Map<BaseTableInfo, Set<String>>, Map<BaseColInfo, Set<String>>> invalidPartitions
= super.calcInvalidPartitions(queryUsedBaseTablePartitionMap, rewrittenPlan, cascadesContext,
materializationContext);
if (PartitionCompensator.needUnionRewrite(invalidPartitions, cascadesContext)) {
// if query use some invalid partition in mv, bail out
if (PartitionCompensator.hasPartitionCompensation(invalidPartitions)) {
// Aggregate-on-non-aggregate rewrite does not support partition compensation.
return null;
}
return invalidPartitions;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -218,11 +218,10 @@ private static Pair<Pair<BaseTableInfo, Set<String>>, Pair<BaseColInfo, Set<Stri
return Pair.of(mvPartitionNeedRemoveNameMap, baseTablePartitionNeedUnionNameMap);
}

public static boolean needUnionRewrite(
Pair<Map<BaseTableInfo, Set<String>>, Map<BaseColInfo, Set<String>>> invalidPartitions,
CascadesContext cascadesContext) {
public static boolean hasPartitionCompensation(
Pair<Map<BaseTableInfo, Set<String>>, Map<BaseColInfo, Set<String>>> invalidPartitions) {
return invalidPartitions != null
&& (!invalidPartitions.key().values().isEmpty() || !invalidPartitions.value().values().isEmpty());
&& (!invalidPartitions.key().isEmpty() || !invalidPartitions.value().isEmpty());
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@

import com.google.common.collect.ArrayListMultimap;
import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableMap;
import com.google.common.collect.ImmutableSet;
import com.google.common.collect.Multimap;
import org.junit.jupiter.api.Assertions;
Expand Down Expand Up @@ -290,6 +291,20 @@ public void testNotNeedUnionRewriteWhenAllPartitions() throws Exception {
Assertions.assertFalse(PartitionCompensator.needUnionRewrite(ctx, sc));
}

@Test
public void testHasPartitionCompensation() {
BaseTableInfo mvTableInfo = newBaseTableInfo();
BaseColInfo baseColInfo = new BaseColInfo("c", newBaseTableInfo());

Assertions.assertFalse(PartitionCompensator.hasPartitionCompensation(null));
Assertions.assertFalse(PartitionCompensator.hasPartitionCompensation(
Pair.of(ImmutableMap.of(), ImmutableMap.of())));
Assertions.assertTrue(PartitionCompensator.hasPartitionCompensation(
Pair.of(ImmutableMap.of(mvTableInfo, ImmutableSet.of("mv_p1")), ImmutableMap.of())));
Assertions.assertTrue(PartitionCompensator.hasPartitionCompensation(
Pair.of(ImmutableMap.of(), ImmutableMap.of(baseColInfo, ImmutableSet.of("base_p1")))));
}

@Test
public void testGetQueryUsedPartitionsAllAndPartial() {
// Prepare qualifiers
Expand Down
Loading