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

924060929 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 b7a97c26252 [fix](search) Rewrite SEARCH to slots under materialized 
CTE (#68468)
b7a97c26252 is described below

commit b7a97c262525ecc26697c040cc9ce9c3af6d0e0e
Author: 924060929 <[email protected]>
AuthorDate: Mon Sep 28 17:21:18 2026 +0800

    [fix](search) Rewrite SEARCH to slots under materialized CTE (#68468)
    
    ### What problem does this PR solve?
    
    Problem Summary:
    `RewriteSearchToSlots` is the only rule that converts the raw `Search`
    scalar function into a
    `SearchExpression` with bound slot children, which is what lets BE
    evaluate the predicate through
    the inverted index. The whole-tree instance of that rule is registered
    inside
    `notTraverseChildrenOf(ImmutableSet.of(LogicalCTEAnchor.class))`, so it
    does not descend into
    `LogicalCTEAnchor` subtrees. When a statement contains a CTE referenced
    more than
    `inline_cte_referenced_threshold` times (default 1), `CTEInline` keeps
    the anchor materialized and
    `PullUpCteAnchor` hoists all anchors above the whole plan. A SEARCH
    predicate then sits below an
    anchor, is never rewritten, reaches BE as `FunctionSearch`, and is
    rejected with
    `only inverted index queries are supported`.
    
    Minimal reproduction (table `t` has a VARCHAR column `title` with an
    `USING INVERTED` index):
    
        WITH c AS (SELECT id FROM t WHERE search('title:hello'))
        SELECT a.id FROM c a JOIN c b ON a.id = b.id;
    
    `c` is referenced twice, so it is materialized instead of inlined. After
    rewrite the producer still
    contains the raw `search('title:hello')` filter and the scan has no
    operative column, so the query
    fails on BE. The same happens when the SEARCH is written in a
    once-referenced (inlined) CTE while
    any sibling CTE is materialized.
    
    Fix:
    Also register `RewriteSearchToSlots` in
    `CTE_CHILDREN_REWRITE_JOBS_BEFORE_SUB_PATH_PUSH_DOWN`,
    before `QueryColumnCollector` and before
    `ColumnPruning`/`OperativeColumnDerive` of the
    after-push-down job list. SEARCH predicates in CTE subtrees and in a
    main body below a hoisted
    anchor are then rewritten to `SearchExpression` the same way as the
    regular whole-tree path.
    
    ### Release note
    
    Fixed: a query that combines the `SEARCH` function with a materialized
    CTE (a CTE referenced more
    than `inline_cte_referenced_threshold` times) no longer fails with
    `only inverted index queries are supported`.
---
 .../doris/nereids/jobs/executor/Rewriter.java      | 146 ++++++++++-----------
 .../rules/rewrite/RecordPlanForMvPreRewrite.java   |   8 ++
 .../rules/rewrite/SearchCteRewriteTest.java        | 128 ++++++++++++++++++
 3 files changed, 206 insertions(+), 76 deletions(-)

diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/executor/Rewriter.java 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/executor/Rewriter.java
index dbbef606244..6cf3e3e7f82 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/executor/Rewriter.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/executor/Rewriter.java
@@ -189,7 +189,6 @@ import org.apache.doris.nereids.util.MoreFieldsThread;
 
 import com.google.common.collect.ImmutableList;
 import com.google.common.collect.ImmutableSet;
-import com.google.common.collect.Lists;
 
 import java.util.List;
 
@@ -737,6 +736,10 @@ public class Rewriter extends AbstractBatchJobExecutor {
                         ),
                         bottomUp(RuleSet.PUSH_DOWN_FILTERS)
                 ),
+                // Bind SEARCH to slots. RewriteCteChildren runs this list on 
every non-anchor subtree
+                // and, when the plan has no LogicalCTEAnchor, on the whole 
plan as well. It must stay
+                // before the after-push-down 
ColumnPruning/OperativeColumnDerive.
+                bottomUp(new RewriteSearchToSlots()),
                 custom(RuleType.ELIMINATE_UNNECESSARY_PROJECT, 
EliminateUnnecessaryProject::new),
                 topic("adjust preagg status",
                         custom(RuleType.SET_PREAGG_STATUS, 
SetPreAggStatus::new)
@@ -891,82 +894,73 @@ public class Rewriter extends AbstractBatchJobExecutor {
         if (includeNormalizePlanJobs) {
             builder.addAll(NORMALIZE_PLAN_JOBS);
         }
-        builder.addAll(notTraverseChildrenOf(
-                ImmutableSet.of(LogicalCTEAnchor.class),
-                () -> {
-                    List<RewriteJob> rewriteJobs = 
Lists.newArrayListWithExpectedSize(300);
-                    rewriteJobs.addAll(jobs(
-                            topic("cte inline and pull up all cte anchor",
-                                    custom(RuleType.PULL_UP_CTE_ANCHOR, 
PullUpCteAnchor::new),
-                                    custom(RuleType.CTE_INLINE, CTEInline::new)
-                            ),
-                            topic("process limit session variables",
-                                    custom(RuleType.ADD_DEFAULT_LIMIT, 
AddDefaultLimit::new)
-                            ),
-                            topic("record query tmp plan for mv pre rewrite",
-                                    
custom(RuleType.RECORD_PLAN_FOR_MV_PRE_REWRITE, RecordPlanForMvPreRewrite::new)
-                            ),
-                            topic("rewrite cte sub-tree before sub path push 
down",
-                                    custom(RuleType.REWRITE_CTE_CHILDREN,
-                                            () -> new 
RewriteCteChildren(beforePushDownJobs, runCboRules)
-                                    )
-                            )));
-                    rewriteJobs.addAll(jobs(topic("convert outer join to anti",
-                            custom(RuleType.CONVERT_OUTER_JOIN_TO_ANTI, 
ConvertOuterJoinToAntiJoin::new))));
-                    rewriteJobs.addAll(jobs(topic("eliminate Aggregate 
according to fd items",
-                            cascadesContext -> 
cascadesContext.rewritePlanContainsTypes(LogicalAggregate.class)
-                                    || 
cascadesContext.rewritePlanContainsTypes(LogicalJoin.class)
-                                    || 
cascadesContext.rewritePlanContainsTypes(LogicalUnion.class),
-                            custom(RuleType.ELIMINATE_GROUP_BY_KEY, 
EliminateGroupByKey::new))));
-                    rewriteJobs.addAll(jobs(topic("eliminate group by key by 
uniform",
-                            custom(RuleType.ELIMINATE_GROUP_BY_KEY_BY_UNIFORM, 
EliminateGroupByKeyByUniform::new))));
-                    if (needOrExpansion) {
-                        rewriteJobs.addAll(jobs(topic("or expansion",
-                                custom(RuleType.OR_EXPANSION, () -> 
OrExpansion.INSTANCE))));
-                    }
-                    rewriteJobs.add(topic("repeat rewrite",
-                            custom(RuleType.DECOMPOSE_REPEAT, () -> 
DecomposeRepeatWithPreAggregation.INSTANCE)));
-
-                    rewriteJobs.addAll(jobs(topic("split multi distinct",
-                            custom(RuleType.DISTINCT_AGG_STRATEGY_SELECTOR,
-                                    () -> 
DistinctAggStrategySelector.INSTANCE))));
-
-                    // Rewrite search function before VariantSubPathPruning
-                    // so that ElementAt expressions from search can be 
processed
-                    rewriteJobs.addAll(jobs(
-                            bottomUp(new RewriteSearchToSlots())
-                    ));
-
-                    if (needSubPathPushDown) {
-                        rewriteJobs.addAll(jobs(
-                                topic("variant element_at push down",
-                                        
custom(RuleType.VARIANT_SUB_PATH_PRUNING, VariantSubPathPruning::new)
+        builder.addAll(jobs(
+                topic("cte inline and pull up all cte anchor",
+                        custom(RuleType.PULL_UP_CTE_ANCHOR, 
PullUpCteAnchor::new),
+                        custom(RuleType.CTE_INLINE, CTEInline::new)
+                ),
+                topic("process limit session variables",
+                        custom(RuleType.ADD_DEFAULT_LIMIT, 
AddDefaultLimit::new)
+                ),
+                topic("record query tmp plan for mv pre rewrite",
+                        custom(RuleType.RECORD_PLAN_FOR_MV_PRE_REWRITE, 
RecordPlanForMvPreRewrite::new)
+                ),
+                topic("rewrite cte sub-tree before sub path push down",
+                        custom(RuleType.REWRITE_CTE_CHILDREN,
+                                () -> new 
RewriteCteChildren(beforePushDownJobs, runCboRules)
+                        )
+                )));
+        builder.addAll(jobs(topic("convert outer join to anti",
+                custom(RuleType.CONVERT_OUTER_JOIN_TO_ANTI, 
ConvertOuterJoinToAntiJoin::new))));
+        builder.addAll(jobs(topic("eliminate Aggregate according to fd items",
+                cascadesContext -> 
cascadesContext.rewritePlanContainsTypes(LogicalAggregate.class)
+                        || 
cascadesContext.rewritePlanContainsTypes(LogicalJoin.class)
+                        || 
cascadesContext.rewritePlanContainsTypes(LogicalUnion.class),
+                custom(RuleType.ELIMINATE_GROUP_BY_KEY, 
EliminateGroupByKey::new))));
+        builder.addAll(jobs(topic("eliminate group by key by uniform",
+                custom(RuleType.ELIMINATE_GROUP_BY_KEY_BY_UNIFORM, 
EliminateGroupByKeyByUniform::new))));
+        if (needOrExpansion) {
+            builder.addAll(jobs(topic("or expansion",
+                    custom(RuleType.OR_EXPANSION, () -> 
OrExpansion.INSTANCE))));
+        }
+        builder.add(topic("repeat rewrite",
+                custom(RuleType.DECOMPOSE_REPEAT, () -> 
DecomposeRepeatWithPreAggregation.INSTANCE)));
+
+        builder.addAll(jobs(topic("split multi distinct",
+                custom(RuleType.DISTINCT_AGG_STRATEGY_SELECTOR,
+                        () -> DistinctAggStrategySelector.INSTANCE))));
+
+        // RewriteSearchToSlots is registered in 
CTE_CHILDREN_REWRITE_JOBS_BEFORE_SUB_PATH_PUSH_DOWN.
+        // The CTE children rewriter runs that list on every non-anchor 
subtree and also on the whole
+        // plan when there is no LogicalCTEAnchor, so SEARCH is bound in both 
cases from one place.
+
+        if (needSubPathPushDown) {
+            builder.addAll(jobs(
+                    topic("variant element_at push down",
+                            custom(RuleType.VARIANT_SUB_PATH_PRUNING, 
VariantSubPathPruning::new)
+                    )
+            ));
+        }
+        builder.add(
+                topic("nested column prune",
+                        custom(RuleType.NESTED_COLUMN_PRUNING, 
NestedColumnPruning::new)
+                )
+        );
+        builder.addAll(jobs(
+                        topic("rewrite cte sub-tree after sub path push down",
+                                custom(RuleType.CLEAR_CONTEXT_STATUS, 
ClearContextStatus::new),
+                                custom(RuleType.REWRITE_CTE_CHILDREN,
+                                        () -> new 
RewriteCteChildren(afterPushDownJobs, runCboRules)
                                 )
-                        ));
-                    }
-                    rewriteJobs.add(
-                            topic("nested column prune",
-                                    custom(RuleType.NESTED_COLUMN_PRUNING, 
NestedColumnPruning::new)
-                            )
-                    );
-                    rewriteJobs.addAll(jobs(
-                                    topic("rewrite cte sub-tree after sub path 
push down",
-                                            
custom(RuleType.CLEAR_CONTEXT_STATUS, ClearContextStatus::new),
-                                            
custom(RuleType.REWRITE_CTE_CHILDREN,
-                                                    () -> new 
RewriteCteChildren(afterPushDownJobs, runCboRules)
-                                            )
-                                    ),
-                                    topic("whole plan check",
-                                            custom(RuleType.ADJUST_NULLABLE, 
() -> new AdjustNullable(false))
-                                    ),
-                                    // NullableDependentExpressionRewrite need 
to be done after nullable fixed
-                                    topic("condition function", 
bottomUp(ImmutableList.of(
-                                            new 
NullableDependentExpressionRewrite())))
-                            )
-                    );
-                    return rewriteJobs;
-                }
-        ));
+                        ),
+                        topic("whole plan check",
+                                custom(RuleType.ADJUST_NULLABLE, () -> new 
AdjustNullable(false))
+                        ),
+                        // NullableDependentExpressionRewrite need to be done 
after nullable fixed
+                        topic("condition function", bottomUp(ImmutableList.of(
+                                new NullableDependentExpressionRewrite())))
+                )
+        );
         return builder.build();
     }
 
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/RecordPlanForMvPreRewrite.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/RecordPlanForMvPreRewrite.java
index ede510e1ba2..c7b38a7feae 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/RecordPlanForMvPreRewrite.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/RecordPlanForMvPreRewrite.java
@@ -19,6 +19,7 @@ package org.apache.doris.nereids.rules.rewrite;
 
 import org.apache.doris.nereids.CascadesContext;
 import org.apache.doris.nereids.StatementContext;
+import org.apache.doris.nereids.StatementContext.CteEnvironmentSnapshot;
 import org.apache.doris.nereids.jobs.JobContext;
 import org.apache.doris.nereids.jobs.executor.Rewriter;
 import org.apache.doris.nereids.rules.RuleType;
@@ -46,6 +47,11 @@ public class RecordPlanForMvPreRewrite extends 
DefaultPlanRewriter<Void> impleme
         if 
(!PreMaterializedViewRewriter.needRecordTmpPlanForRewrite(cascadesContext)) {
             return plan;
         }
+        // The temporary rewrite below shares this StatementContext, so 
snapshot the CTE environment
+        // and restore it afterwards. Otherwise the CTE producers/consumers 
rewritten here would be
+        // reused by the regular RewriteCteChildren pass, which would then 
skip its CTE children job
+        // list for those CTEs (for example RewriteSearchToSlots in the 
before-push-down list).
+        CteEnvironmentSnapshot cteEnvironmentSnapshot = 
statementContext.cacheCteEnvironment();
         // plan pre normalize
         Plan finalPlan;
         try {
@@ -63,6 +69,8 @@ public class RecordPlanForMvPreRewrite extends 
DefaultPlanRewriter<Void> impleme
         } catch (Exception e) {
             LOG.error("mv rewrite in rbo rewrite pre normalize fail, query id 
is {}",
                     cascadesContext.getConnectContext().getQueryIdentifier(), 
e);
+        } finally {
+            statementContext.restoreCteEnvironment(cteEnvironmentSnapshot);
         }
         return plan;
     }
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/rewrite/SearchCteRewriteTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/rewrite/SearchCteRewriteTest.java
new file mode 100644
index 00000000000..3ee4d861ff3
--- /dev/null
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/rewrite/SearchCteRewriteTest.java
@@ -0,0 +1,128 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.nereids.rules.rewrite;
+
+import org.apache.doris.nereids.trees.expressions.Expression;
+import org.apache.doris.nereids.trees.expressions.SearchExpression;
+import org.apache.doris.nereids.trees.expressions.functions.scalar.Search;
+import org.apache.doris.nereids.trees.plans.Plan;
+import org.apache.doris.nereids.util.PlanChecker;
+import org.apache.doris.utframe.TestWithFeService;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.ArrayList;
+import java.util.List;
+
+/**
+ * Regression test for SEARCH under a materialized CTE.
+ *
+ * <p>RewriteSearchToSlots is the only rule that turns the raw Search scalar 
function into a
+ * SearchExpression with bound slot children. The whole-tree instance of the 
rule is configured
+ * with notTraverseChildrenOf(LogicalCTEAnchor), so a SEARCH that ends up 
under a materialized CTE
+ * (or in a main body below a hoisted anchor) used to keep the raw Search 
function. It was then
+ * translated to FunctionSearch and BE rejected it with "only inverted index 
queries are supported".
+ */
+public class SearchCteRewriteTest extends TestWithFeService {
+
+    @Override
+    protected void runBeforeAll() throws Exception {
+        createDatabase("test");
+        connectContext.setDatabase("test");
+        // The tables below stay empty, so keep the OlapScan instead of 
replacing it by LogicalEmptyRelation.
+        
connectContext.getSessionVariable().setDisableNereidsRules("PRUNE_EMPTY_PARTITION");
+        createTable("CREATE TABLE t_search (\n"
+                + "  id INT,\n"
+                + "  title VARCHAR(255),\n"
+                + "  INDEX idx_title(title) USING INVERTED\n"
+                + ") ENGINE=OLAP\n"
+                + "DUPLICATE KEY(id)\n"
+                + "DISTRIBUTED BY HASH(id) BUCKETS 1\n"
+                + "PROPERTIES ('replication_num' = '1');");
+        createTable("CREATE TABLE t_dim (\n"
+                + "  id INT\n"
+                + ") ENGINE=OLAP\n"
+                + "DUPLICATE KEY(id)\n"
+                + "DISTRIBUTED BY HASH(id) BUCKETS 1\n"
+                + "PROPERTIES ('replication_num' = '1');");
+    }
+
+    private static List<Expression> allExpressions(Plan plan) {
+        List<Expression> result = new ArrayList<>(plan.getExpressions());
+        for (Plan child : plan.children()) {
+            result.addAll(allExpressions(child));
+        }
+        return result;
+    }
+
+    private static long countSearchExpression(Plan plan) {
+        return allExpressions(plan).stream()
+                .filter(expr -> expr.anyMatch(e -> e instanceof 
SearchExpression))
+                .count();
+    }
+
+    private static long countRawSearch(Plan plan) {
+        return allExpressions(plan).stream()
+                .filter(expr -> expr.anyMatch(e -> e instanceof Search))
+                .count();
+    }
+
+    private void assertSearchRewritten(Plan plan) {
+        Assertions.assertEquals(1, countSearchExpression(plan),
+                "SEARCH must be rewritten to SearchExpression, plan:\n" + 
plan.treeString());
+        Assertions.assertEquals(0, countRawSearch(plan),
+                "raw Search scalar function must not survive in the plan:\n" + 
plan.treeString());
+    }
+
+    @Test
+    public void testSearchInsideMaterializedCte() {
+        Plan plan = PlanChecker.from(connectContext)
+                .analyze("WITH c AS (SELECT id FROM t_search WHERE 
search('title:hello'))\n"
+                        + "SELECT a.id FROM c a JOIN c b ON a.id = b.id")
+                .rewrite().getPlan();
+        assertSearchRewritten(plan);
+    }
+
+    @Test
+    public void testSearchInMaterializedCteWithPreMvRecord() {
+        PlanChecker checker = PlanChecker.from(connectContext)
+                .analyze("WITH c AS (SELECT id FROM t_search WHERE 
search('title:hello'))\n"
+                        + "SELECT a.id FROM c a JOIN c b ON a.id = b.id");
+        // Force RecordPlanForMvPreRewrite, which runs a temporary 
RewriteCteChildren pass and may
+        // leave its rewritten CTE producers/consumers in the shared 
StatementContext.
+        
checker.getCascadesContext().getStatementContext().setForceRecordTmpPlan(true);
+        checker.rewrite();
+        List<Plan> recordedPlans = 
checker.getCascadesContext().getStatementContext().getTmpPlanForMvRewrite();
+        Assertions.assertEquals(1, recordedPlans.size(), "pre-MV rewrite must 
record exactly one temporary plan");
+        Assertions.assertNotNull(recordedPlans.get(0), "the recorded pre-MV 
plan must not be null");
+        assertSearchRewritten(checker.getPlan());
+    }
+
+    @Test
+    public void testSearchInInlinedCteWithMaterializedSibling() {
+        // c1 is referenced once and inlined, c3 is referenced twice and 
materialized. PullUpCteAnchor
+        // hoists c3's anchor above the whole plan, so the inlined SEARCH ends 
up below an anchor.
+        Plan plan = PlanChecker.from(connectContext)
+                .analyze("WITH c1 AS (SELECT id FROM t_search WHERE 
search('title:hello')),\n"
+                        + "     c3 AS (SELECT id FROM t_dim)\n"
+                        + "SELECT a.id FROM c1 a JOIN c3 x ON a.id = x.id JOIN 
c3 y ON a.id = y.id")
+                .rewrite().getPlan();
+        assertSearchRewritten(plan);
+    }
+}


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

Reply via email to