yujun777 commented on code in PR #68648:
URL: https://github.com/apache/doris/pull/68648#discussion_r4219759316
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/UpdateMvByPartitionCommand.java:
##########
@@ -284,6 +440,55 @@ public Plan
visitLogicalSubQueryAlias(LogicalSubQueryAlias<? extends Plan> subQu
return super.visitLogicalSubQueryAlias(subQueryAlias, predicates);
}
+ /**
+ * The ranges of the MV partitions this compensation removes, as a
predicate on the base table's
+ * partition column, or nothing when they cannot be written on one
column.
+ *
+ * <p>A default partition's rows belong to whichever MV partition
their own key falls in, so this is
+ * what they are read through: pinning them to the partition's own key
-- the sentinel those rows were
+ * placed by -- reads none of them, and reading them whole would add
the rows of the MV partitions the
+ * plan still has. Only an MV partitioned by this one column can be
written that way here.
+ */
+ private static Set<Expression>
mvPartitionsToReadThrough(PredicateAddContext predicates,
+ BaseColInfo relatedTableColumnInfo, Slot partitionSlot) {
+ Set<Expression> res = Sets.newHashSet();
+ if (predicates.getMvPartitionsToRemove().isEmpty()) {
+ return res;
+ }
+ for (Map.Entry<BaseTableInfo, Set<String>> entry
+ : predicates.getMvPartitionsToRemove().entrySet()) {
+ try {
+ TableIf table = MTMVUtil.getTable(entry.getKey());
+ if (!(table instanceof MTMV)) {
+ continue;
+ }
+ MTMV mtmv = (MTMV) table;
+ if (mtmv.getPartitionColumns().size() != 1
+ || !mtmv.getPartitionColumns().get(0).getName()
+
.equalsIgnoreCase(relatedTableColumnInfo.getColName())) {
+ continue;
+ }
+ Type columnType =
mtmv.getPartitionColumns().get(0).getType();
+ for (String partitionName : entry.getValue()) {
+ PartitionItem item =
mtmv.getPartitionItemOrAnalysisException(partitionName);
+ if (!(item instanceof ListPartitionItem)) {
+ continue;
+ }
+ List<Expression> values = ((ListPartitionItem)
item).getItems().stream()
+ .map(key -> convertPartitionKeyToLiteral(key,
0, Optional.of(columnType)))
+ .collect(Collectors.toList());
+ res.add(new InPredicate(partitionSlot, values));
+ }
Review Comment:
Fixed in 0dbc4082600. The default-partition compensation now builds its
predicate with the same NULL-aware list conversion the scoped read uses
(`convertListPartitionToIn`), so a key that is NULL is asked for as `IS NULL`
instead of `k IN (NULL)`: a nullable `LIST(k)` MV partition that is stale as
NULL no longer has the union branch's base read come back empty for it.
The conversion itself is the one `UpdateMvByPartitionCommandTest#testNull`
pins for the scoped read; the end-to-end reach of the compensation branch is
the same one the earlier threads left unverified -- a partly usable MV is
refused by `checkMaterializationPattern` in this build, so what I ran is the
unit tests and the rewrite regression suites, not this case on a live rewrite.
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/UpdateMvByPartitionCommand.java:
##########
@@ -350,10 +555,33 @@ public Plan
visitLogicalCatalogRelation(LogicalCatalogRelation catalogRelation,
predicates.setNeedAddFilter(false);
}
if (!partitionHasDataItems.isEmpty()) {
+ // The partitions are pinned the way a refresh pins
them: the whole key of each, at
+ // the partition column's own type. A predicate on the
MV's partition column alone
+ // reads the partitions that differ in the other keys
too, and those rows are ones
+ // this branch must not add -- the MV branch of the
union already supplies them, or
+ // the compensation would count them twice -- while a
key written with a scale has to
+ // be compared at that scale or its rows are read as
none.
+ // A list partitioned table's default partition is not
pinned the way the others
+ // are: its own key is the sentinel the rows no other
partition claims were placed
+ // by, so pinning it to that key reads none of them.
It is read the way it was
+ // before the whole key was pinned -- on the MV's
partition column alone -- which
+ // can be seen to be too wide rather than one that
drops its rows.
+ boolean hasDefaultPartition =
partitionHasDataItems.stream()
+ .anyMatch(PartitionItem::isDefaultPartition);
+ Set<Expression> mvPartitionPredicates =
hasDefaultPartition
+ ? mvPartitionsToReadThrough(predicates,
relatedTableColumnInfo, partitionSlot)
+ : Sets.newHashSet();
+ Set<Expression> preds;
+ if (!mvPartitionPredicates.isEmpty()) {
+ preds = mvPartitionPredicates;
+ } else if (targetTable instanceof OlapTable &&
!hasDefaultPartition) {
+ preds =
constructPredicatesOfBasePartitions(partitionHasDataItems,
+ (OlapTable) targetTable,
relatedTableColumnInfo.getColName());
Review Comment:
Fixed in 0dbc4082600, in the direction you point at. The whole-key predicate
is now written against the slots the relation itself reads
(`partitionColumnSlots`, from `catalogRelation.getOutput()`), which is what the
compensation used before it grew the whole-key form -- so the filter no longer
fails a plan that nothing binds again.
One case needed a decision: a partition column the relation does not read
has no slot to write the key against, and pinning the read to the MV's
partition column alone is what the whole-key form replaced (it reads the
partitions differing in the other keys, whose rows the MV branch of the union
already supplies). That case is refused -- `handleSuccess = false`, the
candidate is dropped and the query is answered from the base table -- rather
than narrowed to a guess.
`UpdateMvByPartitionCommandTest#testACompensationFilterIsWrittenWithTheRelationSlots`
pins the bound-slot requirement: it fails on an unbound slot, which is how
this read as a failed candidate rather than as a narrowed filter.
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/RefreshMTMVCommand.java:
##########
@@ -198,8 +198,11 @@ protected LogicalPlan createRefreshCommand(MTMV mtmv,
StatementContext statement
statementContext.setIvmRewriteContext(Optional.of(IvmRewriteContext.fullExplain(mtmv)));
}
statementContext.setExcludedTriggerTables(mtmv.getExcludedTriggerTables());
+ // Explained through the MV partitions' own key ranges: the
read a refresh narrows to its
+ // partition mapping is decided by that refresh's context,
which a plan built here has not.
return UpdateMvByPartitionCommand.from(
- mtmv, getCompleteRefreshPartitions(mtmv),
getIncrementalTableMap(mtmv), statementContext);
+ mtmv, getCompleteRefreshPartitions(mtmv),
getIncrementalTableMap(mtmv), statementContext,
+ null);
Review Comment:
Fixed in 0dbc4082600. `EXPLAIN REFRESH ... COMPLETE` now builds the same
base-partition scope the refresh reads, from the same mapping, so an explain of
a windowed MV shows the partitions the refresh reads rather than the MV
partitions' own (wider) key ranges. An IVM MV is passed no scope, as in the
refresh.
The scope helper is now shared -- `MTMVPartitionUtil#mappedBasePartitions`,
which `MTMVTask` also calls -- so the explain and the refresh cannot drift
apart again: that is what made this a P2 worth taking rather than a documented
difference.
##########
fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java:
##########
@@ -1608,9 +1608,11 @@ public Map<String, Map<MTMVRelatedTableIf, Set<String>>>
calculatePartitionMappi
Map<PartitionKeyDesc, Map<MTMVRelatedTableIf, Set<String>>>
pctPartitionDescs = MTMVPartitionUtil
.generateRelatedPartitionDescs(mvPartitionInfo, mvProperties,
getPartitionColumns(),
effectiveFilter, pinnedSnapshots);
+ Map<MTMVRelatedTableIf, String> defaultListPartitions =
defaultListPartitionsOf();
for (Entry<String, PartitionItem> entry : mvPartitionItems.entrySet())
{
- res.put(entry.getKey(),
-
pctPartitionDescs.getOrDefault(entry.getValue().toPartitionKeyDesc(),
Maps.newHashMap()));
+ res.put(entry.getKey(), withDefaultListPartitions(
Review Comment:
Fixed in 0dbc4082600, on the side that cannot drop rows. Preserving the
whole inverse set made every MV partition that reads a base partition a
candidate, but the eligibility still answered a query when only *some* of those
partitions were valid. With `enable_materialized_view_union_rewrite` off there
is nothing to compensate the others: the MV alone answers the query and reads
the partitions that are not valid as they are.
`getMTMVCanRewritePartitions` now returns no partition for a query unless
every MV partition its base partitions are mapped from is valid, when the union
rewrite is off. Read as a subset rather than as a count, since a partition
inside its grace period is answered usable without having been compared and is
not necessarily one the query reads. With the union rewrite on (the default) a
partially valid answer is still a complete one, and the valid partitions still
answer.
##########
fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRewriteUtil.java:
##########
@@ -166,25 +166,35 @@ private static Set<String>
getMtmvPartitionsByRelatedPartitions(MTMV mtmv, MTMVR
}
Set<String> pctPartitions = entry.getValue();
for (String pctPartition : pctPartitions) {
- String mvPartition = relatedToMv.get(Pair.of(tableIf,
pctPartition));
- if (mvPartition != null) {
- res.add(mvPartition);
+ Set<String> mvPartitions = relatedToMv.get(Pair.of(tableIf,
pctPartition));
+ if (mvPartitions != null) {
Review Comment:
Fixed in 0dbc4082600. You are right that the grace period was the way around
the guard: a partition inside its grace period was answered usable and
`continue`d before `getMtmvPartitionsByRelatedPartitions` was ever called, so
with every partition in grace the mapping was never read and the "no MV
partition is mapped from a base partition the query reads" answer never came.
The mapping is now read before the grace period answers, and an empty answer
returns no partition at all -- not the grace-period partitions accumulated so
far, which would still have answered the query.
One consequence to be straight about: a query against a recently refreshed
MV partition now builds the refresh context, which the grace period used to
skip. That is the cost of the check; the guard it is for is the one this thread
and a13f2939f3a are about.
##########
fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java:
##########
@@ -1619,6 +1621,58 @@ public Map<String, Map<MTMVRelatedTableIf, Set<String>>>
calculatePartitionMappi
return res;
}
+ /**
+ * The list partition each base table of this MV has that takes the rows
no other partition of it claims,
+ * by table, or none for a table that has no such partition.
+ *
+ * <p>Read once per mapping rather than per MV partition: the mapping
describes every MV partition and the
+ * answer is the table's, not the partition's. The partition metadata is
read without a lock, like the
+ * rest of the mapping this is part of.
+ */
+ private Map<MTMVRelatedTableIf, String> defaultListPartitionsOf() throws
AnalysisException {
+ Map<MTMVRelatedTableIf, String> res = Maps.newHashMap();
+ for (MTMVRelatedTableIf pctTable : mvPartitionInfo.getPctTables()) {
+ if (!(pctTable instanceof OlapTable)) {
+ continue;
+ }
+ OlapTable olapTable = (OlapTable) pctTable;
+ if (!(olapTable.getPartitionInfo() instanceof ListPartitionInfo)) {
+ continue;
+ }
+ for (String partitionName : olapTable.getPartitionNames()) {
+ if
(olapTable.getPartitionItemOrAnalysisException(partitionName).isDefaultPartition())
{
+ res.put(pctTable, partitionName);
+ break;
+ }
+ }
+ }
+ return res;
+ }
+
+ /**
+ * One MV partition's mapping, with every base table's default list
partition named in it.
+ *
+ * <p>Such a partition holds rows for every key its table can be read by,
so it belongs to every MV
+ * partition that reads the table -- not only to the one its own key, the
sentinel those rows were placed
+ * by, maps to. Naming it everywhere is what the read and the record have
to agree on: the refresh reads
+ * the rows of it that belong to the MV partition being refreshed, and the
partition is recorded among the
+ * ones that partition is read through, so an insert into it leaves that
MV partition out of sync instead
+ * of changing nothing the MV compares.
+ */
+ private Map<MTMVRelatedTableIf, Set<String>> withDefaultListPartitions(
+ Map<MTMVRelatedTableIf, Set<String>> mapping,
Map<MTMVRelatedTableIf, String> defaultListPartitions) {
+ if (defaultListPartitions.isEmpty()) {
+ return mapping;
+ }
+ Map<MTMVRelatedTableIf, Set<String>> res = Maps.newHashMap(mapping);
+ for (Entry<MTMVRelatedTableIf, String> entry :
defaultListPartitions.entrySet()) {
+ Set<String> partitions =
Sets.newHashSet(res.getOrDefault(entry.getKey(), Sets.newHashSet()));
+ partitions.add(entry.getValue());
Review Comment:
Partly fixed in 0dbc4082600, and one part still deliberately left out.
What is fixed is the case where the base key is not the MV's own column
name: `mvPartitionsToReadThrough` identified the column by comparing the MV
partition column's name with the base column's, so an MV that aliases it (the
`l.k AS k` shape) was skipped and its default partition was read through the
partition's own sentinel key -- none of its rows. The table is now identified
by the base table the MV's partition info names, so the aliased shape goes
through the MV partition's ranges like the unaliased one.
What is still out is the multi-column MV (`getPartitionColumns().size() >
1`), as in my earlier reply: the key position of this base column among the
MV's partition columns is not something the partition info carries, and writing
it at the wrong position would narrow the read to the wrong values. That is
still the next step, and it is why the guard is kept rather than guessed.
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/UpdateMvByPartitionCommand.java:
##########
@@ -130,17 +139,66 @@ private static List<String>
constructPartsForMv(Set<String> partitionNames) {
return Lists.newArrayList(partitionNames);
}
+ /**
+ * The predicate every base table of the MV definition is read through.
+ *
+ * <p>A table the caller scopes is read from exactly the base partitions
it named. Those are the ones
+ * the refresh is about to record as this MV partition's, and the read is
what has to match the record:
+ * reading the MV partition's own key range instead also reads base
partitions no snapshot describes,
+ * and a later silent change to one of them -- dropped, with the base
partition set back to what it
+ * was -- leaves the rows it put in this MV partition behind while the
partition is still judged
+ * synchronized, so the transparent rewrite serves them and no refresh
plans it again.
+ *
+ * <p>Every other table keeps the MV partition's own key range, which is
what the tables the caller
+ * does not scope were always read through. Scoped tables are olap ones;
the partition names are
+ * looked up on one, see the caller.
+ */
private static Map<TableIf, Set<Expression>>
constructTableWithPredicates(MTMV mv,
- Set<String> partitionNames, Map<TableIf, String> tableWithPartKey)
throws AnalysisException {
- Set<PartitionItem> items = Sets.newHashSet();
+ Set<String> partitionNames, Map<TableIf, String> tableWithPartKey,
+ Map<BaseTableInfo, Set<String>> readableBasePartitions) throws
AnalysisException {
+ Set<PartitionItem> mvItems = Sets.newHashSet();
for (String partitionName : partitionNames) {
- PartitionItem partitionItem =
mv.getPartitionItemOrAnalysisException(partitionName);
- items.add(partitionItem);
+ mvItems.add(mv.getPartitionItemOrAnalysisException(partitionName));
}
ImmutableMap.Builder<TableIf, Set<Expression>> builder = new
ImmutableMap.Builder<>();
- tableWithPartKey.forEach((table, colName) ->
- builder.put(table, constructPredicates(items, colName))
- );
+ for (Map.Entry<TableIf, String> entry : tableWithPartKey.entrySet()) {
+ TableIf table = entry.getKey();
+ String colName = entry.getValue();
+ Set<String> readable = readableBasePartitions == null ? null
+ : readableBasePartitions.get(new BaseTableInfo(table));
+ if (readable == null) {
+ builder.put(table, constructPredicates(mvItems, colName));
+ continue;
+ }
+ OlapTable olapTable = (OlapTable) table;
+ Set<PartitionItem> items = Sets.newHashSet();
+ for (String partitionName : readable) {
+
items.add(olapTable.getPartitionItemOrAnalysisException(partitionName));
+ }
+ if (items.stream().anyMatch(PartitionItem::isDefaultPartition)) {
+ // One of the partitions this MV partition is recorded with is
a list partitioned table's
+ // default partition, which takes the rows no other partition
of it claims. Those rows are
+ // the ones the MV partition's own key range names, wherever
the base table put them, and a
+ // partition of the MV takes them by that key rather than by
the partition they were placed
+ // in. So a table whose mapped partitions include one is read
the way an unscoped one is:
+ // the MV partition's key range, at the partition column's own
type. That read can be seen to
+ // be too wide -- it is the one this scope exists to narrow --
rather than one that drops
+ // rows belonging to the MV partition being refreshed. The
mapping names the default
+ // partition in every MV partition that reads the table, so
this is reached for each of them
+ // and not only for the one the sentinel key maps to.
+ builder.put(table, constructPredicates(mvItems, colName,
Review Comment:
Not fixed yet, and I want to be straight about why, because the direction I
took does not work.
I tried the direction that keeps the read and the record one set by taking
such a table out of the window altogether (a table whose list partitions have a
default one cannot be windowed for the same reason the read cannot be narrowed:
the default partition's rows and an expired explicit partition's rows are not
told apart by a predicate on the partition columns). It does not hold up: with
the table out of the window, an MV of this shape cannot be built at all. On
this tree:
```
CREATE TABLE t (d DATE NOT NULL, k INT NOT NULL, amount BIGINT) DUPLICATE
KEY(d, k)
PARTITION BY LIST(d, k) (
PARTITION p_expired VALUES IN (("2020-01-01", 1)),
PARTITION p_kept VALUES IN (("2020-01-01", 2), ("2038-01-01", 2)),
PARTITION p_default)
...
CREATE MATERIALIZED VIEW mv BUILD IMMEDIATE REFRESH COMPLETE ON MANUAL
PARTITION BY (d)
PROPERTIES("partition_sync_limit" = "2", "partition_sync_time_unit" =
"YEAR") AS ...
```
fails with
`Invalid list value format: The partition key[('2020-01-01')] in partition
item[('2020-01-01')] is conflict with current
partitionKeys[(('2020-01-01'),('2038-01-01'))]`.
The desc pipeline turns the two base partitions into two MV partition items
whose keys at the MV's column overlap -- `p_expired`'s `d` is inside `p_kept`'s
`d` set. The same table and MV with **no** `partition_sync_limit` fail
identically, so the overlap is not something this PR's scope created: the
window is what currently hides it, by removing the expired partition before the
MV's partitions are built.
So there are two ways to close this, and I would rather ask than guess:
1. Keep the window and keep the MV partition shape, and exclude the expired
explicit partitions from the *read*: the default term becomes `mvRange AND
NOT(<the whole keys of the explicit partitions the mapping does not name>)`,
comparing null-safely (`NullSafeEqual`) so a NULL key is not excluded by
three-valued logic. Exact, and it preserves both the window and the partitions;
the cost is one disjunct per base partition inside the MV partition's key range
(a year rollup over days is ~365 of them) and that `NOT` cannot drive partition
pruning.
2. Take the table out of the window and let the desc pipeline merge
overlapping descs into one MV partition (the union of their keys, which is the
partition the windowed pipeline already produces for this shape today). Cheap
at read time and it is the "read and record are one set" answer, but it changes
the partition model for these tables and the merge has to be right on both the
create and the add-partition paths.
I have implemented neither, because (2) without the merge is strictly worse
than the bug -- the MV cannot be created. Tell me which one you want and I will
do it; (1) is the smaller change but (2) is the one that makes the mapping and
the read agree for this table by construction.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]