This is an automated email from the ASF dual-hosted git repository.

morrySnow pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new 0d9ffb05f35 [fix](fe) Snapshot RF partition metadata during planning 
(#65919)
0d9ffb05f35 is described below

commit 0d9ffb05f35a4bfb55c79a4e343f6987715d9f0c
Author: Pxl <[email protected]>
AuthorDate: Tue Jul 28 15:37:24 2026 +0800

    [fix](fe) Snapshot RF partition metadata during planning (#65919)
    
    A query can finish planning and then wait in a Workload Group queue
    before its plan is serialized. If a selected RANGE or LIST partition is
    replaced during that wait, runtime-filter target metadata still
    references the planned partition ID, while `OlapScanNode` previously
    rebuilt runtime-filter partition boundaries from the current catalog
    during Thrift serialization. The removed partition ID was therefore
    absent from the boundaries sent to BE, which violated the runtime-filter
    partition-pruning invariant and failed the query with
    `boundary_partition_ids.contains(partition_id)`.
    
    This change snapshots runtime-filter partition boundaries during
    planning, from the same catalog view used to select partitions and build
    scan ranges. It also snapshots selected partition names so
    profile/explain generation does not look up a partition that was
    replaced after planning. No BE or Thrift protocol changes are required.
    
    The FE unit test snapshots a RANGE-partitioned scan, mutates the live
    `PartitionInfo` and partition map to simulate `REPLACE PARTITION`, and
    verifies that serialization and explain output continue to use the
    planned partition IDs and names.
    
    ### Release note
    
    Fix queued runtime-filter queries failing when selected partitions are
    replaced after planning.
---
 .../glue/translator/RuntimeFilterTranslator.java   | 25 ++++---
 .../org/apache/doris/planner/OlapScanNode.java     | 66 +++++++++++++++---
 .../translator/PhysicalPlanTranslatorTest.java     | 47 +++++++++++++
 .../org/apache/doris/planner/OlapScanNodeTest.java | 81 ++++++++++++++++++++++
 4 files changed, 198 insertions(+), 21 deletions(-)

diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/RuntimeFilterTranslator.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/RuntimeFilterTranslator.java
index 286330d79ea..6558541f0d7 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/RuntimeFilterTranslator.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/RuntimeFilterTranslator.java
@@ -35,6 +35,7 @@ import org.apache.doris.planner.CTEScanNode;
 import org.apache.doris.planner.DataStreamSink;
 import org.apache.doris.planner.DistributionMode;
 import org.apache.doris.planner.HashJoinNode;
+import org.apache.doris.planner.OlapScanNode;
 import org.apache.doris.planner.PlanNode;
 import org.apache.doris.planner.RuntimeFilter.RuntimeFilterTarget;
 import org.apache.doris.planner.ScanNode;
@@ -247,11 +248,7 @@ public class RuntimeFilterTranslator {
                     RuntimeFilterPartitionPruneClassifier.Classification 
classification =
                             RuntimeFilterPartitionPruneClassifier.classify(
                                     head.getType(), targetExpr, 
nereidsTargetExprList.get(i), scanNode);
-                    if (classification.canPrunePartitions()) {
-                        
origFilter.markTargetCanPrunePartitions(scanNode.getId());
-                    }
-                    origFilter.setTargetPartitionMonotonicity(
-                            scanNode.getId(), 
classification.getPartitionMonotonicity());
+                    setPartitionPruningMetadata(origFilter, scanNode, 
classification);
                 }
                 
origFilter.setBloomFilterSizeCalculatedByNdv(head.isBloomFilterSizeCalculatedByNdv());
                 setWaitTimeMs(origFilter, head.isNonBlocking(), isLocalTarget);
@@ -354,11 +351,7 @@ public class RuntimeFilterTranslator {
                     RuntimeFilterPartitionPruneClassifier.Classification 
classification =
                             RuntimeFilterPartitionPruneClassifier.classify(
                                     filter.getType(), targetExpr, 
filter.getTargetExpressions().get(i), scanNode);
-                    if (classification.canPrunePartitions()) {
-                        
origFilter.markTargetCanPrunePartitions(scanNode.getId());
-                    }
-                    origFilter.setTargetPartitionMonotonicity(
-                            scanNode.getId(), 
classification.getPartitionMonotonicity());
+                    setPartitionPruningMetadata(origFilter, scanNode, 
classification);
                 }
                 
origFilter.setBloomFilterSizeCalculatedByNdv(filter.isBloomFilterSizeCalculatedByNdv());
                 setWaitTimeMs(origFilter, filter.isNonBlocking(), 
isLocalTarget);
@@ -388,6 +381,18 @@ public class RuntimeFilterTranslator {
         return origFilter;
     }
 
+    private void 
setPartitionPruningMetadata(org.apache.doris.planner.RuntimeFilter 
runtimeFilter,
+            ScanNode scanNode, 
RuntimeFilterPartitionPruneClassifier.Classification classification) {
+        if (classification.canPrunePartitions()) {
+            Preconditions.checkState(scanNode instanceof OlapScanNode,
+                    "partition-pruning runtime filter target must be an 
OlapScanNode");
+            runtimeFilter.markTargetCanPrunePartitions(scanNode.getId());
+            ((OlapScanNode) 
scanNode).snapshotPartitionBoundariesForRuntimeFilter();
+        }
+        runtimeFilter.setTargetPartitionMonotonicity(
+                scanNode.getId(), classification.getPartitionMonotonicity());
+    }
+
     private void setWaitTimeMs(org.apache.doris.planner.RuntimeFilter filter,
             boolean isNonBlocking, boolean isLocalTarget) {
         if (isNonBlocking) {
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/planner/OlapScanNode.java 
b/fe/fe-core/src/main/java/org/apache/doris/planner/OlapScanNode.java
index 2535961a76f..d17c8b7dd63 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/planner/OlapScanNode.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/planner/OlapScanNode.java
@@ -182,6 +182,11 @@ public class OlapScanNode extends ScanNode {
     private long totalTabletsNum = 0;
     private long selectedIndexId = -1;
     private Collection<Long> selectedPartitionIds = Lists.newArrayList();
+    private Map<Long, String> selectedPartitionNames = Collections.emptyMap();
+    // Partition boundaries must come from the same planning snapshot as 
runtime-filter
+    // partition monotonicity. Thrift serialization can happen after query 
queueing, when
+    // the catalog may already contain replacement partition IDs.
+    private List<TPartitionBoundary> runtimeFilterPartitionBoundaries;
     private long totalBytes = 0;
     // tablet id to single replica bytes
     private Map<Long, Long> tabletBytes = Maps.newLinkedHashMap();
@@ -355,6 +360,7 @@ public class OlapScanNode extends ScanNode {
      */
     public void init() throws UserException {
         selectedPartitionNum = selectedPartitionIds.size();
+        snapshotSelectedPartitionNames();
         try {
             createScanRangeLocations();
         } catch (AnalysisException e) {
@@ -802,6 +808,7 @@ public class OlapScanNode extends ScanNode {
                         partition.getName(), "RESTORING");
             }
         }
+        snapshotSelectedPartitionNames();
         if (LOG.isDebugEnabled()) {
             LOG.debug("partition prune cost: {} ms, partitions: {}",
                     (System.currentTimeMillis() - start), 
selectedPartitionIds);
@@ -1136,9 +1143,7 @@ public class OlapScanNode extends ScanNode {
                     .append(expr.accept(ExprToSqlVisitor.INSTANCE, 
ToSqlParams.WITH_TABLE)).append("\n");
         }
 
-        String selectedPartitions = getSelectedPartitionIds().stream().sorted()
-                .map(id -> olapTable.getPartition(id).getName())
-                .collect(Collectors.joining(","));
+        String selectedPartitions = getSelectedPartitionNamesForExplain();
         output.append(prefix).append(String.format("partitions=%s/%s (%s)", 
selectedPartitionNum,
                 olapTable.getPartitions().size(), 
selectedPartitions)).append("\n");
         output.append(prefix).append(String.format("tablets=%s/%s", 
selectedSplitNum, totalTabletsNum));
@@ -1169,6 +1174,14 @@ public class OlapScanNode extends ScanNode {
         return output.toString();
     }
 
+    @VisibleForTesting
+    String getSelectedPartitionNamesForExplain() {
+        return getSelectedPartitionIds().stream().sorted()
+                .map(id -> 
Preconditions.checkNotNull(selectedPartitionNames.get(id),
+                        "missing snapshotted name for selected partition %d", 
id))
+                .collect(Collectors.joining(","));
+    }
+
     private String getExtraKeyColumnExplainName(Integer slotId) {
         SlotDescriptor extraKeySlot = desc.getSlots().stream()
                 .filter(slot -> slot.getId().asInt() == slotId)
@@ -1355,7 +1368,7 @@ public class OlapScanNode extends ScanNode {
         if (rfPruneCtx != null
                 && 
rfPruneCtx.getSessionVariable().isEnableRuntimeFilterPartitionPrune()
                 && hasRfDrivingPartitionPruning()) {
-            setPartitionBoundaries(msg.olap_scan_node);
+            setPartitionBoundariesForRuntimeFilter(msg.olap_scan_node);
         }
 
         super.toThrift(msg);
@@ -1374,15 +1387,48 @@ public class OlapScanNode extends ScanNode {
         return false;
     }
 
-    private void setPartitionBoundaries(TOlapScanNode olapScanNode) {
+    /**
+     * Snapshot partition boundaries while the query plan still owns its 
catalog snapshot.
+     * RuntimeFilterTranslator calls this when a target is classified as 
capable of
+     * partition pruning. A scan can be targeted by multiple runtime filters, 
so retain
+     * the first snapshot.
+     */
+    public void snapshotPartitionBoundariesForRuntimeFilter() {
+        if (runtimeFilterPartitionBoundaries != null) {
+            return;
+        }
+        runtimeFilterPartitionBoundaries = 
buildPartitionBoundariesForRuntimeFilter();
+    }
+
+    @VisibleForTesting
+    void snapshotSelectedPartitionNames() {
+        Map<Long, String> partitionNames = 
Maps.newHashMapWithExpectedSize(selectedPartitionIds.size());
+        for (Long partitionId : selectedPartitionIds) {
+            Partition partition = 
Preconditions.checkNotNull(olapTable.getPartition(partitionId),
+                    "missing selected partition %s during planning", 
partitionId);
+            partitionNames.put(partitionId, partition.getName());
+        }
+        selectedPartitionNames = partitionNames;
+    }
+
+    @VisibleForTesting
+    void setPartitionBoundariesForRuntimeFilter(TOlapScanNode olapScanNode) {
+        Preconditions.checkNotNull(runtimeFilterPartitionBoundaries,
+                "runtime-filter partition boundaries must be snapshotted 
during planning");
+        if (!runtimeFilterPartitionBoundaries.isEmpty()) {
+            olapScanNode.setPartitionBoundaries(new 
ArrayList<>(runtimeFilterPartitionBoundaries));
+        }
+    }
+
+    private List<TPartitionBoundary> 
buildPartitionBoundariesForRuntimeFilter() {
         PartitionInfo partitionInfo = olapTable.getPartitionInfo();
         PartitionType partType = partitionInfo.getType();
         if (partType != PartitionType.RANGE && partType != PartitionType.LIST) 
{
-            return;
+            return Collections.emptyList();
         }
         List<Column> partColumns = partitionInfo.getPartitionColumns();
         if (partColumns.isEmpty()) {
-            return;
+            return Collections.emptyList();
         }
 
         // Build partition column name → slot ID mapping
@@ -1399,7 +1445,7 @@ public class OlapScanNode extends ScanNode {
             }
         }
         if (partColToSlotId.isEmpty()) {
-            return;
+            return Collections.emptyList();
         }
 
         List<TPartitionBoundary> boundaries = new ArrayList<>();
@@ -1416,9 +1462,7 @@ public class OlapScanNode extends ScanNode {
                         partColumns, partColToSlotId);
             }
         }
-        if (!boundaries.isEmpty()) {
-            olapScanNode.setPartitionBoundaries(boundaries);
-        }
+        return boundaries;
     }
 
     private void addRangeBoundaries(List<TPartitionBoundary> boundaries, long 
partitionId,
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/glue/translator/PhysicalPlanTranslatorTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/glue/translator/PhysicalPlanTranslatorTest.java
index db88d61545c..918f9313c28 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/glue/translator/PhysicalPlanTranslatorTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/glue/translator/PhysicalPlanTranslatorTest.java
@@ -59,6 +59,7 @@ import org.apache.doris.planner.ScanContext;
 import org.apache.doris.planner.ScanNode;
 import org.apache.doris.thrift.TExplainLevel;
 import org.apache.doris.thrift.TPlanNode;
+import org.apache.doris.thrift.TRuntimeFilterType;
 import org.apache.doris.thrift.TScanRangeLocations;
 import org.apache.doris.utframe.TestWithFeService;
 
@@ -78,6 +79,7 @@ import java.util.Optional;
 import java.util.Set;
 import java.util.function.Consumer;
 import java.util.function.Supplier;
+import java.util.stream.Collectors;
 
 public class PhysicalPlanTranslatorTest extends TestWithFeService {
 
@@ -217,6 +219,51 @@ public class PhysicalPlanTranslatorTest extends 
TestWithFeService {
         
Assertions.assertFalse(nonPartitionFilteredScanNode.hasPartitionPredicate());
     }
 
+    @Test
+    public void testRfPartitionPruneSnapshotSurvivesEnablementChange() throws 
Exception {
+        int oldRuntimeFilterType = 
connectContext.getSessionVariable().getRuntimeFilterType();
+        boolean oldEnablePartitionPrune =
+                
connectContext.getSessionVariable().isEnableRuntimeFilterPartitionPrune();
+        boolean oldEnableRuntimeFilterPrune =
+                
connectContext.getSessionVariable().isEnableRuntimeFilterPrune();
+        boolean oldDisableJoinReorder = 
connectContext.getSessionVariable().isDisableJoinReorder();
+        try {
+            
connectContext.getSessionVariable().setRuntimeFilterType(TRuntimeFilterType.MIN_MAX.getValue());
+            
connectContext.getSessionVariable().setEnableRuntimeFilterPartitionPrune(false);
+            
connectContext.getSessionVariable().setEnableRuntimeFilterPrune(false);
+            connectContext.getSessionVariable().setDisableJoinReorder(true);
+
+            Planner planner = getSQLPlanner("select p.* from 
test_db.partitioned_t p "
+                    + "join [broadcast] test_db.t d on p.p1 = d.a");
+            List<OlapScanNode> scanNodes = new ArrayList<>();
+            for (PlanFragment fragment : planner.getFragments()) {
+                PlanNode root = fragment.getPlanRoot();
+                if (root != null) {
+                    root.collect(OlapScanNode.class, scanNodes);
+                }
+            }
+            OlapScanNode partitionedScan = scanNodes.stream()
+                    .filter(scan -> 
scan.getOlapTable().getName().equals("partitioned_t"))
+                    .findFirst()
+                    .orElseThrow();
+            Assertions.assertEquals(2, 
partitionedScan.getSelectedPartitionIds().size());
+
+            
connectContext.getSessionVariable().setEnableRuntimeFilterPartitionPrune(true);
+            TPlanNode thriftScanNode = 
partitionedScan.treeToThrift().getNodes().get(0);
+
+            
Assertions.assertTrue(thriftScanNode.olap_scan_node.isSetPartitionBoundaries());
+            
Assertions.assertEquals(Sets.newHashSet(partitionedScan.getSelectedPartitionIds()),
+                    
thriftScanNode.olap_scan_node.getPartitionBoundaries().stream()
+                            .map(boundary -> boundary.getPartitionId())
+                            .collect(Collectors.toSet()));
+        } finally {
+            
connectContext.getSessionVariable().setRuntimeFilterType(oldRuntimeFilterType);
+            
connectContext.getSessionVariable().setEnableRuntimeFilterPartitionPrune(oldEnablePartitionPrune);
+            
connectContext.getSessionVariable().setEnableRuntimeFilterPrune(oldEnableRuntimeFilterPrune);
+            
connectContext.getSessionVariable().setDisableJoinReorder(oldDisableJoinReorder);
+        }
+    }
+
     @Test
     public void testNereidsFileScanCarryPartitionPredicateSignal() throws 
Exception {
         PhysicalFileScan fileScan = Mockito.mock(PhysicalFileScan.class);
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/planner/OlapScanNodeTest.java 
b/fe/fe-core/src/test/java/org/apache/doris/planner/OlapScanNodeTest.java
index 136df60e3a2..89d734e7d08 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/planner/OlapScanNodeTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/planner/OlapScanNodeTest.java
@@ -21,28 +21,40 @@ import org.apache.doris.analysis.BinaryPredicate;
 import org.apache.doris.analysis.Expr;
 import org.apache.doris.analysis.InPredicate;
 import org.apache.doris.analysis.IntLiteral;
+import org.apache.doris.analysis.PartitionValue;
 import org.apache.doris.analysis.SlotDescriptor;
 import org.apache.doris.analysis.SlotId;
 import org.apache.doris.analysis.SlotRef;
 import org.apache.doris.analysis.TupleDescriptor;
 import org.apache.doris.analysis.TupleId;
 import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.OlapTable;
+import org.apache.doris.catalog.Partition;
 import org.apache.doris.catalog.PartitionKey;
 import org.apache.doris.catalog.PrimitiveType;
+import org.apache.doris.catalog.RangePartitionInfo;
+import org.apache.doris.catalog.RangePartitionItem;
 import org.apache.doris.catalog.info.TableNameInfo;
 import org.apache.doris.common.AnalysisException;
 import org.apache.doris.datasource.InternalCatalog;
+import org.apache.doris.thrift.TOlapScanNode;
+import org.apache.doris.thrift.TPartitionBoundary;
 
 import com.google.common.collect.Lists;
 import com.google.common.collect.Maps;
+import com.google.common.collect.Range;
 import org.apache.commons.collections4.map.CaseInsensitiveMap;
 import org.junit.Assert;
 import org.junit.Test;
+import org.mockito.Mockito;
 
 import java.util.ArrayList;
 import java.util.Collection;
+import java.util.Collections;
+import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.stream.Collectors;
 
 public class OlapScanNodeTest {
     // columnA in (1) hashmode=3
@@ -212,6 +224,75 @@ public class OlapScanNodeTest {
                 Lists.newArrayList(partitionSlot.getColumn()), 
tupleDescriptor, conjuncts, null));
     }
 
+    @Test
+    public void testRuntimeFilterPartitionBoundariesUsePlanningSnapshot() 
throws AnalysisException {
+        long oldTargetPartitionId = 1L;
+        long afterPartitionId = 2L;
+        long replacementPartitionId = 3L;
+        Column partitionColumn = new Column("event_date", PrimitiveType.INT);
+        RangePartitionInfo partitionInfo = new 
RangePartitionInfo(Lists.newArrayList(partitionColumn));
+        setRangePartitionItem(partitionInfo, oldTargetPartitionId, "20260721", 
"20260722");
+        setRangePartitionItem(partitionInfo, afterPartitionId, "20260722", 
"20260723");
+
+        OlapTable table = Mockito.mock(OlapTable.class);
+        Mockito.when(table.getName()).thenReturn("rfpp_queue_range_fact");
+        
Mockito.when(table.getDistributionColumnNames()).thenReturn(Collections.emptySet());
+        Mockito.when(table.getPartitionInfo()).thenReturn(partitionInfo);
+
+        Map<Long, Partition> livePartitions = new HashMap<>();
+        livePartitions.put(oldTargetPartitionId, mockPartition("p_target"));
+        livePartitions.put(afterPartitionId, mockPartition("p_after"));
+        Mockito.when(table.getPartition(Mockito.anyLong()))
+                .thenAnswer(invocation -> 
livePartitions.get(invocation.getArgument(0)));
+        Mockito.when(table.getPartitions()).thenAnswer(invocation -> 
livePartitions.values());
+
+        TupleDescriptor tupleDescriptor = new TupleDescriptor(new TupleId(1));
+        tupleDescriptor.setTable(table);
+        SlotDescriptor partitionSlot = new SlotDescriptor(new SlotId(1), 
tupleDescriptor.getId());
+        partitionSlot.setColumn(partitionColumn);
+        partitionSlot.setType(partitionColumn.getType());
+        tupleDescriptor.addSlot(partitionSlot);
+
+        OlapScanNode scanNode = new OlapScanNode(
+                new PlanNodeId(1), tupleDescriptor, "rfppScanNode", 
ScanContext.EMPTY);
+        
scanNode.setSelectedPartitionIds(Lists.newArrayList(oldTargetPartitionId, 
afterPartitionId));
+        scanNode.snapshotSelectedPartitionNames();
+        scanNode.snapshotPartitionBoundariesForRuntimeFilter();
+
+        // Simulate REPLACE PARTITION after planning but before Thrift 
serialization.
+        partitionInfo.dropPartition(oldTargetPartitionId);
+        setRangePartitionItem(partitionInfo, replacementPartitionId, 
"20260721", "20260722");
+        livePartitions.remove(oldTargetPartitionId);
+        livePartitions.put(replacementPartitionId, mockPartition("p_target"));
+
+        scanNode.snapshotPartitionBoundariesForRuntimeFilter();
+        TOlapScanNode thriftScanNode = new TOlapScanNode();
+        scanNode.setPartitionBoundariesForRuntimeFilter(thriftScanNode);
+        List<Long> serializedPartitionIds = 
thriftScanNode.getPartitionBoundaries().stream()
+                .map(TPartitionBoundary::getPartitionId)
+                .collect(Collectors.toList());
+
+        Assert.assertEquals(Lists.newArrayList(oldTargetPartitionId, 
afterPartitionId), serializedPartitionIds);
+
+        Assert.assertEquals("p_target,p_after", 
scanNode.getSelectedPartitionNamesForExplain());
+    }
+
+    private Partition mockPartition(String name) {
+        Partition partition = Mockito.mock(Partition.class);
+        Mockito.when(partition.getName()).thenReturn(name);
+        return partition;
+    }
+
+    private void setRangePartitionItem(RangePartitionInfo partitionInfo, long 
partitionId,
+            String lowerValue, String upperValue) throws AnalysisException {
+        List<Column> partitionColumns = partitionInfo.getPartitionColumns();
+        PartitionKey lower = PartitionKey.createPartitionKey(
+                Lists.newArrayList(new PartitionValue(lowerValue)), 
partitionColumns);
+        PartitionKey upper = PartitionKey.createPartitionKey(
+                Lists.newArrayList(new PartitionValue(upperValue)), 
partitionColumns);
+        partitionInfo.setItem(partitionId, false, new 
RangePartitionItem(Range.closedOpen(lower, upper)));
+    }
+
     private SlotDescriptor addSlot(TupleDescriptor tupleDescriptor, int 
slotId, String columnName) {
         SlotDescriptor slotDescriptor = new SlotDescriptor(new SlotId(slotId), 
tupleDescriptor.getId());
         slotDescriptor.setColumn(new Column(columnName, PrimitiveType.BIGINT));


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to