Skip to content
Merged
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 @@ -49,6 +49,7 @@
import org.apache.doris.nereids.trees.plans.logical.LogicalFilter;
import org.apache.doris.nereids.trees.plans.logical.LogicalJoin;
import org.apache.doris.nereids.trees.plans.logical.LogicalLimit;
import org.apache.doris.nereids.trees.plans.logical.LogicalPartitionTopN;
import org.apache.doris.nereids.trees.plans.logical.LogicalProject;
import org.apache.doris.nereids.trees.plans.logical.LogicalRelation;
import org.apache.doris.nereids.trees.plans.logical.LogicalResultSink;
Expand All @@ -61,6 +62,7 @@
import org.apache.doris.nereids.trees.plans.visitor.DefaultPlanVisitor;
import org.apache.doris.nereids.trees.plans.visitor.ExpressionLineageReplacer;
import org.apache.doris.nereids.types.DataType;
import org.apache.doris.nereids.util.ImmutableEqualSet;

import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableSet;
Expand Down Expand Up @@ -430,6 +432,9 @@ && new BaseTableInfo(((LogicalCatalogRelation) relation).getTable())
@Override
public Void visitLogicalAggregate(LogicalAggregate<? extends Plan> aggregate,
PartitionIncrementCheckContext context) {
if (!planOutputContainsPartitionColumnToCheck(aggregate, context)) {
return super.visitLogicalAggregate(aggregate, context);
}
Set<Expression> groupByExprSet = new HashSet<>(aggregate.getGroupByExpressions());
if (groupByExprSet.isEmpty()) {
context.addFailReason("group by sets is empty, doesn't contain the target partition");
Expand All @@ -446,6 +451,9 @@ public Void visitLogicalAggregate(LogicalAggregate<? extends Plan> aggregate,

@Override
public Void visitLogicalWindow(LogicalWindow<? extends Plan> window, PartitionIncrementCheckContext context) {
if (!planOutputContainsPartitionColumnToCheck(window, context)) {
return super.visitLogicalWindow(window, context);
}
List<NamedExpression> windowExpressions = window.getWindowExpressions();
if (windowExpressions.isEmpty()) {
context.addFailReason("window expression is empty, doesn't contain the target partition");
Expand All @@ -454,7 +462,7 @@ public Void visitLogicalWindow(LogicalWindow<? extends Plan> window, PartitionIn
return visit(window, context);
}
for (NamedExpression namedExpression : windowExpressions) {
if (!checkWindowPartition(namedExpression, context)) {
if (!checkWindowPartition(namedExpression, window, context)) {
context.addFailReason("window partition sets doesn't contain the target partition");
context.collectFailedTableSet(window);
context.setFailFast(true);
Expand All @@ -464,6 +472,29 @@ public Void visitLogicalWindow(LogicalWindow<? extends Plan> window, PartitionIn
return super.visitLogicalWindow(window, context);
}

@Override
public Void visitLogicalPartitionTopN(LogicalPartitionTopN<? extends Plan> partitionTopN,
PartitionIncrementCheckContext context) {
if (!planOutputContainsPartitionColumnToCheck(partitionTopN, context)) {
return super.visitLogicalPartitionTopN(partitionTopN, context);
}
if (partitionTopN.hasGlobalLimit()) {
// A global limit/topN selects the top rows across all partitions rather than within
// each partition, so a change in one source partition can move the global winner and
// affect multiple MV partitions. That breaks partition-local maintenance, so reject it.
context.addFailReason("partition topN has global limit, which is not partition local");
context.collectFailedTableSet(partitionTopN);
context.setFailFast(true);
return super.visitLogicalPartitionTopN(partitionTopN, context);
}
if (!checkPartitionKeysContainPartitionToCheck(partitionTopN.getPartitionKeys(), partitionTopN, context)) {
context.addFailReason("partition topN partition keys doesn't contain the target partition");
context.collectFailedTableSet(partitionTopN);
context.setFailFast(true);
}
return super.visitLogicalPartitionTopN(partitionTopN, context);
}

@Override
public Void visit(Plan plan, PartitionIncrementCheckContext context) {
if (plan instanceof LogicalProject
Expand All @@ -476,6 +507,7 @@ public Void visit(Plan plan, PartitionIncrementCheckContext context) {
|| plan instanceof LogicalWindow
|| (plan instanceof LogicalUnion
&& ((LogicalUnion) plan).getQualifier() == SetOperation.Qualifier.ALL)
|| plan instanceof LogicalPartitionTopN
|| plan instanceof LogicalCTEAnchor
|| plan instanceof LogicalCTEConsumer
|| plan instanceof LogicalCTEProducer
Expand All @@ -491,30 +523,58 @@ public Void visit(Plan plan, PartitionIncrementCheckContext context) {
return super.visit(plan, context);
}

private boolean checkWindowPartition(Expression expression, PartitionIncrementCheckContext context) {
private boolean checkWindowPartition(Expression expression, Plan currentPlan,
PartitionIncrementCheckContext context) {
List<Object> windowExpressions =
expression.collectToList(expressionTreeNode -> expressionTreeNode instanceof WindowExpression);
for (Object windowExpressionObj : windowExpressions) {
WindowExpression windowExpression = (WindowExpression) windowExpressionObj;
List<Expression> partitionKeys = windowExpression.getPartitionKeys();
Set<Column> originalPartitionbyExprSet = new HashSet<>();
partitionKeys.forEach(groupExpr -> {
if (groupExpr instanceof SlotReference && groupExpr.isColumnFromTable()) {
originalPartitionbyExprSet.add(((SlotReference) groupExpr).getOriginalColumn().get());
}
});
Set<SlotReference> contextPartitionColumnSet = getPartitionColumnsToCheck(context);
if (contextPartitionColumnSet.isEmpty()) {
return false;
}
if (contextPartitionColumnSet.stream().noneMatch(
partition -> originalPartitionbyExprSet.contains(partition.getOriginalColumn().get()))) {
if (!checkPartitionKeysContainPartitionToCheck(windowExpression.getPartitionKeys(),
currentPlan, context)) {
return false;
}
}
return true;
}

private boolean checkPartitionKeysContainPartitionToCheck(List<Expression> partitionKeys,
Plan currentPlan, PartitionIncrementCheckContext context) {
// Match by slot exprId, not catalog Column: Column.equals ignores the owning table, so an
// unrelated same-schema column (e.g. l.p vs r.p) would be wrongly accepted as the tracked key.
Set<SlotReference> partitionKeySlotSet = new HashSet<>();
partitionKeys.forEach(partitionKey -> {
if (partitionKey instanceof SlotReference && partitionKey.isColumnFromTable()) {
partitionKeySlotSet.add((SlotReference) partitionKey);
}
});
Set<SlotReference> contextPartitionColumnSet = getPartitionColumnsToCheck(context);
if (contextPartitionColumnSet.isEmpty()) {
return false;
}
if (contextPartitionColumnSet.stream().anyMatch(partitionKeySlotSet::contains)) {
return true;
}
// Also accept a key that is a different slot but proven equal to the tracked column, e.g.
// window over (partition by r.p) on an inner join l.p = r.p, or a forwarding alias p AS p_alias.
// currentPlan's DataTrait carries these equalities (bottom-up) even though this top-down check
// runs before the join equalities reach the context.
ImmutableEqualSet<Slot> equalSet = currentPlan.getLogicalProperties().getTrait().getEqualSet();
return contextPartitionColumnSet.stream().anyMatch(contextSlot ->
partitionKeySlotSet.stream().anyMatch(keySlot -> equalSet.isEqual(contextSlot, keySlot)));
}

private boolean planOutputContainsPartitionColumnToCheck(Plan plan, PartitionIncrementCheckContext context) {
Set<Slot> outputSet = plan.getOutputSet();
for (NamedExpression namedExpression : context.getPartitionAndRefExpressionMap().keySet()) {
// Plan outputs are slots, so only tracked partition expressions already resolved
// to slots can match here.
if (namedExpression instanceof Slot && outputSet.contains(namedExpression)) {
return true;
}
}
return false;
}

private Set<SlotReference> getPartitionColumnsToCheck(PartitionIncrementCheckContext context) {
Set<NamedExpression> partitionExpressionSet = context.getPartitionAndRefExpressionMap().keySet();
Set<SlotReference> partitionSlotSet = new HashSet<>();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -644,6 +644,80 @@ public void getRelatedTableInfoTestWithWindowButNotPartitionTest() {
});
}

@Test
public void getRelatedTableInfoTestWithUnrelatedRightWindowTest() {
PlanChecker.from(connectContext)
.checkExplain("SELECT l.L_SHIPDATE, l.L_ORDERKEY, o.O_ORDERDATE "
+ "FROM lineitem as l "
+ "LEFT JOIN ("
+ "SELECT O_ORDERKEY, O_ORDERDATE, "
+ "ROW_NUMBER() OVER (PARTITION BY O_ORDERKEY ORDER BY O_ORDERDATE DESC) AS rn "
+ "FROM orders"
+ ") as o "
+ "ON l.L_ORDERKEY = o.O_ORDERKEY AND o.rn = 1",
nereidsPlanner -> {
Plan rewrittenPlan = nereidsPlanner.getRewrittenPlan();
RelatedTableInfo relatedTableInfo =
MaterializedViewUtils.getRelatedTableInfo("L_SHIPDATE", null,
rewrittenPlan, nereidsPlanner.getCascadesContext());
Assertions.assertTrue(relatedTableInfo.isPctPossible(), relatedTableInfo.getFailReason());
checkRelatedTableInfo(relatedTableInfo,
"lineitem",
"L_SHIPDATE",
true);
});
}

@Test
public void getRelatedTableInfoTestWithUnrelatedSameSchemaWindowTest() {
// t1 and t2 come from the same table, so t1.L_SHIPDATE and t2.L_SHIPDATE share an identical
// catalog Column definition but belong to different table instances. The MV partitions by
// t1.L_SHIPDATE while the row_number() partitions by the unrelated t2.L_SHIPDATE, and the join
// condition is on L_ORDERKEY so the two shipdate slots are NOT in the same equal set. Partition
// tracking must reject this: matching by bare Column would wrongly treat t2.L_SHIPDATE as the
// tracked partition key and allow stale rows after a partition-only refresh.
PlanChecker.from(connectContext)
.checkExplain("SELECT t1.L_SHIPDATE, t1.L_ORDERKEY "
+ "FROM lineitem t1 "
+ "LEFT JOIN ("
+ "SELECT L_ORDERKEY, L_SHIPDATE, "
+ "ROW_NUMBER() OVER (PARTITION BY L_SHIPDATE ORDER BY L_PARTKEY DESC) AS rn "
+ "FROM lineitem"
+ ") as t2 "
+ "ON t1.L_ORDERKEY = t2.L_ORDERKEY AND t2.rn = 1",
nereidsPlanner -> {
Plan rewrittenPlan = nereidsPlanner.getRewrittenPlan();
RelatedTableInfo relatedTableInfo =
MaterializedViewUtils.getRelatedTableInfo("L_SHIPDATE", null,
rewrittenPlan, nereidsPlanner.getCascadesContext());
Assertions.assertFalse(relatedTableInfo.isPctPossible());
});
}

@Test
public void getRelatedTableInfoTestWithUnrelatedRightAggregateTest() {
PlanChecker.from(connectContext)
.checkExplain("SELECT l.L_SHIPDATE, l.L_ORDERKEY, o.max_orderdate "
+ "FROM lineitem as l "
+ "LEFT JOIN ("
+ "SELECT O_ORDERKEY, max(O_ORDERDATE) AS max_orderdate "
+ "FROM orders "
+ "GROUP BY O_ORDERKEY"
+ ") as o "
+ "ON l.L_ORDERKEY = o.O_ORDERKEY",
nereidsPlanner -> {
Plan rewrittenPlan = nereidsPlanner.getRewrittenPlan();
RelatedTableInfo relatedTableInfo =
MaterializedViewUtils.getRelatedTableInfo("L_SHIPDATE", null,
rewrittenPlan, nereidsPlanner.getCascadesContext());
Assertions.assertTrue(relatedTableInfo.isPctPossible(), relatedTableInfo.getFailReason());
checkRelatedTableInfo(relatedTableInfo,
"lineitem",
"L_SHIPDATE",
true);
});
}

@Test
public void getRelatedTableInfoTestWithLimitTest() {
PlanChecker.from(connectContext)
Expand Down Expand Up @@ -882,9 +956,11 @@ public void getRelatedTableInfoWhenMultiBaseTablePartition() {
RelatedTableInfo relatedTableInfo =
MaterializedViewUtils.getRelatedTableInfo("upgrade_day", null,
rewrittenPlan, nereidsPlanner.getCascadesContext());
Assertions.assertTrue(relatedTableInfo.getFailReason().contains(
"partition column is not in group by or window partition by"));
Assertions.assertFalse(relatedTableInfo.isPctPossible());
Assertions.assertTrue(relatedTableInfo.isPctPossible(), relatedTableInfo.getFailReason());
checkRelatedTableInfo(relatedTableInfo,
"test1",
"upgrade_day",
true);
});
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2232,6 +2232,21 @@ class Suite implements GroovyInterceptable {
logger.info("index stats: " + stats.toString())
}

void waitingMTMVTaskFinishedWithoutAnalyze(String jobName) {
// Wait for the newly submitted MTMV task to become visible in tasks().
Thread.sleep(2000);
String showTasks = """
select TaskId, Status from tasks('type'='mv')
where JobName = '${jobName}' order by CreateTime DESC limit 1
"""
List<Object> taskRow = waitMTMVTaskTerminal(showTasks, "waitingMTMVTaskFinishedWithoutAnalyze")
String status = taskRow == null ? "NULL" : taskRow.get(1).toString()
if (status != "SUCCESS") {
logger.info("status is not success")
}
Assert.assertEquals("SUCCESS", status)
}

void waitingMTMVTaskFinishedNotNeedSuccess(String jobName) {
// Wait for the newly submitted MTMV task to become visible in tasks().
Thread.sleep(2000);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -345,12 +345,12 @@ suite("cross_join_list_str_increment_create") {

def sql_all_list = [mv_sql_1, mv_sql_3, mv_sql_4, mv_sql_6, mv_sql_7, mv_sql_8, mv_sql_9, mv_sql_10, mv_sql_11, mv_sql_12,
mv_sql_13, mv_sql_14, mv_sql_15, mv_sql_16, mv_sql_17, mv_sql_18]
def sql_increment_list = [mv_sql_1, mv_sql_3, mv_sql_4, mv_sql_6, mv_sql_8]
def sql_increment_list = [mv_sql_1, mv_sql_3, mv_sql_4, mv_sql_6, mv_sql_8, mv_sql_14]
def sql_complete_list = []

// change left table data
// create mv base on left table with partition col
def sql_error_list = [mv_sql_7, mv_sql_9, mv_sql_10, mv_sql_11, mv_sql_12, mv_sql_13, mv_sql_14, mv_sql_15, mv_sql_16, mv_sql_17, mv_sql_18]
def sql_error_list = [mv_sql_7, mv_sql_9, mv_sql_10, mv_sql_11, mv_sql_12, mv_sql_13, mv_sql_15, mv_sql_16, mv_sql_17, mv_sql_18]
list_judgement(sql_all_list, sql_increment_list, sql_complete_list, sql_error_list,
partition_by_part_col, primary_tb_change, is_complete_change)

Expand All @@ -375,9 +375,9 @@ suite("cross_join_list_str_increment_create") {

// change right table data
// create mv base on left table with partition col
sql_error_list = [mv_sql_7, mv_sql_9, mv_sql_10, mv_sql_11, mv_sql_12, mv_sql_13, mv_sql_14, mv_sql_15, mv_sql_16, mv_sql_17, mv_sql_18]
sql_error_list = [mv_sql_7, mv_sql_9, mv_sql_10, mv_sql_11, mv_sql_12, mv_sql_13, mv_sql_15, mv_sql_16, mv_sql_17, mv_sql_18]
sql_increment_list = []
sql_complete_list = [mv_sql_1, mv_sql_3, mv_sql_4, mv_sql_6, mv_sql_8]
sql_complete_list = [mv_sql_1, mv_sql_3, mv_sql_4, mv_sql_6, mv_sql_8, mv_sql_14]
list_judgement(sql_all_list, sql_increment_list, sql_complete_list, sql_error_list, partition_by_part_col, slave_tb_change, is_complete_change)

// create mv base on left table with no partition col
Expand Down
Loading
Loading