yujun777 commented on code in PR #68648:
URL: https://github.com/apache/doris/pull/68648#discussion_r4225092137
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/UpdateMvByPartitionCommand.java:
##########
@@ -160,30 +218,218 @@ public static Set<Expression>
constructPredicates(Set<PartitionItem> partitions,
*/
@VisibleForTesting
public static Set<Expression> constructPredicates(Set<PartitionItem>
partitions, Slot colSlot) {
+ return constructPredicates(partitions, colSlot, Optional.empty());
+ }
+
+ private static Set<Expression> constructPredicates(Set<PartitionItem>
partitions, Slot colSlot,
+ Optional<Type> columnType) {
Set<Expression> predicates = new HashSet<>();
if (partitions.isEmpty()) {
return Sets.newHashSet(BooleanLiteral.TRUE);
}
if (partitions.iterator().next() instanceof ListPartitionItem) {
for (PartitionItem item : partitions) {
- predicates.add(convertListPartitionToIn(item, colSlot));
+ predicates.add(convertListPartitionToIn(item, colSlot,
columnType));
}
} else {
for (PartitionItem item : partitions) {
- predicates.add(convertRangePartitionToCompare(item, colSlot));
+ predicates.add(convertRangePartitionToCompare(item, colSlot,
columnType));
+ }
+ }
+ return predicates;
+ }
+
+ /**
+ * The predicate a base table is read through when the refresh is to read
exactly these partitions of it.
+ *
+ * <p>A partition of a list partitioned table holds one key per partition
column, and the column the MV
+ * partition is named by is only one of them. A predicate on that column
alone also reaches the
+ * partitions whose other keys differ -- a table partitioned by (d,
region) has one partition of
+ * (d0, 'US') and one of (d0, 'EU'), and `d = d0` reaches both, while only
the second is a partition
+ * this refresh is to read; a later drop of the first would then leave its
rows in the MV partition
+ * while the snapshot, which names only the second, still calls it
synchronized. So a list partition is
+ * pinned to its whole key. A range partition is pinned to its bounds,
which is the same thing: a base
+ * table partitioned by range has a single partition column, see
+ * {@code RangePartitionItem#toPartitionKeyDesc(int)}.
+ *
+ * <p>The partitions are never empty: a table the caller scopes with no
partition is read as nothing
+ * before this is reached, see {@code constructTableWithPredicates}.
+ */
+ private static Set<Expression>
constructPredicatesOfBasePartitions(Set<PartitionItem> partitions,
+ OlapTable baseTable, String colName) {
+ return constructPredicatesOfBasePartitions(partitions, baseTable,
colName, UnboundSlot::new);
+ }
+
+ /**
+ * The same, with the partition columns read through the slots the caller
names them by. A caller that
+ * writes these predicates into a plan that is already bound and is never
bound again -- the union
+ * compensation -- has to pass that plan's own slots: an unbound slot
there fails the rewrite instead of
+ * narrowing the read.
+ */
+ private static Set<Expression>
constructPredicatesOfBasePartitions(Set<PartitionItem> partitions,
+ OlapTable baseTable, String colName, Function<String, Slot>
slotOfColumn) {
+ List<Column> partitionColumns = baseTable.getPartitionColumns();
+ List<Type> partitionColumnTypes = Lists.transform(partitionColumns,
Column::getType);
+ if (!(partitions.iterator().next() instanceof ListPartitionItem)) {
+ Set<Expression> predicates = new HashSet<>();
+ for (PartitionItem item : partitions) {
+ predicates.add(convertRangePartitionToCompare(item,
slotOfColumn.apply(colName),
+ Optional.of(partitionColumnTypes.get(0))));
}
+ return predicates;
+ }
+ List<Slot> partitionSlots = Lists.newArrayList();
+ for (Column partitionColumn : partitionColumns) {
+ partitionSlots.add(slotOfColumn.apply(partitionColumn.getName()));
+ }
+ Set<Expression> predicates = new HashSet<>();
+ for (PartitionItem item : partitions) {
+ predicates.add(convertListPartitionToKey(item, partitionSlots,
partitionColumnTypes));
}
return predicates;
}
- private static Expression convertPartitionKeyToLiteral(PartitionKey key) {
- return Literal.fromLegacyLiteral(key.getKeys().get(0),
- Type.fromPrimitiveType(key.getTypes().get(0)));
+ /**
+ * The predicate for a table whose mapped partitions include a list
partition's default one.
+ *
+ * <p>That partition takes the rows no other partition of the table
claims, and an MV partition takes the
+ * rows whose own key falls in it wherever the base table placed them --
so those rows are the ones the MV
+ * partition's own key range names, and no predicate on the partition
columns picks them out on its own: a
+ * test on those columns reaches the rows an explicit partition holds as
readily as it reaches theirs.
+ * The read is therefore written the other way round, as the record's own
partitions plus what the key
+ * range holds beyond them:
+ *
+ * <ul>
+ * <li>the partitions this MV partition is recorded with, pinned by
their whole key, and</li>
+ * <li>the rows of the MV partitions' key range that no explicit
partition the mapping does not name
+ * holds, which is the default partition's rows in that range and
nothing else.</li>
+ * </ul>
+ *
+ * <p>An explicit partition the window left out shares the MV partition's
key range with a retained one --
+ * the shape this scope is about, a list partitioned table's partitions
holding several keys of the MV's
+ * partition column -- and reading it back would put its rows in the MV
partition while no snapshot names
+ * it, leaving them there through a later drop of it. Subtracting the
partitions the mapping does not name
+ * is what keeps the read to the record; the partitions it does name are
added back by the first part
+ * whether or not their keys fall in the range, the way every other scoped
read pins them.
+ */
+ private static Set<Expression>
constructPredicatesOfADefaultPartitionTable(Set<PartitionItem> mvItems,
+ Set<PartitionItem> mappedItems, Set<String> readable, OlapTable
baseTable, String colName)
+ throws AnalysisException {
+ List<Column> partitionColumns = baseTable.getPartitionColumns();
+ List<Type> partitionColumnTypes = Lists.transform(partitionColumns,
Column::getType);
+ List<Slot> partitionSlots = Lists.newArrayList();
+ for (Column partitionColumn : partitionColumns) {
+ partitionSlots.add(new UnboundSlot(partitionColumn.getName()));
+ }
+ Set<PartitionItem> notMapped = Sets.newHashSet();
+ for (String partitionName : baseTable.getPartitionNames()) {
+ if (readable.contains(partitionName)) {
+ continue;
+ }
+ PartitionItem item =
baseTable.getPartitionItemOrAnalysisException(partitionName);
+ if (!item.isDefaultPartition()) {
+ notMapped.add(item);
+ }
+ }
+ Expression inRange = ExpressionUtils.or(constructPredicates(mvItems,
colName,
+ Optional.of(partitionColumnType(baseTable, colName))));
+ if (!notMapped.isEmpty()) {
+ inRange = ExpressionUtils.and(inRange, new Not(ExpressionUtils.or(
Review Comment:
Confirmed, and the direction it killed is gone in 2caf7c584b7.
Your scenario is right: `ADD PARTITION` makes an empty partition for a key
and does not move the rows the default partition already holds under it, so a
row can sit in `p_default` under a key an explicit partition claims -- and one
the window dropped. The key test cannot tell the two sides apart, so the
exclusion removed a row the snapshot names and a COMPLETE refresh dropped
committed data while calling the partition synchronized.
So the exclusion is gone, and the record is what is made to match the read
instead: a table whose partitions include a list partition's default one is
left out of the window, and the mapping names every partition the MV
partitions' key ranges cover, which is what the read takes. That needed one
more piece, because leaving the table unwindowed puts two base partitions
describing keys that meet into the MV's own partition set -- it cannot have two
partitions repeating a key, and the same shape with no `partition_sync_limit`
failed to create at all before this:
`MTMVRelatedPartitionDescTransferGenerator` merges descs whose keys meet into
the one partition whose keys are their union (the partition the windowed
pipeline already produced for such a shape), and leaves disjoint descs and
non-list descs exactly as they were.
Measured on your shape -- the row inserted as `(2020-01-01, 1, 1)` before
`ADD PARTITION p_expired VALUES IN ((2020-01-01, 1))`, an MV `PARTITION BY (d)`
with a 2 YEAR window: the MV partition holds that row, and inserting into
`p_expired` afterwards leaves the MV partition not synchronized, since it is a
partition the MV partition is recorded with. It is `list_default_scope` in
`test_mtmv_base_partition_read_scope`, and
`MTMVRelatedPartitionDescGeneratorTest` pins the merge.
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/UpdateMvByPartitionCommand.java:
##########
@@ -160,30 +218,218 @@ public static Set<Expression>
constructPredicates(Set<PartitionItem> partitions,
*/
@VisibleForTesting
public static Set<Expression> constructPredicates(Set<PartitionItem>
partitions, Slot colSlot) {
+ return constructPredicates(partitions, colSlot, Optional.empty());
+ }
+
+ private static Set<Expression> constructPredicates(Set<PartitionItem>
partitions, Slot colSlot,
+ Optional<Type> columnType) {
Set<Expression> predicates = new HashSet<>();
if (partitions.isEmpty()) {
return Sets.newHashSet(BooleanLiteral.TRUE);
}
if (partitions.iterator().next() instanceof ListPartitionItem) {
for (PartitionItem item : partitions) {
- predicates.add(convertListPartitionToIn(item, colSlot));
+ predicates.add(convertListPartitionToIn(item, colSlot,
columnType));
}
} else {
for (PartitionItem item : partitions) {
- predicates.add(convertRangePartitionToCompare(item, colSlot));
+ predicates.add(convertRangePartitionToCompare(item, colSlot,
columnType));
+ }
+ }
+ return predicates;
+ }
+
+ /**
+ * The predicate a base table is read through when the refresh is to read
exactly these partitions of it.
+ *
+ * <p>A partition of a list partitioned table holds one key per partition
column, and the column the MV
+ * partition is named by is only one of them. A predicate on that column
alone also reaches the
+ * partitions whose other keys differ -- a table partitioned by (d,
region) has one partition of
+ * (d0, 'US') and one of (d0, 'EU'), and `d = d0` reaches both, while only
the second is a partition
+ * this refresh is to read; a later drop of the first would then leave its
rows in the MV partition
+ * while the snapshot, which names only the second, still calls it
synchronized. So a list partition is
+ * pinned to its whole key. A range partition is pinned to its bounds,
which is the same thing: a base
+ * table partitioned by range has a single partition column, see
+ * {@code RangePartitionItem#toPartitionKeyDesc(int)}.
+ *
+ * <p>The partitions are never empty: a table the caller scopes with no
partition is read as nothing
+ * before this is reached, see {@code constructTableWithPredicates}.
+ */
+ private static Set<Expression>
constructPredicatesOfBasePartitions(Set<PartitionItem> partitions,
+ OlapTable baseTable, String colName) {
+ return constructPredicatesOfBasePartitions(partitions, baseTable,
colName, UnboundSlot::new);
+ }
+
+ /**
+ * The same, with the partition columns read through the slots the caller
names them by. A caller that
+ * writes these predicates into a plan that is already bound and is never
bound again -- the union
+ * compensation -- has to pass that plan's own slots: an unbound slot
there fails the rewrite instead of
+ * narrowing the read.
+ */
+ private static Set<Expression>
constructPredicatesOfBasePartitions(Set<PartitionItem> partitions,
+ OlapTable baseTable, String colName, Function<String, Slot>
slotOfColumn) {
+ List<Column> partitionColumns = baseTable.getPartitionColumns();
+ List<Type> partitionColumnTypes = Lists.transform(partitionColumns,
Column::getType);
+ if (!(partitions.iterator().next() instanceof ListPartitionItem)) {
+ Set<Expression> predicates = new HashSet<>();
+ for (PartitionItem item : partitions) {
+ predicates.add(convertRangePartitionToCompare(item,
slotOfColumn.apply(colName),
+ Optional.of(partitionColumnTypes.get(0))));
}
+ return predicates;
+ }
+ List<Slot> partitionSlots = Lists.newArrayList();
+ for (Column partitionColumn : partitionColumns) {
+ partitionSlots.add(slotOfColumn.apply(partitionColumn.getName()));
+ }
+ Set<Expression> predicates = new HashSet<>();
+ for (PartitionItem item : partitions) {
+ predicates.add(convertListPartitionToKey(item, partitionSlots,
partitionColumnTypes));
}
return predicates;
}
- private static Expression convertPartitionKeyToLiteral(PartitionKey key) {
- return Literal.fromLegacyLiteral(key.getKeys().get(0),
- Type.fromPrimitiveType(key.getTypes().get(0)));
+ /**
+ * The predicate for a table whose mapped partitions include a list
partition's default one.
+ *
+ * <p>That partition takes the rows no other partition of the table
claims, and an MV partition takes the
+ * rows whose own key falls in it wherever the base table placed them --
so those rows are the ones the MV
+ * partition's own key range names, and no predicate on the partition
columns picks them out on its own: a
+ * test on those columns reaches the rows an explicit partition holds as
readily as it reaches theirs.
+ * The read is therefore written the other way round, as the record's own
partitions plus what the key
+ * range holds beyond them:
+ *
+ * <ul>
+ * <li>the partitions this MV partition is recorded with, pinned by
their whole key, and</li>
+ * <li>the rows of the MV partitions' key range that no explicit
partition the mapping does not name
+ * holds, which is the default partition's rows in that range and
nothing else.</li>
+ * </ul>
+ *
+ * <p>An explicit partition the window left out shares the MV partition's
key range with a retained one --
+ * the shape this scope is about, a list partitioned table's partitions
holding several keys of the MV's
+ * partition column -- and reading it back would put its rows in the MV
partition while no snapshot names
+ * it, leaving them there through a later drop of it. Subtracting the
partitions the mapping does not name
+ * is what keeps the read to the record; the partitions it does name are
added back by the first part
+ * whether or not their keys fall in the range, the way every other scoped
read pins them.
+ */
+ private static Set<Expression>
constructPredicatesOfADefaultPartitionTable(Set<PartitionItem> mvItems,
+ Set<PartitionItem> mappedItems, Set<String> readable, OlapTable
baseTable, String colName)
+ throws AnalysisException {
+ List<Column> partitionColumns = baseTable.getPartitionColumns();
+ List<Type> partitionColumnTypes = Lists.transform(partitionColumns,
Column::getType);
+ List<Slot> partitionSlots = Lists.newArrayList();
+ for (Column partitionColumn : partitionColumns) {
+ partitionSlots.add(new UnboundSlot(partitionColumn.getName()));
+ }
+ Set<PartitionItem> notMapped = Sets.newHashSet();
Review Comment:
Gone in 2caf7c584b7 with the exclusion it was about: the read of such a
table is not narrowed per key any more, so there is no `NOT(OR(...))` term per
batch and nothing quadratic to prune. The MV partition's key range is the whole
predicate for these tables.
##########
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()) {
Review Comment:
Fixed in 2caf7c584b7, along the first line you suggest:
`MTMV.defaultListPartitionsOf` now holds the table's read lock over the name
walk and the item lookup, so a concurrent `ADD`/`DROP PARTITION` cannot be seen
half applied and this walk is no longer one more reader of a tree another
thread is modifying.
I did not move the fact into the partition-item copy the generators take,
which is the other half of what you point at: the walk and that copy can still
be taken at two different moments, so the mapping can carry a default partition
that the descs did not come from for one refresh. That is the same one-refresh
window the mapping already documents for partition sync
(`generateRelatedBasePartitionIds`: "A base partition that partition sync adds
or drops at the same moment can therefore sit outside the returned set for one
refresh"), and it fails the refresh loudly rather than serving anything, since
the mapping would name a partition the table no longer has. Say the word if you
want it removed rather than bounded and I will thread the fact through
`RelatedPartitionDescResult` from the init generator, which is where the locked
copy already is.
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/UpdateMvByPartitionCommand.java:
##########
@@ -160,30 +218,218 @@ public static Set<Expression>
constructPredicates(Set<PartitionItem> partitions,
*/
@VisibleForTesting
public static Set<Expression> constructPredicates(Set<PartitionItem>
partitions, Slot colSlot) {
+ return constructPredicates(partitions, colSlot, Optional.empty());
+ }
+
+ private static Set<Expression> constructPredicates(Set<PartitionItem>
partitions, Slot colSlot,
+ Optional<Type> columnType) {
Set<Expression> predicates = new HashSet<>();
if (partitions.isEmpty()) {
return Sets.newHashSet(BooleanLiteral.TRUE);
}
if (partitions.iterator().next() instanceof ListPartitionItem) {
for (PartitionItem item : partitions) {
- predicates.add(convertListPartitionToIn(item, colSlot));
+ predicates.add(convertListPartitionToIn(item, colSlot,
columnType));
}
} else {
for (PartitionItem item : partitions) {
- predicates.add(convertRangePartitionToCompare(item, colSlot));
+ predicates.add(convertRangePartitionToCompare(item, colSlot,
columnType));
+ }
+ }
+ return predicates;
+ }
+
+ /**
+ * The predicate a base table is read through when the refresh is to read
exactly these partitions of it.
+ *
+ * <p>A partition of a list partitioned table holds one key per partition
column, and the column the MV
+ * partition is named by is only one of them. A predicate on that column
alone also reaches the
+ * partitions whose other keys differ -- a table partitioned by (d,
region) has one partition of
+ * (d0, 'US') and one of (d0, 'EU'), and `d = d0` reaches both, while only
the second is a partition
+ * this refresh is to read; a later drop of the first would then leave its
rows in the MV partition
+ * while the snapshot, which names only the second, still calls it
synchronized. So a list partition is
+ * pinned to its whole key. A range partition is pinned to its bounds,
which is the same thing: a base
+ * table partitioned by range has a single partition column, see
+ * {@code RangePartitionItem#toPartitionKeyDesc(int)}.
+ *
+ * <p>The partitions are never empty: a table the caller scopes with no
partition is read as nothing
+ * before this is reached, see {@code constructTableWithPredicates}.
+ */
+ private static Set<Expression>
constructPredicatesOfBasePartitions(Set<PartitionItem> partitions,
+ OlapTable baseTable, String colName) {
+ return constructPredicatesOfBasePartitions(partitions, baseTable,
colName, UnboundSlot::new);
+ }
+
+ /**
+ * The same, with the partition columns read through the slots the caller
names them by. A caller that
+ * writes these predicates into a plan that is already bound and is never
bound again -- the union
+ * compensation -- has to pass that plan's own slots: an unbound slot
there fails the rewrite instead of
+ * narrowing the read.
+ */
+ private static Set<Expression>
constructPredicatesOfBasePartitions(Set<PartitionItem> partitions,
+ OlapTable baseTable, String colName, Function<String, Slot>
slotOfColumn) {
+ List<Column> partitionColumns = baseTable.getPartitionColumns();
+ List<Type> partitionColumnTypes = Lists.transform(partitionColumns,
Column::getType);
+ if (!(partitions.iterator().next() instanceof ListPartitionItem)) {
+ Set<Expression> predicates = new HashSet<>();
+ for (PartitionItem item : partitions) {
+ predicates.add(convertRangePartitionToCompare(item,
slotOfColumn.apply(colName),
+ Optional.of(partitionColumnTypes.get(0))));
}
+ return predicates;
+ }
+ List<Slot> partitionSlots = Lists.newArrayList();
+ for (Column partitionColumn : partitionColumns) {
+ partitionSlots.add(slotOfColumn.apply(partitionColumn.getName()));
+ }
+ Set<Expression> predicates = new HashSet<>();
+ for (PartitionItem item : partitions) {
+ predicates.add(convertListPartitionToKey(item, partitionSlots,
partitionColumnTypes));
}
return predicates;
}
- private static Expression convertPartitionKeyToLiteral(PartitionKey key) {
- return Literal.fromLegacyLiteral(key.getKeys().get(0),
- Type.fromPrimitiveType(key.getTypes().get(0)));
+ /**
+ * The predicate for a table whose mapped partitions include a list
partition's default one.
+ *
+ * <p>That partition takes the rows no other partition of the table
claims, and an MV partition takes the
+ * rows whose own key falls in it wherever the base table placed them --
so those rows are the ones the MV
+ * partition's own key range names, and no predicate on the partition
columns picks them out on its own: a
+ * test on those columns reaches the rows an explicit partition holds as
readily as it reaches theirs.
+ * The read is therefore written the other way round, as the record's own
partitions plus what the key
+ * range holds beyond them:
+ *
+ * <ul>
+ * <li>the partitions this MV partition is recorded with, pinned by
their whole key, and</li>
+ * <li>the rows of the MV partitions' key range that no explicit
partition the mapping does not name
+ * holds, which is the default partition's rows in that range and
nothing else.</li>
+ * </ul>
+ *
+ * <p>An explicit partition the window left out shares the MV partition's
key range with a retained one --
+ * the shape this scope is about, a list partitioned table's partitions
holding several keys of the MV's
+ * partition column -- and reading it back would put its rows in the MV
partition while no snapshot names
+ * it, leaving them there through a later drop of it. Subtracting the
partitions the mapping does not name
+ * is what keeps the read to the record; the partitions it does name are
added back by the first part
+ * whether or not their keys fall in the range, the way every other scoped
read pins them.
+ */
+ private static Set<Expression>
constructPredicatesOfADefaultPartitionTable(Set<PartitionItem> mvItems,
+ Set<PartitionItem> mappedItems, Set<String> readable, OlapTable
baseTable, String colName)
+ throws AnalysisException {
+ List<Column> partitionColumns = baseTable.getPartitionColumns();
+ List<Type> partitionColumnTypes = Lists.transform(partitionColumns,
Column::getType);
+ List<Slot> partitionSlots = Lists.newArrayList();
+ for (Column partitionColumn : partitionColumns) {
+ partitionSlots.add(new UnboundSlot(partitionColumn.getName()));
+ }
+ Set<PartitionItem> notMapped = Sets.newHashSet();
+ for (String partitionName : baseTable.getPartitionNames()) {
Review Comment:
Fixed in 2caf7c584b7 by removal, mostly: the exclusion list is gone, so
there is no second metadata read to be consistent with. The mapping now reads
the partition items through the generators' own `getAndCopyPartitionItems`
copies (each of which takes the table's read lock) and the default-partition
name under the table's read lock, which is the one place it is read outside
those copies.
On the second half -- a base partition added after the mapping is built,
whose rows the scan could read while the snapshot does not name it -- that is
bounded the same way for these tables: the read is the MV partitions' key
range, the mapping is recomputed from the current metadata on every comparison,
and a partition it names that the snapshot does not is exactly what the sync
check reports, so the partition is refreshed rather than served. It cannot stay
unnamed, which is what the window did.
--
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]