github-actions[bot] commented on code in PR #67725:
URL: https://github.com/apache/doris/pull/67725#discussion_r4214178158
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/PruneFileScanPartition.java:
##########
@@ -114,21 +120,70 @@ private SelectedPartitions
pruneExternalPartitions(ExternalTable externalTable,
}
Map<String, PartitionItem> nameToPartitionItem =
scan.getSelectedPartitions().selectedPartitions;
+ boolean connectorFilteredPartitions = false;
+ Optional<FilterApplicationResult<ConnectorTableHandle>>
connectorFilterResult = Optional.empty();
+ Optional<MvccSnapshot> snapshot =
ctx.getStatementContext().getSnapshot(externalTable,
+ scan.getTableSnapshot(), scan.getScanParams());
+ if (nameToPartitionItem.isEmpty()
+ && scan.getSelectedPartitions().isDeferredPartitionPruning()
+ && externalTable instanceof PluginDrivenExternalTable
+ && ((PluginDrivenExternalTable)
externalTable).supportsConnectorPartitionPruning()) {
+ ConnectorExpression connectorPredicate =
+
NereidsToConnectorExpressionConverter.convert(filter.getPredicate());
+ if (connectorPredicate != null) {
+
Optional<PluginDrivenExternalTable.ConnectorFilteredPartitionView>
connectorPartitions =
+ ((PluginDrivenExternalTable) externalTable)
+ .applyPartitionFilterForScan(snapshot,
connectorPredicate);
+ if (connectorPartitions.isPresent()) {
+ nameToPartitionItem = connectorPartitions.get().getItems();
+ // Carry the handle the selection was materialized from
into physical planning: applying
+ // the same predicate a second time could observe a
different remote generation and mix
+ // this name set with another handle's partition metadata.
+ connectorFilterResult =
Optional.of(connectorPartitions.get().getFilterResult());
+ connectorFilteredPartitions = true;
+ }
+ }
+ }
+ if (!connectorFilteredPartitions && nameToPartitionItem.isEmpty()
+ && (scan.getSelectedPartitions().isNotPruned()
+ || scan.getSelectedPartitions().isDeferredPartitionPruning()))
{
+ // The plugin-driven scan path degrades instead of failing: a
connector entry that cannot be
+ // represented as a Doris partition item makes the whole view
unavailable, and pruning is then
+ // skipped (every partition is read) rather than pruning against a
partial/empty set.
+ //
+ // The view is resolved ONCE per statement and table reference
(resolveScanPartitionView): the MV
+ // partition collector records the same view, and its compensation
union is restricted to those
+ // names, so a second enumeration here could hand the scan a
generation the compensation never saw.
+ Optional<Map<String, PartitionItem>> scanView =
+ externalTable instanceof PluginDrivenExternalTable
+ ?
ctx.getStatementContext().resolveScanPartitionView(externalTable,
+ scan.getTableSnapshot(),
scan.getScanParams(),
+ () -> ((PluginDrivenExternalTable)
externalTable)
+
.getNameToPartitionItemsForScan(snapshot))
+ :
Optional.of(externalTable.getNameToPartitionItems(snapshot));
+ if (!scanView.isPresent()) {
+ return SelectedPartitions.NOT_PRUNED;
+ }
+ nameToPartitionItem = scanView.get();
+ }
+ final Map<String, PartitionItem> partitionItems = nameToPartitionItem;
Optional<SortedPartitionRanges<String>> sortedPartitionRanges =
Optional.empty();
boolean enableBinarySearch = ctx.getConnectContext() == null
||
ctx.getConnectContext().getSessionVariable().enableBinarySearchFilteringPartitions;
- if (enableBinarySearch && !nameToPartitionItem.isEmpty()) {
- sortedPartitionRanges =
scan.getSelectedPartitions().sortedPartitionRanges
- .or(() -> (Optional)
externalTable.getSortedPartitionRanges(scan))
- .or(() ->
Optional.ofNullable(SortedPartitionRanges.build(nameToPartitionItem)));
+ if (enableBinarySearch && !partitionItems.isEmpty()) {
+ sortedPartitionRanges = connectorFilteredPartitions
+ ?
Optional.ofNullable(SortedPartitionRanges.build(partitionItems))
+ : scan.getSelectedPartitions().sortedPartitionRanges
Review Comment:
[P1] Build fallback ranges from the recorded partition view. For
`Filter(h.year > 2024) -> FileScan(Hive h)`, Hive declines the connector
filter, so lines 157-167 take the statement's full map A. This sorted-cache
lookup re-enumerates the DEFERRED Hive pin to build ranges B; binary-search
pruning returns B's names, and a matching partition added between A and B makes
line 193 abort planning. The cache also re-enumerates for its version token and
can store B's ranges under a later C name set, silently excluding new matching
rows on a cache hit. Build ranges from `partitionItems` (or cache an immutable
map and version from the same view) in this fallback too, and test a mutation
between listings. The existing sorted-range thread addressed only
connector-filtered views.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/mvcc/PluginDrivenMvccExternalTable.java:
##########
@@ -714,9 +689,87 @@ static boolean schemaCacheDisabled(Connector connector) {
@Override
public Map<String, PartitionItem>
getNameToPartitionItems(Optional<MvccSnapshot> snapshot) {
+ if (supportsConnectorPartitionPruning()) {
+ PluginDrivenMvccSnapshot pin = getOrMaterialize(snapshot);
+ if (pin.isPartitionViewMaterialized()) {
+ return pin.getNameToPartitionItem();
+ }
Review Comment:
[P1] Reuse the pre-lock view for query-time MV consumers. With
`enableMaterializedViewRewriteWhenBaseTableUnawareness=true` and an unfiltered
Hive base scan, validating each partitioned Hive-based async MV calls
`getAndCopyPartitionItems`, and union compensation then calls
`getNameToPartitionItems(...).get(name)` inside its loop over K invalid base
partitions. The latest Hive pin is DEFERRED, so this `super` call rebuilds all
N partition items for each candidate and each compensated partition while the
planner holds internal table read locks. A disabled or expired Hive
partition-view cache also repeats full HMS listings. The statement already
preloads one full view; share it with both readers outside the lock window. The
older comment on this line concerned scheduled MTMV refresh, which now
materializes its pin before locking; these are query-time rewrite paths.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/scan/PluginDrivenScanNode.java:
##########
@@ -1160,18 +1166,67 @@ protected TFileAttributes getFileAttributes() throws
UserException {
protected void doFinalize() throws UserException {
scanNodeProperties = null;
cachedPropertiesResult = null;
+ materializeDeferredSelectedPartitions();
// Nereids prunes scan slots between init and finalize; fencing the
init-time table-wide
// tuple would reject old backends even when the executable scan no
longer carries Variant.
checkVariantBackendCompatibilityForCurrentScan(backendPolicy.getBackends());
super.doFinalize();
}
+ private void materializeDeferredSelectedPartitions() throws UserException {
+ // A null selection is the "nothing selected" state this node handles
everywhere else (see
+ // resolveRequiredPartitions, displayPartitionCounts,
shouldUseBatchMode and numApproximateSplits);
+ // there is no deferred view to materialize for it.
+ if (selectedPartitions == null ||
!selectedPartitions.isDeferredPartitionPruning()) {
+ return;
+ }
+ // A logical filter materializes this state earlier in
PruneFileScanPartition. Reaching finalize still
+ // deferred therefore means a no-filter full scan, which must recover
the complete map before the
+ // batch-mode gate so it keeps the legacy asynchronous
split-generation path.
+ PluginDrivenExternalTable table = (PluginDrivenExternalTable)
getTargetTable();
+ // The map is resolved per statement and table reference rather than
enumerated here: the MV partition
+ // compensator recorded the view it reasoned about, and its union
branch is restricted to exactly those
+ // partition names, so reading a later generation here would read
partitions that neither the MV branch
+ // nor the compensation union covers - rows silently missing from a
rewritten query.
+ Optional<TableSnapshot> tableSnapshot =
Optional.ofNullable(getQueryTableSnapshot());
+ Optional<TableScanParams> scanParams =
Optional.ofNullable(getScanParams());
+ Optional<MvccSnapshot> snapshot =
MvccUtil.getSnapshotFromContext(table, tableSnapshot, scanParams);
+ ConnectContext connectContext = ConnectContext.get();
+ StatementContext statementContext = connectContext == null ? null :
connectContext.getStatementContext();
+ Optional<Map<String, PartitionItem>> partitions = statementContext ==
null
+ ? table.getNameToPartitionItemsForScan(snapshot)
+ : statementContext.resolveScanPartitionView(table,
tableSnapshot, scanParams,
+ () -> table.getNameToPartitionItemsForScan(snapshot));
+ selectedPartitions =
materializeDeferredSelectedPartitions(selectedPartitions, partitions);
+ }
+
+ static SelectedPartitions
materializeDeferredSelectedPartitions(SelectedPartitions selectedPartitions,
+ Optional<Map<String, PartitionItem>> partitions) {
+ if (!selectedPartitions.isDeferredPartitionPruning()) {
+ return selectedPartitions;
+ }
+ // An UNAVAILABLE view (a connector entry that cannot be represented
as a Doris partition item) must not
+ // become an empty selection: keep NOT_PRUNED so the scan reads every
partition instead of none.
+ return partitions.map(items -> new SelectedPartitions(items.size(),
items, false))
+ .orElse(SelectedPartitions.NOT_PRUNED);
+ }
+
@Override
protected void convertPredicate() {
// Attempt filter pushdown via the connector SPI
if (conjuncts == null || conjuncts.isEmpty()) {
return;
}
+ // Reuse the connector filter result logical pruning already obtained
for THIS scan: the partition
+ // selection was materialized from that exact handle, so applying the
predicate again could observe a
Review Comment:
[P1] Keep predicates added after the connector filter was applied. In an
enabled async-MV union rewrite, a copied `Filter(p IN (1,2)) -> FileScan` can
gain an inner compensation `Filter(p=2)` while retaining the original
connector-filtered selection and handle H for p1+p2. The deep copier
invalidates pruning only for OlapScan, so this scan skips re-pruning; physical
translation attaches p=2 to its conjuncts. For a valid connector that returned
`remainingFilter=null` after consuming only the original predicate, reusing H
here clears every current conjunct, including p=2. The base branch then reads
p1+p2 while the MV branch reads p1, duplicating p1 rows in UNION ALL. Reuse the
result only for the predicate set it covered, or retain later conjuncts for BE
evaluation, and test this compensation case.
--
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]