github-actions[bot] commented on code in PR #68499: URL: https://github.com/apache/doris/pull/68499#discussion_r4219172384
########## fe/fe-core/src/main/java/org/apache/doris/nereids/spm/builder/SPMPlan2SQLBuilder.java: ########## @@ -0,0 +1,4352 @@ +// 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.spm.builder; + +import org.apache.doris.analysis.TableScanParams; +import org.apache.doris.analysis.TableSnapshot; +import org.apache.doris.catalog.Column; +import org.apache.doris.catalog.OlapTable; +import org.apache.doris.catalog.Partition; +import org.apache.doris.common.Pair; +import org.apache.doris.nereids.parser.NereidsParser; +import org.apache.doris.nereids.properties.DistributionSpec; +import org.apache.doris.nereids.properties.DistributionSpecReplicated; +import org.apache.doris.nereids.trees.TableSample; +import org.apache.doris.nereids.trees.expressions.AggregateExpression; +import org.apache.doris.nereids.trees.expressions.Alias; +import org.apache.doris.nereids.trees.expressions.CTEId; +import org.apache.doris.nereids.trees.expressions.EqualTo; +import org.apache.doris.nereids.trees.expressions.ExprId; +import org.apache.doris.nereids.trees.expressions.Expression; +import org.apache.doris.nereids.trees.expressions.MarkJoinSlotReference; +import org.apache.doris.nereids.trees.expressions.NamedExpression; +import org.apache.doris.nereids.trees.expressions.Slot; +import org.apache.doris.nereids.trees.expressions.SlotReference; +import org.apache.doris.nereids.trees.expressions.WindowExpression; +import org.apache.doris.nereids.trees.expressions.functions.Function; +import org.apache.doris.nereids.trees.expressions.functions.agg.AggregateFunction; +import org.apache.doris.nereids.trees.expressions.functions.agg.GroupConcat; +import org.apache.doris.nereids.trees.expressions.functions.agg.MultiDistinctGroupConcat; +import org.apache.doris.nereids.trees.expressions.functions.scalar.Grouping; +import org.apache.doris.nereids.trees.expressions.functions.table.TableValuedFunction; +import org.apache.doris.nereids.trees.plans.AggPhase; +import org.apache.doris.nereids.trees.plans.JoinType; +import org.apache.doris.nereids.trees.plans.LimitPhase; +import org.apache.doris.nereids.trees.plans.Plan; +import org.apache.doris.nereids.trees.plans.algebra.SetOperation.Qualifier; +import org.apache.doris.nereids.trees.plans.physical.AbstractPhysicalJoin; +import org.apache.doris.nereids.trees.plans.physical.PhysicalAssertNumRows; +import org.apache.doris.nereids.trees.plans.physical.PhysicalCTEAnchor; +import org.apache.doris.nereids.trees.plans.physical.PhysicalCTEConsumer; +import org.apache.doris.nereids.trees.plans.physical.PhysicalCTEProducer; +import org.apache.doris.nereids.trees.plans.physical.PhysicalCatalogRelation; +import org.apache.doris.nereids.trees.plans.physical.PhysicalDistribute; +import org.apache.doris.nereids.trees.plans.physical.PhysicalEmptyRelation; +import org.apache.doris.nereids.trees.plans.physical.PhysicalExcept; +import org.apache.doris.nereids.trees.plans.physical.PhysicalFileScan; +import org.apache.doris.nereids.trees.plans.physical.PhysicalFilter; +import org.apache.doris.nereids.trees.plans.physical.PhysicalGenerate; +import org.apache.doris.nereids.trees.plans.physical.PhysicalHashAggregate; +import org.apache.doris.nereids.trees.plans.physical.PhysicalHashJoin; +import org.apache.doris.nereids.trees.plans.physical.PhysicalIntersect; +import org.apache.doris.nereids.trees.plans.physical.PhysicalLazyMaterialize; +import org.apache.doris.nereids.trees.plans.physical.PhysicalLazyMaterializeFileScan; +import org.apache.doris.nereids.trees.plans.physical.PhysicalLazyMaterializeOlapScan; +import org.apache.doris.nereids.trees.plans.physical.PhysicalLazyMaterializeTVFScan; +import org.apache.doris.nereids.trees.plans.physical.PhysicalLimit; +import org.apache.doris.nereids.trees.plans.physical.PhysicalNestedLoopJoin; +import org.apache.doris.nereids.trees.plans.physical.PhysicalOlapScan; +import org.apache.doris.nereids.trees.plans.physical.PhysicalOneRowRelation; +import org.apache.doris.nereids.trees.plans.physical.PhysicalPartitionTopN; +import org.apache.doris.nereids.trees.plans.physical.PhysicalProject; +import org.apache.doris.nereids.trees.plans.physical.PhysicalQuickSort; +import org.apache.doris.nereids.trees.plans.physical.PhysicalRecursiveUnion; +import org.apache.doris.nereids.trees.plans.physical.PhysicalRecursiveUnionAnchor; +import org.apache.doris.nereids.trees.plans.physical.PhysicalRecursiveUnionProducer; +import org.apache.doris.nereids.trees.plans.physical.PhysicalRelation; +import org.apache.doris.nereids.trees.plans.physical.PhysicalRepeat; +import org.apache.doris.nereids.trees.plans.physical.PhysicalResultSink; +import org.apache.doris.nereids.trees.plans.physical.PhysicalSetOperation; +import org.apache.doris.nereids.trees.plans.physical.PhysicalSink; +import org.apache.doris.nereids.trees.plans.physical.PhysicalStorageLayerAggregate; +import org.apache.doris.nereids.trees.plans.physical.PhysicalTVFRelation; +import org.apache.doris.nereids.trees.plans.physical.PhysicalTopN; +import org.apache.doris.nereids.trees.plans.physical.PhysicalUnion; +import org.apache.doris.nereids.trees.plans.physical.PhysicalWindow; +import org.apache.doris.nereids.trees.plans.physical.PhysicalWorkTableReference; +import org.apache.doris.nereids.trees.plans.visitor.PlanVisitor; + +import com.google.common.annotations.VisibleForTesting; +import com.google.common.collect.Lists; +import org.apache.commons.lang3.StringUtils; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.LinkedHashSet; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.Optional; +import java.util.Set; +import java.util.TreeMap; +import java.util.regex.Matcher; +import java.util.regex.Pattern; +import java.util.stream.Collectors; + +/** + * Physical-plan decompiler (M1). + * + * Decompiles the optimal physical plan tree output by the optimizer into one + * semantically equivalent standard SQL (planSql). planSql freezes the optimizer's key + * decisions (JOIN order, aggregate structure, sort semantics) through SQL structure and + * HINTs, so the plan can be replayed during query rewrite. This is the core engine of + * the SPM "plan-to-SQL" approach. + * + * Design (design doc 6.2): + * + * - Pass-through: physical infrastructure nodes such as PhysicalDistribute and + * PhysicalResultSink have no SQL equivalent and return the child result directly. + * - Merge: a local aggregate (LOCAL) registers the aggregate function names onto the + * child relation without adding nesting. + * - Wrap: Filter / Join / global aggregate / TopN / Project / Window create a new + * SQLRelation and set an alias; parent operators reference it as a subquery through + * toRelationSQL(). + * + * Known M1 simplifications: + * + * - Scan uses the real column names (no c_N normalization; column-name collision + * scenarios are handled in a later milestone). + */ +public class SPMPlan2SQLBuilder extends PlanVisitor<SQLRelation, Void> { + + private static final Logger LOG = LogManager.getLogger(SPMPlan2SQLBuilder.class); + + /** Generated output names ("c_" + ExprId) produced by re-export labels and dedupe renames. */ + private static final Pattern GENERATED_NAME_PATTERN = Pattern.compile("\\bc_\\d+\\b"); + + /** JOIN distribution HINT prefix constants. */ + private static final String HINT_JOIN_BROADCAST = "BROADCAST"; + private static final String HINT_JOIN_SHUFFLE = "SHUFFLE"; + + /** + * Re-export labels a hoist step created one level below ("X AS c_N" items): maps the + * generated label back to the ExprId of the item it re-exports, so a later hoist step + * can re-export the SAME label through one more projection level (tpcds q78: the + * ORDER BY hoisted out of the 16-item projection references c_14/c_15/c_16, and the + * 11-item projection above it must therefore re-export them). + */ + private final Map<String, ExprId> reExportedLabels = new HashMap<>(); + + /** Expression printer (carries the columnNames mapping of SQLRelation). */ + private final SPMExprSqlBuilder exprSqlBuilder = new SPMExprSqlBuilder(); + + /** Per-decompile alias sequence for LATERAL VIEW clauses whose output slot has no + * user qualifier (reset in toSQL). */ + private int lateralViewSeq = 0; + + /** CTE body relations keyed by CTEId (registered by visitPhysicalCTEProducer). The + * consumer references the CTE by its alias instead of inlining the body, so the + * decompiled planSql keeps the WITH structure - one definition shared by every + * consumer - like StarRocks does. */ + private final Map<CTEId, SQLRelation> cteBodies = new HashMap<>(); + + /** CTEId -> the generated alias (t_N) that references the WITH definition. */ + private final Map<CTEId, String> cteAliases = new HashMap<>(); + + /** + * The WITH entries (alias AS (body)) in producer visit order - the definition + * order of the emitted WITH clause. They are attached to the ROOT relation in + * toSQL(): every consumer lies inside the statement, so the definitions are + * visible everywhere regardless of how deeply the CTE anchors are nested in the + * physical tree (attaching at each anchor instead would place a definition inside + * one join branch while a consumer sits in another branch). Producers are visited + * in dependency order - a consumer requires its producer to be registered already - + * so the list is a valid SQL definition order (a CTE used by another CTE is + * defined before it). + */ + private final List<String> cteDefinitions = new ArrayList<>(); + + /** Recursive CTE output column names keyed by CTE name. Registered from the anchor + * branch of a PhysicalRecursiveUnion; used to name the recursive work-table self + * reference and the outer reference so the decompiled WITH RECURSIVE stays + * re-parseable. */ + private final Map<String, List<String>> recursiveCteColumns = new HashMap<>(); + + /** Local (partial) aggregate output column -> the aggregate input expression it + * aggregates, e.g. local partial_sum(x)#N registers N -> x. A global + * aggregate references such a column as sum(partial_sum(x)#N); the global + * decompile rewrites it back to sum(x) so the decompiled SQL stays a single + * logical aggregate (no partial_ intermediate functions). */ + private final Map<ExprId, Expression> localAggParams = new HashMap<>(); + + /** + * Buffer columns produced by the distinct-dedup branch: their defining partial + * aggregate consumed a DATA column (e.g. DISTINCT_LOCAL's partial_count(key)), + * not another partial buffer. In a DISTINCT_GLOBAL merge stage a count(...) over + * such a buffer is the user's count(DISTINCT key). + * + * A plan may mix plain aggregates with a distinct one (TPCDS q28: + * avg(x), count(x), count(DISTINCT x)): the plain count rides along the same + * DISTINCT_GLOBAL stage, but its buffer is a merge chain (its defining partial + * aggregated another partial buffer, e.g. partial_count(partial_count(x))), so it + * must NOT be rendered with DISTINCT. Both chains otherwise resolve to the same + * data column, which is why the provenance has to be recorded while the + * intermediate stages are eliminated. + */ + private final Set<ExprId> distinctMergeBuffers = new HashSet<>(); + + /** The OUTERMOST projection of the decompiled tree (by identity). It alone prunes its + * output to the final user-visible columns; every intermediate projection outputs the + * full child column set plus its own expressions, so an upper layer (filter / join / + * aggregate / projection) can always resolve the columns it references. */ + private final Set<Plan> outputProjects = + Collections.newSetFromMap(new java.util.IdentityHashMap<>()); + + /** + * Identity map: physical plan node -> the ExprIds of its output columns that upper + * layers actually consume (the final result columns plus every column referenced by + * a decompiled expression above). An entry ABSENT, or mapped to null, means "no + * pruning for this node and everything below it" (the pre-pruning behaviour). + * + * The decompiled planSql inflates when every intermediate projection re-emits the + * whole child column set (a deep TPCDS join stack re-lists hundreds of c_N columns + * per level). This map drives live-column pruning: each SELECT list is filtered to + * the live columns only, so a column that is never referenced above and is not part + * of the final result disappears from every intermediate projection. Equivalence is + * preserved because only dead columns are dropped. + */ + private final Map<Plan, Set<ExprId>> neededOutputs = + new java.util.IdentityHashMap<>(); + + /** + * CTEId -> the CTE body output columns referenced by every inlined consumer (the + * consumer output slots mapped back to their producer slots). Filled while the main + * query is propagated top-down; the CTE producer subtree is propagated afterwards + * with exactly this column set, so a big CTE body (a deep TPCDS join stack) only + * keeps the columns its consumers actually read. + */ + private final Map<CTEId, Set<ExprId>> cteConsumerNeeds = new HashMap<>(); + + /** + * Per-decompile generated-column-name state: when a decompiled output column has no + * clean user alias it is exported under a generated c_N name (c_1, c_2, ...), + * assigned in decompile order and memoized by the output ExprId so every reference + * to the same column prints the same alias. A fresh per-decompile sequence (instead + * of embedding the analyzer's ExprId) keeps the generated names small, readable and + * independent of how many internal ExprIds the optimizer allocated. + */ + private final Map<ExprId, String> generatedColumnNames = new HashMap<>(); + private int generatedColumnSeq = 0; + + /** + * Normalized names already visible in the CURRENT decompile: user aliases and the + * preserved source columns of the relations a generated name gets exported next to. + * A generated c_N must avoid them - the counter alone only keeps the + * GENERATED names apart, so a source column literally named c_1 (or a group-by + * item / user alias of that name) could end up next to sum(v) AS c_1 or + * MARK_SLOT c_1: two ExprIds registered under one visible name make the + * enclosing projection / result sink read an AMBIGUOUS column from the derived + * relation after reload. Project / window outputs repair duplicates afterwards + * (dedupeSelectOutputNames); the join (MARK_SLOT / explicit projection) and aggregate + * exports do not, so they reserve here instead. + */ + private final Set<String> reservedOutputNames = new HashSet<>(); + + /** + * Rejects freezing when any expression of the plan carries a + * SessionVarGuardExpr (see containsSessionVarGuard): the guard holds the + * alias-UDF DEFINITION's saved session variables and has no SQL rendering. + * + * @param plan the physical plan about to be decompiled + */ + public static void rejectSessionVarGuardedExpressions(Plan plan) { + org.apache.doris.nereids.spm.SPMPlanTreeSupport.<RuntimeException>walkPlans(plan, node -> { + for (org.apache.doris.nereids.trees.expressions.Expression expr + : node.getExpressions()) { + if (containsSessionVarGuard(expr)) { + throw new UnsupportedOperationException("SPM cannot freeze an expression" + + " carrying a session-variable guard: " + expr); + } + } + }); + } + + private static boolean containsSessionVarGuard( + org.apache.doris.nereids.trees.expressions.Expression expr) { + if (expr instanceof org.apache.doris.nereids.trees.expressions.SessionVarGuardExpr) { + return true; + } + for (org.apache.doris.nereids.trees.expressions.Expression child : expr.children()) { + if (containsSessionVarGuard(child)) { + return true; + } + } + return false; + } + + /** + * Marks the OUTERMOST projection of the decompiled tree: the first PhysicalProject + * reached from the root along single-child pass-through nodes (ResultSink / Sort / + * Distribute / ...). Walking stops at an aggregate - an aggregate that feeds the + * result carries the final SELECT list itself, so there is no outermost projection + * and every projection below it is an intermediate one (full child output). + */ + private void markOutputProject(Plan plan) { + Plan current = plan; + while (current != null) { + if (current instanceof PhysicalProject) { + outputProjects.add(current); + return; + } + if (current instanceof PhysicalHashAggregate || current.arity() != 1) { + return; + } + current = current.child(0); + } + } + + // ==================== live-column analysis (dead-column pruning) ==================== + + /** Whether this subtree may be pruned by the live-column analysis. Non-linear + * decompile shapes (CTE bodies, set operations, recursive unions, grouping sets) + * are left untouched: their select lists are kept whole. */ + private static boolean isPrunableNode(Plan node) { + return node instanceof PhysicalProject + || node instanceof PhysicalWindow + || node instanceof PhysicalHashJoin + || node instanceof PhysicalNestedLoopJoin + || node instanceof PhysicalHashAggregate + || node instanceof PhysicalFilter + || node instanceof PhysicalTopN + || node instanceof PhysicalQuickSort + || node instanceof PhysicalLimit + || node instanceof PhysicalDistribute + || node instanceof PhysicalLazyMaterialize + || node instanceof PhysicalLazyMaterializeOlapScan + || node instanceof PhysicalPartitionTopN + || node instanceof PhysicalAssertNumRows + || node instanceof PhysicalResultSink + || node instanceof PhysicalCTEAnchor; + } + + /** Pre-analysis entry: fills neededOutputs top-down from the root. */ + private void computeNeeded(Plan root) { + neededOutputs.clear(); + cteConsumerNeeds.clear(); + Set<ExprId> rootNeed = new HashSet<>(); + for (Slot slot : root.getOutput()) { + rootNeed.add(slot.getExprId()); + } + propagateNeed(root, rootNeed); + } + + /** + * Top-down live-column propagation. need is the set of output ExprIds of + * node that upper layers consume; null means "keep everything" (never prune + * below). Every parent hands each child the child's output columns that must stay + * alive: the columns the parent passes through and the columns the parent's own + * decompiled expressions reference. + */ + private void propagateNeed(Plan node, Set<ExprId> need) { + if (node == null) { + return; + } + if (node instanceof PhysicalCTEConsumer) { + // a CTE consumer has no children; record which body columns its live output + // columns map back to, so the producer subtree can be pruned to exactly them + if (need != null) { + PhysicalCTEConsumer consumer = (PhysicalCTEConsumer) node; + Set<ExprId> bodyNeeds = cteConsumerNeeds.computeIfAbsent( + consumer.getCteId(), k -> new HashSet<>()); + for (Slot slot : consumer.getOutput()) { + if (need.contains(slot.getExprId())) { + bodyNeeds.add(consumer.getProducerSlot(slot).getExprId()); + } + } + } + return; + } + if (need == null || !isPrunableNode(node)) { + // keep everything on this node and below (no pruning boundary) + for (Plan child : node.children()) { + propagateNeed(child, null); + } + return; + } + Set<ExprId> nodeNeed = new HashSet<>(need); + neededOutputs.put(node, nodeNeed); + + if (node instanceof PhysicalCTEAnchor) { + // child(0) is the CTE producer whose body every consumer inlines. The main + // query (child(1)) is propagated FIRST so every consumer records which body + // columns it reads; the body subtree is then pruned to exactly those columns. + // A body with no consumer (or one whose consumers never resolve) keeps every + // column. + PhysicalCTEAnchor<? extends Plan, ? extends Plan> anchor = + (PhysicalCTEAnchor<? extends Plan, ? extends Plan>) node; + if (anchor.child(1) != null) { + propagateNeed(anchor.child(1), nodeNeed); + } + Plan producer = anchor.child(0); + if (producer == null) { + return; + } + if (producer instanceof PhysicalCTEProducer) { + PhysicalCTEProducer<? extends Plan> cteProducer = + (PhysicalCTEProducer<? extends Plan>) producer; + Set<ExprId> bodyNeed = cteConsumerNeeds.get(cteProducer.getCteId()); + if (bodyNeed != null && !bodyNeed.isEmpty()) { + propagateNeed(cteProducer.child(0), bodyNeed); + } else { + propagateNeed(cteProducer.child(0), null); + } + } else { + propagateNeed(producer, null); + } + return; + } + if (node instanceof PhysicalProject) { + propagateProjectNeed((PhysicalProject<? extends Plan>) node, nodeNeed); + } else if (node instanceof PhysicalHashJoin || node instanceof PhysicalNestedLoopJoin) { + propagateJoinNeed((AbstractPhysicalJoin<? extends Plan, ? extends Plan>) node, nodeNeed); + } else if (node instanceof PhysicalHashAggregate) { + propagateAggregateNeed((PhysicalHashAggregate<? extends Plan>) node, nodeNeed); + } else if (node instanceof PhysicalWindow) { + propagateWindowNeed((PhysicalWindow<? extends Plan>) node, nodeNeed); + } else if (node instanceof PhysicalFilter) { + PhysicalFilter<? extends Plan> filter = (PhysicalFilter<? extends Plan>) node; + Set<ExprId> childNeed = new HashSet<>(nodeNeed); + addExprSlots(filter.getPredicate(), childNeed); + propagateNeed(filter.child(0), childNeed); + } else if (node instanceof PhysicalTopN) { + PhysicalTopN<? extends Plan> topN = (PhysicalTopN<? extends Plan>) node; + Set<ExprId> childNeed = new HashSet<>(nodeNeed); + for (org.apache.doris.nereids.properties.OrderKey key : topN.getOrderKeys()) { + addExprSlots(key.getExpr(), childNeed); + } + propagateNeed(topN.child(0), childNeed); + } else if (node instanceof PhysicalQuickSort) { + PhysicalQuickSort<? extends Plan> sort = (PhysicalQuickSort<? extends Plan>) node; + Set<ExprId> childNeed = new HashSet<>(nodeNeed); + for (org.apache.doris.nereids.properties.OrderKey key : sort.getOrderKeys()) { + addExprSlots(key.getExpr(), childNeed); + } + propagateNeed(sort.child(0), childNeed); + } else { + // pure pass-through nodes (limit / distribute / lazy materialize / ...): + // their output slots are the child slots, so the same need flows down + for (Plan child : node.children()) { + propagateNeed(child, nodeNeed); + } + } + } + + /** PhysicalProject: live pass-through columns plus the columns referenced by the + * project's own expressions whose output is itself live. */ + private void propagateProjectNeed(PhysicalProject<? extends Plan> project, Set<ExprId> need) { + Plan child = project.child(0); + if (child == null) { + return; + } + Set<ExprId> childNeed = new HashSet<>(); + boolean finalProject = outputProjects.contains(project); + if (finalProject) { + // the outermost projection emits its full projection list: every referenced + // child column must stay alive + for (NamedExpression projectExpr : project.getProjects()) { + addExprSlots(projectExpr, childNeed); + } + } else { + // pass-through: the child columns that are still live above this projection + Set<ExprId> childOut = outputIdSet(child); + for (ExprId id : need) { + if (childOut.contains(id)) { + childNeed.add(id); + } + } + // plus the columns referenced by this projection's own live expressions + for (NamedExpression projectExpr : project.getProjects()) { + if (need.contains(projectExpr.getExprId())) { + addExprSlots(projectExpr, childNeed); + } + } + } + propagateNeed(child, childNeed); + } + + /** Join: the preserved output columns flow to the side that produces them; the + * hash / other / mark conjuncts are always decompiled, so every column they + * reference stays alive on its own side. */ + private void propagateJoinNeed(AbstractPhysicalJoin<? extends Plan, ? extends Plan> join, + Set<ExprId> need) { + Plan left = join.left(); + Plan right = join.right(); + if (left == null || right == null) { + return; + } + Set<ExprId> leftOut = outputIdSet(left); + Set<ExprId> rightOut = outputIdSet(right); + Set<ExprId> conjRefs = new HashSet<>(); + for (Expression conjunct : join.getHashJoinConjuncts()) { + addExprSlots(conjunct, conjRefs); + } + for (Expression conjunct : join.getOtherJoinConjuncts()) { + addExprSlots(conjunct, conjRefs); + } + for (Expression conjunct : join.getMarkJoinConjuncts()) { + addExprSlots(conjunct, conjRefs); + } + Set<ExprId> leftNeed = new HashSet<>(); + Set<ExprId> rightNeed = new HashSet<>(); + for (ExprId id : need) { + if (leftOut.contains(id)) { + leftNeed.add(id); + } + if (rightOut.contains(id)) { + rightNeed.add(id); + } + } + for (ExprId id : conjRefs) { + if (leftOut.contains(id)) { + leftNeed.add(id); + } + if (rightOut.contains(id)) { + rightNeed.add(id); + } + } + propagateNeed(left, leftNeed); + propagateNeed(right, rightNeed); + } + + /** Aggregate: the global stage keeps its whole output list (GROUP BY keys and + * aggregate functions), but the input columns its expressions reference must stay + * alive below. Local / intermediate execution stages are pass-through nodes. */ + private void propagateAggregateNeed(PhysicalHashAggregate<? extends Plan> agg, Set<ExprId> need) { + AggPhase phase = agg.getAggPhase(); + if (phase.isLocal() || isIntermediateAggStage(agg)) { + for (Plan child : agg.children()) { + propagateNeed(child, need); + } + return; + } + Plan child = agg.child(0); + if (child == null) { + return; + } + Set<ExprId> childNeed = new HashSet<>(); + for (Expression groupBy : agg.getGroupByExpressions()) { + addExprSlots(groupBy, childNeed); + } + for (NamedExpression output : agg.getOutputExpressions()) { + addAggregateOutputRefs(output, childNeed); + } + propagateNeed(child, childNeed); + } + + /** Window: live pass-through columns plus the input columns of the live window + * expressions. */ + private void propagateWindowNeed(PhysicalWindow<? extends Plan> window, Set<ExprId> need) { + Plan child = window.child(0); + if (child == null) { + return; + } + Set<ExprId> childNeed = new HashSet<>(); + Set<ExprId> childOut = outputIdSet(child); + for (ExprId id : need) { + if (childOut.contains(id)) { + childNeed.add(id); + } + } + for (NamedExpression windowExpr : window.getWindowExpressions()) { + if (need.contains(windowExpr.getExprId())) { + addExprSlots(windowExpr, childNeed); + } + } + propagateNeed(child, childNeed); + } + + /** The columns referenced by one global-aggregate output expression, with every + * local partial-buffer reference resolved down to its data columns (mirrors + * appendAggSelect / resolveBufferSlots so the pruned SELECT lists keep exactly the + * columns the decompiled aggregate prints). */ + private void addAggregateOutputRefs(NamedExpression output, Set<ExprId> out) { + Expression inner = output instanceof Alias ? ((Alias) output).child() : output; + if (inner instanceof SlotReference) { + // group-by key pass-through + out.add(((SlotReference) inner).getExprId()); + return; + } + if (inner instanceof AggregateExpression) { + AggregateExpression aggExpr = (AggregateExpression) inner; + List<Expression> args = aggExpr.getFunction().children().isEmpty() + ? new ArrayList<>(aggExpr.children()) : aggExpr.getFunction().children(); + for (Expression arg : args) { + addResolvedExprSlots(arg, out); + } + return; + } + addExprSlots(inner, out); + } + + /** Collects every slot of expr, resolving partial-buffer slots through + * localAggParams down to their data columns (mirrors resolveBufferSlots). */ + private void addResolvedExprSlots(Expression expr, Set<ExprId> out) { + if (expr instanceof SlotReference) { + Expression param = localAggParams.get(((SlotReference) expr).getExprId()); + if (param != null) { + addResolvedExprSlots(param, out); + return; + } + out.add(((SlotReference) expr).getExprId()); + return; + } + for (Expression childExpr : expr.children()) { + addResolvedExprSlots(childExpr, out); + } + } + + /** Output ExprId set of a plan node. */ + private static Set<ExprId> outputIdSet(Plan node) { + Set<ExprId> ids = new HashSet<>(); + for (Slot slot : node.getOutput()) { + ids.add(slot.getExprId()); + } + return ids; + } + + /** Adds every slot ExprId used by an expression. */ + private static void addExprSlots(Expression expr, Set<ExprId> out) { + if (expr == null) { + return; + } + collectAllSlotIds(expr, out); + } + + /** Pre-walk that fills localAggParams bottom-up so the live-column analysis + * (which runs before the decompile walk) can resolve partial-buffer references. The + * decompile walk re-fills the same map while it descends (harmless duplicate). */ + private void collectLocalAggParams(Plan node) { + for (Plan child : node.children()) { + collectLocalAggParams(child); + } + if (node instanceof PhysicalHashAggregate) { + PhysicalHashAggregate<? extends Plan> agg = (PhysicalHashAggregate<? extends Plan>) node; + if (agg.getAggPhase().isLocal() || isIntermediateAggStage(agg)) { + recordLocalAggStage(agg); + } + } + } + + /** Registers one local (partial) aggregate stage's buffer outputs (see the field + * comment of localAggParams). */ + private void recordLocalAggStage(PhysicalHashAggregate<? extends Plan> agg) { + for (NamedExpression output : agg.getOutputExpressions()) { + Expression inner = output instanceof Alias ? ((Alias) output).child() : output; + if (inner instanceof AggregateExpression && isPartialAggregate((AggregateExpression) inner)) { + Expression rawParam = extractPartialParam((AggregateExpression) inner); + if (agg.getGroupByExpressions().isEmpty() + && isDistinctMergeContribution(((AggregateExpression) inner).children())) { + distinctMergeBuffers.add(output.getExprId()); + } + Expression param = rawParam; + if (param == null) { + // count(*) has no data argument (extractPartialParam -> null). Map the + // buffer to the partial expression itself so the enclosing + // merge-finalize count(*) resolves its argument back to the star, which + // appendAggSelect collapses into count(*) (isNestedNoArgCount) - without + // this the raw buffer slot (named e.g. "partial_count(*)") would leak + // into the decompiled SQL as the invalid count(partial_count(*)). + param = inner; + } + while (param instanceof SlotReference) { + Expression resolved = localAggParams.get(((SlotReference) param).getExprId()); + if (resolved == null) { + break; + } + param = resolved; + } + localAggParams.put(output.getExprId(), param); + } + } + } + + /** + * Whether an output of an eliminated DISTINCT_LOCAL stage is the distinct-dedup + * contribution (see distinctMergeBuffers). + * + * The distinction has to be read from the aggregate expression's own children - + * the physical input of the buffer - because the merge function's argument is + * normalized to the data column for EVERY buffer: + * + * - merge-chain buffer (plain aggregate riding along): its input is another + * partial buffer, already recorded in localAggParams, or a count-star function + * with no argument at all; + * - distinct-dedup buffer: its input is the dedup KEY data column, either as the + * bare slot (partial_count(key)) or wrapped as the partial function's argument + * (count(key)). + */ + private boolean isDistinctMergeContribution(List<Expression> exprChildren) { + if (exprChildren.isEmpty()) { + return false; + } + Expression bufferArg = exprChildren.get(0); + if (bufferArg instanceof SlotReference) { + return !localAggParams.containsKey(((SlotReference) bufferArg).getExprId()); + } + List<Expression> argChildren = bufferArg.children(); + if (argChildren.isEmpty()) { + // e.g. the count(*) star buffer: raw-row counting, never the distinct merge + return false; + } + for (Expression child : argChildren) { + if (!(child instanceof SlotReference) + || localAggParams.containsKey(((SlotReference) child).getExprId())) { + return false; + } + } + return true; + } + + /** The live column filter for one explicit SELECT list: the list a node emits is + * pruned to the node's needed output columns (no entry in the map, or a null need, + * keeps the whole list - e.g. mock plan trees and unpruned subtrees). */ + private void filterLiveSelects(Plan node, List<Pair<ExprId, String>> selects) { + Set<ExprId> need = neededOutputs.get(node); + if (need == null) { + return; + } + selects.removeIf(p -> !need.contains(p.key())); + } + + /** + * Decompile entry: physical plan -> planSql. + * + * @param plan the optimal physical plan + * @return planSql (standard SQL text) + */ + public String toSQL(Plan plan) { + // An alias-UDF expansion computed under the DEFINITION's stored session + // variables (decimalOverflowScale / enable_decimal256 / ...) carries a + // SessionVarGuardExpr; the SQL text has no way to express that guard, so + // printing only its child would replan the arithmetic under the LATER caller's + // variables and could change its type, scale or value while the same bind SQL + // still matches. Reject freezing instead (CREATE keeps the user planSql / + // falls back to the parameterized-plan-tree path). + rejectSessionVarGuardedExpressions(plan); Review Comment: [P1] Keep alias-UDF definition settings when freezing even if the creator currently matches them. AliasUdfBuilder adds SessionVarGuardExpr only when the creator's result-affecting variables differ from the UDF's saved variables. In the equal case this rejection sees no guard and stores the expanded decimal arithmetic as ordinary SQL. A later caller with different enable_decimal256 or decimal_overflow_scale still matches the bind SQL, but frozen replay analyzes that arithmetic under the caller's settings instead of the UDF's saved settings, changing its type, scale, or value. Retain a UDF dependency marker and decline freezing this expansion regardless of creator equality. ########## fe/fe-core/src/main/java/org/apache/doris/nereids/spm/builder/SPMPlan2SQLBuilder.java: ########## @@ -0,0 +1,4352 @@ +// 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.spm.builder; + +import org.apache.doris.analysis.TableScanParams; +import org.apache.doris.analysis.TableSnapshot; +import org.apache.doris.catalog.Column; +import org.apache.doris.catalog.OlapTable; +import org.apache.doris.catalog.Partition; +import org.apache.doris.common.Pair; +import org.apache.doris.nereids.parser.NereidsParser; +import org.apache.doris.nereids.properties.DistributionSpec; +import org.apache.doris.nereids.properties.DistributionSpecReplicated; +import org.apache.doris.nereids.trees.TableSample; +import org.apache.doris.nereids.trees.expressions.AggregateExpression; +import org.apache.doris.nereids.trees.expressions.Alias; +import org.apache.doris.nereids.trees.expressions.CTEId; +import org.apache.doris.nereids.trees.expressions.EqualTo; +import org.apache.doris.nereids.trees.expressions.ExprId; +import org.apache.doris.nereids.trees.expressions.Expression; +import org.apache.doris.nereids.trees.expressions.MarkJoinSlotReference; +import org.apache.doris.nereids.trees.expressions.NamedExpression; +import org.apache.doris.nereids.trees.expressions.Slot; +import org.apache.doris.nereids.trees.expressions.SlotReference; +import org.apache.doris.nereids.trees.expressions.WindowExpression; +import org.apache.doris.nereids.trees.expressions.functions.Function; +import org.apache.doris.nereids.trees.expressions.functions.agg.AggregateFunction; +import org.apache.doris.nereids.trees.expressions.functions.agg.GroupConcat; +import org.apache.doris.nereids.trees.expressions.functions.agg.MultiDistinctGroupConcat; +import org.apache.doris.nereids.trees.expressions.functions.scalar.Grouping; +import org.apache.doris.nereids.trees.expressions.functions.table.TableValuedFunction; +import org.apache.doris.nereids.trees.plans.AggPhase; +import org.apache.doris.nereids.trees.plans.JoinType; +import org.apache.doris.nereids.trees.plans.LimitPhase; +import org.apache.doris.nereids.trees.plans.Plan; +import org.apache.doris.nereids.trees.plans.algebra.SetOperation.Qualifier; +import org.apache.doris.nereids.trees.plans.physical.AbstractPhysicalJoin; +import org.apache.doris.nereids.trees.plans.physical.PhysicalAssertNumRows; +import org.apache.doris.nereids.trees.plans.physical.PhysicalCTEAnchor; +import org.apache.doris.nereids.trees.plans.physical.PhysicalCTEConsumer; +import org.apache.doris.nereids.trees.plans.physical.PhysicalCTEProducer; +import org.apache.doris.nereids.trees.plans.physical.PhysicalCatalogRelation; +import org.apache.doris.nereids.trees.plans.physical.PhysicalDistribute; +import org.apache.doris.nereids.trees.plans.physical.PhysicalEmptyRelation; +import org.apache.doris.nereids.trees.plans.physical.PhysicalExcept; +import org.apache.doris.nereids.trees.plans.physical.PhysicalFileScan; +import org.apache.doris.nereids.trees.plans.physical.PhysicalFilter; +import org.apache.doris.nereids.trees.plans.physical.PhysicalGenerate; +import org.apache.doris.nereids.trees.plans.physical.PhysicalHashAggregate; +import org.apache.doris.nereids.trees.plans.physical.PhysicalHashJoin; +import org.apache.doris.nereids.trees.plans.physical.PhysicalIntersect; +import org.apache.doris.nereids.trees.plans.physical.PhysicalLazyMaterialize; +import org.apache.doris.nereids.trees.plans.physical.PhysicalLazyMaterializeFileScan; +import org.apache.doris.nereids.trees.plans.physical.PhysicalLazyMaterializeOlapScan; +import org.apache.doris.nereids.trees.plans.physical.PhysicalLazyMaterializeTVFScan; +import org.apache.doris.nereids.trees.plans.physical.PhysicalLimit; +import org.apache.doris.nereids.trees.plans.physical.PhysicalNestedLoopJoin; +import org.apache.doris.nereids.trees.plans.physical.PhysicalOlapScan; +import org.apache.doris.nereids.trees.plans.physical.PhysicalOneRowRelation; +import org.apache.doris.nereids.trees.plans.physical.PhysicalPartitionTopN; +import org.apache.doris.nereids.trees.plans.physical.PhysicalProject; +import org.apache.doris.nereids.trees.plans.physical.PhysicalQuickSort; +import org.apache.doris.nereids.trees.plans.physical.PhysicalRecursiveUnion; +import org.apache.doris.nereids.trees.plans.physical.PhysicalRecursiveUnionAnchor; +import org.apache.doris.nereids.trees.plans.physical.PhysicalRecursiveUnionProducer; +import org.apache.doris.nereids.trees.plans.physical.PhysicalRelation; +import org.apache.doris.nereids.trees.plans.physical.PhysicalRepeat; +import org.apache.doris.nereids.trees.plans.physical.PhysicalResultSink; +import org.apache.doris.nereids.trees.plans.physical.PhysicalSetOperation; +import org.apache.doris.nereids.trees.plans.physical.PhysicalSink; +import org.apache.doris.nereids.trees.plans.physical.PhysicalStorageLayerAggregate; +import org.apache.doris.nereids.trees.plans.physical.PhysicalTVFRelation; +import org.apache.doris.nereids.trees.plans.physical.PhysicalTopN; +import org.apache.doris.nereids.trees.plans.physical.PhysicalUnion; +import org.apache.doris.nereids.trees.plans.physical.PhysicalWindow; +import org.apache.doris.nereids.trees.plans.physical.PhysicalWorkTableReference; +import org.apache.doris.nereids.trees.plans.visitor.PlanVisitor; + +import com.google.common.annotations.VisibleForTesting; +import com.google.common.collect.Lists; +import org.apache.commons.lang3.StringUtils; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.LinkedHashSet; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.Optional; +import java.util.Set; +import java.util.TreeMap; +import java.util.regex.Matcher; +import java.util.regex.Pattern; +import java.util.stream.Collectors; + +/** + * Physical-plan decompiler (M1). + * + * Decompiles the optimal physical plan tree output by the optimizer into one + * semantically equivalent standard SQL (planSql). planSql freezes the optimizer's key + * decisions (JOIN order, aggregate structure, sort semantics) through SQL structure and + * HINTs, so the plan can be replayed during query rewrite. This is the core engine of + * the SPM "plan-to-SQL" approach. + * + * Design (design doc 6.2): + * + * - Pass-through: physical infrastructure nodes such as PhysicalDistribute and + * PhysicalResultSink have no SQL equivalent and return the child result directly. + * - Merge: a local aggregate (LOCAL) registers the aggregate function names onto the + * child relation without adding nesting. + * - Wrap: Filter / Join / global aggregate / TopN / Project / Window create a new + * SQLRelation and set an alias; parent operators reference it as a subquery through + * toRelationSQL(). + * + * Known M1 simplifications: + * + * - Scan uses the real column names (no c_N normalization; column-name collision + * scenarios are handled in a later milestone). + */ +public class SPMPlan2SQLBuilder extends PlanVisitor<SQLRelation, Void> { + + private static final Logger LOG = LogManager.getLogger(SPMPlan2SQLBuilder.class); + + /** Generated output names ("c_" + ExprId) produced by re-export labels and dedupe renames. */ + private static final Pattern GENERATED_NAME_PATTERN = Pattern.compile("\\bc_\\d+\\b"); + + /** JOIN distribution HINT prefix constants. */ + private static final String HINT_JOIN_BROADCAST = "BROADCAST"; + private static final String HINT_JOIN_SHUFFLE = "SHUFFLE"; + + /** + * Re-export labels a hoist step created one level below ("X AS c_N" items): maps the + * generated label back to the ExprId of the item it re-exports, so a later hoist step + * can re-export the SAME label through one more projection level (tpcds q78: the + * ORDER BY hoisted out of the 16-item projection references c_14/c_15/c_16, and the + * 11-item projection above it must therefore re-export them). + */ + private final Map<String, ExprId> reExportedLabels = new HashMap<>(); + + /** Expression printer (carries the columnNames mapping of SQLRelation). */ + private final SPMExprSqlBuilder exprSqlBuilder = new SPMExprSqlBuilder(); + + /** Per-decompile alias sequence for LATERAL VIEW clauses whose output slot has no + * user qualifier (reset in toSQL). */ + private int lateralViewSeq = 0; + + /** CTE body relations keyed by CTEId (registered by visitPhysicalCTEProducer). The + * consumer references the CTE by its alias instead of inlining the body, so the + * decompiled planSql keeps the WITH structure - one definition shared by every + * consumer - like StarRocks does. */ + private final Map<CTEId, SQLRelation> cteBodies = new HashMap<>(); + + /** CTEId -> the generated alias (t_N) that references the WITH definition. */ + private final Map<CTEId, String> cteAliases = new HashMap<>(); + + /** + * The WITH entries (alias AS (body)) in producer visit order - the definition + * order of the emitted WITH clause. They are attached to the ROOT relation in + * toSQL(): every consumer lies inside the statement, so the definitions are + * visible everywhere regardless of how deeply the CTE anchors are nested in the + * physical tree (attaching at each anchor instead would place a definition inside + * one join branch while a consumer sits in another branch). Producers are visited + * in dependency order - a consumer requires its producer to be registered already - + * so the list is a valid SQL definition order (a CTE used by another CTE is + * defined before it). + */ + private final List<String> cteDefinitions = new ArrayList<>(); + + /** Recursive CTE output column names keyed by CTE name. Registered from the anchor + * branch of a PhysicalRecursiveUnion; used to name the recursive work-table self + * reference and the outer reference so the decompiled WITH RECURSIVE stays + * re-parseable. */ + private final Map<String, List<String>> recursiveCteColumns = new HashMap<>(); + + /** Local (partial) aggregate output column -> the aggregate input expression it + * aggregates, e.g. local partial_sum(x)#N registers N -> x. A global + * aggregate references such a column as sum(partial_sum(x)#N); the global + * decompile rewrites it back to sum(x) so the decompiled SQL stays a single + * logical aggregate (no partial_ intermediate functions). */ + private final Map<ExprId, Expression> localAggParams = new HashMap<>(); + + /** + * Buffer columns produced by the distinct-dedup branch: their defining partial + * aggregate consumed a DATA column (e.g. DISTINCT_LOCAL's partial_count(key)), + * not another partial buffer. In a DISTINCT_GLOBAL merge stage a count(...) over + * such a buffer is the user's count(DISTINCT key). + * + * A plan may mix plain aggregates with a distinct one (TPCDS q28: + * avg(x), count(x), count(DISTINCT x)): the plain count rides along the same + * DISTINCT_GLOBAL stage, but its buffer is a merge chain (its defining partial + * aggregated another partial buffer, e.g. partial_count(partial_count(x))), so it + * must NOT be rendered with DISTINCT. Both chains otherwise resolve to the same + * data column, which is why the provenance has to be recorded while the + * intermediate stages are eliminated. + */ + private final Set<ExprId> distinctMergeBuffers = new HashSet<>(); + + /** The OUTERMOST projection of the decompiled tree (by identity). It alone prunes its + * output to the final user-visible columns; every intermediate projection outputs the + * full child column set plus its own expressions, so an upper layer (filter / join / + * aggregate / projection) can always resolve the columns it references. */ + private final Set<Plan> outputProjects = + Collections.newSetFromMap(new java.util.IdentityHashMap<>()); + + /** + * Identity map: physical plan node -> the ExprIds of its output columns that upper + * layers actually consume (the final result columns plus every column referenced by + * a decompiled expression above). An entry ABSENT, or mapped to null, means "no + * pruning for this node and everything below it" (the pre-pruning behaviour). + * + * The decompiled planSql inflates when every intermediate projection re-emits the + * whole child column set (a deep TPCDS join stack re-lists hundreds of c_N columns + * per level). This map drives live-column pruning: each SELECT list is filtered to + * the live columns only, so a column that is never referenced above and is not part + * of the final result disappears from every intermediate projection. Equivalence is + * preserved because only dead columns are dropped. + */ + private final Map<Plan, Set<ExprId>> neededOutputs = + new java.util.IdentityHashMap<>(); + + /** + * CTEId -> the CTE body output columns referenced by every inlined consumer (the + * consumer output slots mapped back to their producer slots). Filled while the main + * query is propagated top-down; the CTE producer subtree is propagated afterwards + * with exactly this column set, so a big CTE body (a deep TPCDS join stack) only + * keeps the columns its consumers actually read. + */ + private final Map<CTEId, Set<ExprId>> cteConsumerNeeds = new HashMap<>(); + + /** + * Per-decompile generated-column-name state: when a decompiled output column has no + * clean user alias it is exported under a generated c_N name (c_1, c_2, ...), + * assigned in decompile order and memoized by the output ExprId so every reference + * to the same column prints the same alias. A fresh per-decompile sequence (instead + * of embedding the analyzer's ExprId) keeps the generated names small, readable and + * independent of how many internal ExprIds the optimizer allocated. + */ + private final Map<ExprId, String> generatedColumnNames = new HashMap<>(); + private int generatedColumnSeq = 0; + + /** + * Normalized names already visible in the CURRENT decompile: user aliases and the + * preserved source columns of the relations a generated name gets exported next to. + * A generated c_N must avoid them - the counter alone only keeps the + * GENERATED names apart, so a source column literally named c_1 (or a group-by + * item / user alias of that name) could end up next to sum(v) AS c_1 or + * MARK_SLOT c_1: two ExprIds registered under one visible name make the + * enclosing projection / result sink read an AMBIGUOUS column from the derived + * relation after reload. Project / window outputs repair duplicates afterwards + * (dedupeSelectOutputNames); the join (MARK_SLOT / explicit projection) and aggregate + * exports do not, so they reserve here instead. + */ + private final Set<String> reservedOutputNames = new HashSet<>(); + + /** + * Rejects freezing when any expression of the plan carries a + * SessionVarGuardExpr (see containsSessionVarGuard): the guard holds the + * alias-UDF DEFINITION's saved session variables and has no SQL rendering. + * + * @param plan the physical plan about to be decompiled + */ + public static void rejectSessionVarGuardedExpressions(Plan plan) { + org.apache.doris.nereids.spm.SPMPlanTreeSupport.<RuntimeException>walkPlans(plan, node -> { + for (org.apache.doris.nereids.trees.expressions.Expression expr + : node.getExpressions()) { + if (containsSessionVarGuard(expr)) { + throw new UnsupportedOperationException("SPM cannot freeze an expression" + + " carrying a session-variable guard: " + expr); + } + } + }); + } + + private static boolean containsSessionVarGuard( + org.apache.doris.nereids.trees.expressions.Expression expr) { + if (expr instanceof org.apache.doris.nereids.trees.expressions.SessionVarGuardExpr) { + return true; + } + for (org.apache.doris.nereids.trees.expressions.Expression child : expr.children()) { + if (containsSessionVarGuard(child)) { + return true; + } + } + return false; + } + + /** + * Marks the OUTERMOST projection of the decompiled tree: the first PhysicalProject + * reached from the root along single-child pass-through nodes (ResultSink / Sort / + * Distribute / ...). Walking stops at an aggregate - an aggregate that feeds the + * result carries the final SELECT list itself, so there is no outermost projection + * and every projection below it is an intermediate one (full child output). + */ + private void markOutputProject(Plan plan) { + Plan current = plan; + while (current != null) { + if (current instanceof PhysicalProject) { + outputProjects.add(current); + return; + } + if (current instanceof PhysicalHashAggregate || current.arity() != 1) { + return; + } + current = current.child(0); + } + } + + // ==================== live-column analysis (dead-column pruning) ==================== + + /** Whether this subtree may be pruned by the live-column analysis. Non-linear + * decompile shapes (CTE bodies, set operations, recursive unions, grouping sets) + * are left untouched: their select lists are kept whole. */ + private static boolean isPrunableNode(Plan node) { + return node instanceof PhysicalProject + || node instanceof PhysicalWindow + || node instanceof PhysicalHashJoin + || node instanceof PhysicalNestedLoopJoin + || node instanceof PhysicalHashAggregate + || node instanceof PhysicalFilter + || node instanceof PhysicalTopN + || node instanceof PhysicalQuickSort + || node instanceof PhysicalLimit + || node instanceof PhysicalDistribute + || node instanceof PhysicalLazyMaterialize + || node instanceof PhysicalLazyMaterializeOlapScan + || node instanceof PhysicalPartitionTopN + || node instanceof PhysicalAssertNumRows + || node instanceof PhysicalResultSink + || node instanceof PhysicalCTEAnchor; + } + + /** Pre-analysis entry: fills neededOutputs top-down from the root. */ + private void computeNeeded(Plan root) { + neededOutputs.clear(); + cteConsumerNeeds.clear(); + Set<ExprId> rootNeed = new HashSet<>(); + for (Slot slot : root.getOutput()) { + rootNeed.add(slot.getExprId()); + } + propagateNeed(root, rootNeed); + } + + /** + * Top-down live-column propagation. need is the set of output ExprIds of + * node that upper layers consume; null means "keep everything" (never prune + * below). Every parent hands each child the child's output columns that must stay + * alive: the columns the parent passes through and the columns the parent's own + * decompiled expressions reference. + */ + private void propagateNeed(Plan node, Set<ExprId> need) { + if (node == null) { + return; + } + if (node instanceof PhysicalCTEConsumer) { + // a CTE consumer has no children; record which body columns its live output + // columns map back to, so the producer subtree can be pruned to exactly them + if (need != null) { + PhysicalCTEConsumer consumer = (PhysicalCTEConsumer) node; + Set<ExprId> bodyNeeds = cteConsumerNeeds.computeIfAbsent( + consumer.getCteId(), k -> new HashSet<>()); + for (Slot slot : consumer.getOutput()) { + if (need.contains(slot.getExprId())) { + bodyNeeds.add(consumer.getProducerSlot(slot).getExprId()); + } + } + } + return; + } + if (need == null || !isPrunableNode(node)) { + // keep everything on this node and below (no pruning boundary) + for (Plan child : node.children()) { + propagateNeed(child, null); + } + return; + } + Set<ExprId> nodeNeed = new HashSet<>(need); + neededOutputs.put(node, nodeNeed); + + if (node instanceof PhysicalCTEAnchor) { + // child(0) is the CTE producer whose body every consumer inlines. The main + // query (child(1)) is propagated FIRST so every consumer records which body + // columns it reads; the body subtree is then pruned to exactly those columns. + // A body with no consumer (or one whose consumers never resolve) keeps every + // column. + PhysicalCTEAnchor<? extends Plan, ? extends Plan> anchor = + (PhysicalCTEAnchor<? extends Plan, ? extends Plan>) node; + if (anchor.child(1) != null) { + propagateNeed(anchor.child(1), nodeNeed); + } + Plan producer = anchor.child(0); + if (producer == null) { + return; + } + if (producer instanceof PhysicalCTEProducer) { + PhysicalCTEProducer<? extends Plan> cteProducer = + (PhysicalCTEProducer<? extends Plan>) producer; + Set<ExprId> bodyNeed = cteConsumerNeeds.get(cteProducer.getCteId()); + if (bodyNeed != null && !bodyNeed.isEmpty()) { + propagateNeed(cteProducer.child(0), bodyNeed); + } else { + propagateNeed(cteProducer.child(0), null); + } + } else { + propagateNeed(producer, null); + } + return; + } + if (node instanceof PhysicalProject) { + propagateProjectNeed((PhysicalProject<? extends Plan>) node, nodeNeed); + } else if (node instanceof PhysicalHashJoin || node instanceof PhysicalNestedLoopJoin) { + propagateJoinNeed((AbstractPhysicalJoin<? extends Plan, ? extends Plan>) node, nodeNeed); + } else if (node instanceof PhysicalHashAggregate) { + propagateAggregateNeed((PhysicalHashAggregate<? extends Plan>) node, nodeNeed); + } else if (node instanceof PhysicalWindow) { + propagateWindowNeed((PhysicalWindow<? extends Plan>) node, nodeNeed); + } else if (node instanceof PhysicalFilter) { + PhysicalFilter<? extends Plan> filter = (PhysicalFilter<? extends Plan>) node; + Set<ExprId> childNeed = new HashSet<>(nodeNeed); + addExprSlots(filter.getPredicate(), childNeed); + propagateNeed(filter.child(0), childNeed); + } else if (node instanceof PhysicalTopN) { + PhysicalTopN<? extends Plan> topN = (PhysicalTopN<? extends Plan>) node; + Set<ExprId> childNeed = new HashSet<>(nodeNeed); + for (org.apache.doris.nereids.properties.OrderKey key : topN.getOrderKeys()) { + addExprSlots(key.getExpr(), childNeed); + } + propagateNeed(topN.child(0), childNeed); + } else if (node instanceof PhysicalQuickSort) { + PhysicalQuickSort<? extends Plan> sort = (PhysicalQuickSort<? extends Plan>) node; + Set<ExprId> childNeed = new HashSet<>(nodeNeed); + for (org.apache.doris.nereids.properties.OrderKey key : sort.getOrderKeys()) { + addExprSlots(key.getExpr(), childNeed); + } + propagateNeed(sort.child(0), childNeed); + } else { + // pure pass-through nodes (limit / distribute / lazy materialize / ...): + // their output slots are the child slots, so the same need flows down + for (Plan child : node.children()) { + propagateNeed(child, nodeNeed); + } + } + } + + /** PhysicalProject: live pass-through columns plus the columns referenced by the + * project's own expressions whose output is itself live. */ + private void propagateProjectNeed(PhysicalProject<? extends Plan> project, Set<ExprId> need) { + Plan child = project.child(0); + if (child == null) { + return; + } + Set<ExprId> childNeed = new HashSet<>(); + boolean finalProject = outputProjects.contains(project); + if (finalProject) { + // the outermost projection emits its full projection list: every referenced + // child column must stay alive + for (NamedExpression projectExpr : project.getProjects()) { + addExprSlots(projectExpr, childNeed); + } + } else { + // pass-through: the child columns that are still live above this projection + Set<ExprId> childOut = outputIdSet(child); + for (ExprId id : need) { + if (childOut.contains(id)) { + childNeed.add(id); + } + } + // plus the columns referenced by this projection's own live expressions + for (NamedExpression projectExpr : project.getProjects()) { + if (need.contains(projectExpr.getExprId())) { + addExprSlots(projectExpr, childNeed); + } + } + } + propagateNeed(child, childNeed); + } + + /** Join: the preserved output columns flow to the side that produces them; the + * hash / other / mark conjuncts are always decompiled, so every column they + * reference stays alive on its own side. */ + private void propagateJoinNeed(AbstractPhysicalJoin<? extends Plan, ? extends Plan> join, + Set<ExprId> need) { + Plan left = join.left(); + Plan right = join.right(); + if (left == null || right == null) { + return; + } + Set<ExprId> leftOut = outputIdSet(left); + Set<ExprId> rightOut = outputIdSet(right); + Set<ExprId> conjRefs = new HashSet<>(); + for (Expression conjunct : join.getHashJoinConjuncts()) { + addExprSlots(conjunct, conjRefs); + } + for (Expression conjunct : join.getOtherJoinConjuncts()) { + addExprSlots(conjunct, conjRefs); + } + for (Expression conjunct : join.getMarkJoinConjuncts()) { + addExprSlots(conjunct, conjRefs); + } + Set<ExprId> leftNeed = new HashSet<>(); + Set<ExprId> rightNeed = new HashSet<>(); + for (ExprId id : need) { + if (leftOut.contains(id)) { + leftNeed.add(id); + } + if (rightOut.contains(id)) { + rightNeed.add(id); + } + } + for (ExprId id : conjRefs) { + if (leftOut.contains(id)) { + leftNeed.add(id); + } + if (rightOut.contains(id)) { + rightNeed.add(id); + } + } + propagateNeed(left, leftNeed); + propagateNeed(right, rightNeed); + } + + /** Aggregate: the global stage keeps its whole output list (GROUP BY keys and + * aggregate functions), but the input columns its expressions reference must stay + * alive below. Local / intermediate execution stages are pass-through nodes. */ + private void propagateAggregateNeed(PhysicalHashAggregate<? extends Plan> agg, Set<ExprId> need) { + AggPhase phase = agg.getAggPhase(); + if (phase.isLocal() || isIntermediateAggStage(agg)) { + for (Plan child : agg.children()) { + propagateNeed(child, need); + } + return; + } + Plan child = agg.child(0); + if (child == null) { + return; + } + Set<ExprId> childNeed = new HashSet<>(); + for (Expression groupBy : agg.getGroupByExpressions()) { + addExprSlots(groupBy, childNeed); + } + for (NamedExpression output : agg.getOutputExpressions()) { + addAggregateOutputRefs(output, childNeed); + } + propagateNeed(child, childNeed); + } + + /** Window: live pass-through columns plus the input columns of the live window + * expressions. */ + private void propagateWindowNeed(PhysicalWindow<? extends Plan> window, Set<ExprId> need) { + Plan child = window.child(0); + if (child == null) { + return; + } + Set<ExprId> childNeed = new HashSet<>(); + Set<ExprId> childOut = outputIdSet(child); + for (ExprId id : need) { + if (childOut.contains(id)) { + childNeed.add(id); + } + } + for (NamedExpression windowExpr : window.getWindowExpressions()) { + if (need.contains(windowExpr.getExprId())) { + addExprSlots(windowExpr, childNeed); + } + } + propagateNeed(child, childNeed); + } + + /** The columns referenced by one global-aggregate output expression, with every + * local partial-buffer reference resolved down to its data columns (mirrors + * appendAggSelect / resolveBufferSlots so the pruned SELECT lists keep exactly the + * columns the decompiled aggregate prints). */ + private void addAggregateOutputRefs(NamedExpression output, Set<ExprId> out) { + Expression inner = output instanceof Alias ? ((Alias) output).child() : output; + if (inner instanceof SlotReference) { + // group-by key pass-through + out.add(((SlotReference) inner).getExprId()); + return; + } + if (inner instanceof AggregateExpression) { + AggregateExpression aggExpr = (AggregateExpression) inner; + List<Expression> args = aggExpr.getFunction().children().isEmpty() + ? new ArrayList<>(aggExpr.children()) : aggExpr.getFunction().children(); + for (Expression arg : args) { + addResolvedExprSlots(arg, out); + } + return; + } + addExprSlots(inner, out); + } + + /** Collects every slot of expr, resolving partial-buffer slots through + * localAggParams down to their data columns (mirrors resolveBufferSlots). */ + private void addResolvedExprSlots(Expression expr, Set<ExprId> out) { + if (expr instanceof SlotReference) { + Expression param = localAggParams.get(((SlotReference) expr).getExprId()); + if (param != null) { + addResolvedExprSlots(param, out); + return; + } + out.add(((SlotReference) expr).getExprId()); + return; + } + for (Expression childExpr : expr.children()) { + addResolvedExprSlots(childExpr, out); + } + } + + /** Output ExprId set of a plan node. */ + private static Set<ExprId> outputIdSet(Plan node) { + Set<ExprId> ids = new HashSet<>(); + for (Slot slot : node.getOutput()) { + ids.add(slot.getExprId()); + } + return ids; + } + + /** Adds every slot ExprId used by an expression. */ + private static void addExprSlots(Expression expr, Set<ExprId> out) { + if (expr == null) { + return; + } + collectAllSlotIds(expr, out); + } + + /** Pre-walk that fills localAggParams bottom-up so the live-column analysis + * (which runs before the decompile walk) can resolve partial-buffer references. The + * decompile walk re-fills the same map while it descends (harmless duplicate). */ + private void collectLocalAggParams(Plan node) { + for (Plan child : node.children()) { + collectLocalAggParams(child); + } + if (node instanceof PhysicalHashAggregate) { + PhysicalHashAggregate<? extends Plan> agg = (PhysicalHashAggregate<? extends Plan>) node; + if (agg.getAggPhase().isLocal() || isIntermediateAggStage(agg)) { + recordLocalAggStage(agg); + } + } + } + + /** Registers one local (partial) aggregate stage's buffer outputs (see the field + * comment of localAggParams). */ + private void recordLocalAggStage(PhysicalHashAggregate<? extends Plan> agg) { + for (NamedExpression output : agg.getOutputExpressions()) { + Expression inner = output instanceof Alias ? ((Alias) output).child() : output; + if (inner instanceof AggregateExpression && isPartialAggregate((AggregateExpression) inner)) { + Expression rawParam = extractPartialParam((AggregateExpression) inner); + if (agg.getGroupByExpressions().isEmpty() + && isDistinctMergeContribution(((AggregateExpression) inner).children())) { + distinctMergeBuffers.add(output.getExprId()); + } + Expression param = rawParam; + if (param == null) { + // count(*) has no data argument (extractPartialParam -> null). Map the + // buffer to the partial expression itself so the enclosing + // merge-finalize count(*) resolves its argument back to the star, which + // appendAggSelect collapses into count(*) (isNestedNoArgCount) - without + // this the raw buffer slot (named e.g. "partial_count(*)") would leak + // into the decompiled SQL as the invalid count(partial_count(*)). + param = inner; + } + while (param instanceof SlotReference) { + Expression resolved = localAggParams.get(((SlotReference) param).getExprId()); + if (resolved == null) { + break; + } + param = resolved; + } + localAggParams.put(output.getExprId(), param); + } + } + } + + /** + * Whether an output of an eliminated DISTINCT_LOCAL stage is the distinct-dedup + * contribution (see distinctMergeBuffers). + * + * The distinction has to be read from the aggregate expression's own children - + * the physical input of the buffer - because the merge function's argument is + * normalized to the data column for EVERY buffer: + * + * - merge-chain buffer (plain aggregate riding along): its input is another + * partial buffer, already recorded in localAggParams, or a count-star function + * with no argument at all; + * - distinct-dedup buffer: its input is the dedup KEY data column, either as the + * bare slot (partial_count(key)) or wrapped as the partial function's argument + * (count(key)). + */ + private boolean isDistinctMergeContribution(List<Expression> exprChildren) { + if (exprChildren.isEmpty()) { + return false; + } + Expression bufferArg = exprChildren.get(0); + if (bufferArg instanceof SlotReference) { + return !localAggParams.containsKey(((SlotReference) bufferArg).getExprId()); + } + List<Expression> argChildren = bufferArg.children(); + if (argChildren.isEmpty()) { + // e.g. the count(*) star buffer: raw-row counting, never the distinct merge + return false; + } + for (Expression child : argChildren) { + if (!(child instanceof SlotReference) + || localAggParams.containsKey(((SlotReference) child).getExprId())) { + return false; + } + } + return true; + } + + /** The live column filter for one explicit SELECT list: the list a node emits is + * pruned to the node's needed output columns (no entry in the map, or a null need, + * keeps the whole list - e.g. mock plan trees and unpruned subtrees). */ + private void filterLiveSelects(Plan node, List<Pair<ExprId, String>> selects) { + Set<ExprId> need = neededOutputs.get(node); + if (need == null) { + return; + } + selects.removeIf(p -> !need.contains(p.key())); + } + + /** + * Decompile entry: physical plan -> planSql. + * + * @param plan the optimal physical plan + * @return planSql (standard SQL text) + */ + public String toSQL(Plan plan) { + // An alias-UDF expansion computed under the DEFINITION's stored session + // variables (decimalOverflowScale / enable_decimal256 / ...) carries a + // SessionVarGuardExpr; the SQL text has no way to express that guard, so + // printing only its child would replan the arithmetic under the LATER caller's + // variables and could change its type, scale or value while the same bind SQL + // still matches. Reject freezing instead (CREATE keeps the user planSql / + // falls back to the parameterized-plan-tree path). + rejectSessionVarGuardedExpressions(plan); + // reset the per-decompile alias / generated-name sequences: t_N and c_N only + // need to be unique WITHIN the one produced SQL, so numbering restarts here and + // the decompiled text stays compact across calls + SQLRelation.resetAliasCounter(); + cteBodies.clear(); + cteAliases.clear(); + cteDefinitions.clear(); + generatedColumnNames.clear(); + generatedColumnSeq = 0; + reservedOutputNames.clear(); + lateralViewSeq = 0; + outputProjects.clear(); + localAggParams.clear(); + distinctMergeBuffers.clear(); + markOutputProject(plan); + collectLocalAggParams(plan); + computeNeeded(plan); + SQLRelation relation = plan.accept(this, null); + // attach the collected CTE definitions (WITH) to the outermost relation: they + // are collected in producer dependency order and are visible to the whole + // statement, which is exactly the scope a CTE anchor binds + if (!cteDefinitions.isEmpty()) { + List<String> merged = relation.getCte() == null + ? new ArrayList<>() : new ArrayList<>(relation.getCte()); + merged.addAll(cteDefinitions); + relation.setCte(merged); + } + // A top-level ASSERT_ROWS (e.g. EXISTS / single-row assertion) has no SQL + // representation and is rejected by visitPhysicalAssertNumRows; the decompile then + // falls back to the user planSql. + return relation.toSQL(); + } + + /** + * Handles a child node. + * + * @param plan child physical plan + * @return SQLRelation of the child plan + */ + private SQLRelation process(Plan plan) { + return plan.accept(this, null); + } + + // ==================== default: unsupported operators ==================== + + @Override + public SQLRelation visit(Plan plan, Void context) { + throw new UnsupportedOperationException( + "SPMPlan2SQLBuilder does not support plan: " + plan.getClass().getSimpleName()); + } + + // ==================== pass-through mode ==================== + /** + * PhysicalStorageLayerAggregate: a scan-shaped shortcut that returns COUNT / MIN / MAX + * for the wrapped scan from table or footer metadata instead of reading the data. + * AggregateStrategies keeps the enclosing aggregate (or the constant-only project) on + * TOP of the shortcut and keeps every expression in terms of the wrapped relation's + * slots, so decompiling the wrapped relation publishes exactly those slots and the + * enclosing operators render the same SQL aggregate (count(*) / min(x) / ...) over the + * same table. Replaying that SQL re-derives an equivalent plan; there is no clause that + * describes the shortcut itself and none is needed, because the shortcut is an + * execution strategy of the aggregate rather than a different result. + */ + @Override + public SQLRelation visitPhysicalStorageLayerAggregate( + PhysicalStorageLayerAggregate storageLayerAggregate, Void context) { + return visitPhysicalRelation(storageLayerAggregate.getRelation(), context); + } + + /** + * PhysicalDistribute: a data exchange node with no SQL equivalent; returns the child. + */ + @Override + public SQLRelation visitPhysicalDistribute(PhysicalDistribute<? extends Plan> distribute, Void context) { + return process(distribute.child(0)); + } + + /** + * PhysicalLazyMaterialize: an optimizer lazy-column-materialization wrapper with no + * SQL equivalent; returns the child. + */ + @Override + public SQLRelation visitPhysicalLazyMaterialize( + PhysicalLazyMaterialize<? extends Plan> lazy, Void context) { + return process(lazy.child(0)); + } + + /** + * PhysicalLazyMaterializeOlapScan: an OlapScan wrapped with lazy column + * materialization; decompiled as a normal scan. + */ + @Override + public SQLRelation visitPhysicalLazyMaterializeOlapScan( + PhysicalLazyMaterializeOlapScan scan, Void context) { + return visitPhysicalRelation(scan, context); + } + + /** + * PhysicalLazyMaterializeFileScan: a file scan wrapped with lazy column + * materialization; decompiled as a normal scan (including its scan modifiers - + * the wrapper subclasses PhysicalFileScan, so the TABLESAMPLE / snapshot / scan + * parameter rendering applies unchanged). + */ + @Override + public SQLRelation visitPhysicalLazyMaterializeFileScan( + PhysicalLazyMaterializeFileScan scan, Void context) { + return visitPhysicalRelation(scan, context); + } + + /** + * PhysicalResultSink (and sinks in general): result collection nodes with no SQL + * equivalent; passed through. For the RESULT sink the final output column ORDER is + * forced to the sink's output list (the user SELECT order) - the physical aggregate + * / projection below may emit columns in a different order (e.g. group-by keys + * before aggregates), which would reorder the user-visible result columns. + */ + @Override + public SQLRelation visitPhysicalSink(PhysicalSink<? extends Plan> sink, Void context) { + if (sink instanceof PhysicalResultSink) { + SQLRelation child = process(sink.child(0)); + List<Slot> outputs = sink.getOutput(); + if (outputs.isEmpty()) { + return child; + } + // Final output columns in the USER SELECT order. Each output column is + // emitted as its in-scope reference (child.getColumnNames()), and - when the + // reference is a decompile-internal alias (c_N) that hides the original + // column label - re-labelled with the ResultSink output slot's name + // (output.getName()): that label is exactly what a normal execution of the + // original query shows in the result header (a plain alias such as + // "d_week_seq1", a table column name, or the expression text of an + // un-aliased expression column), so the frozen planSql keeps the SAME output + // column names as the original SQL (SR model: the final SELECT list is + // driven by the logical query's output columns, not by the internal c_N + // aliases the decompiler had to mint for unambiguous references). + List<Pair<ExprId, String>> ordered = new ArrayList<>(); + boolean anyRenamed = false; + for (Slot output : outputs) { + String ref = child.getColumnNames().get(output.getExprId()); + if (ref == null) { + ref = output.getName(); + } + String display = output.getName(); + if (display == null || display.isEmpty() || display.equals(ref)) { + ordered.add(Pair.of(output.getExprId(), ref)); + continue; + } + // Re-expose the column under its original label: "c_3 AS d_week_seq1" + // (or "AS `round(...)`" for an un-aliased expression column, whose header + // text is the expression itself). The quoted name never participates in + // name resolution of the frozen SQL body, so quoting is always safe. + ordered.add(Pair.of(output.getExprId(), ref + " AS " + quoteIdentifier(display))); + anyRenamed = true; + } + if (child.getRelationName() == null) { + // inline relation (plain table scan): SELECT * renders in the table + // schema order and the ResultSink output follows the scan order; + // nothing to reorder. + return child; + } + if (child.getSelects().isEmpty()) { + // SELECT * over a wrapped subquery (TopN/Sort wrapper whose FROM is a + // subquery): the columns are referenceable, so rewrite the outer SELECT + // list to the user column order in place (preserving ORDER BY / LIMIT) - + // UNLESS a re-labelled alias would SHADOW a name the ORDER BY references + // "c_3 AS b" next to "ORDER BY b" rebinds the sort to the + // alias although the clause was rendered against the base column b. Then + // the relabel moves to the wrapper below, where the clause keeps its own + // query block and the aliases cannot capture it. + String inPlaceOrderBy = child.getOrderBy(); + if (inPlaceOrderBy.isEmpty() + || !orderByShadowedByRelabel(child, inPlaceOrderBy, ordered)) { + child.setSelects(ordered); + return child; + } + } + // Explicit projection already emitted (e.g. a bare aggregate + // "sum(...) AS revenue" or a pass-through projection). Overwriting it in + // place would replace the expression with a bare name that is not resolvable + // in the FROM scope, so: + // - if the emitted columns already carry the original labels in the user + // SELECT order, keep the projection untouched; + // - otherwise wrap the child in a subquery and select the ordered + // (re-labelled) columns from it. + if (!anyRenamed && sameExprIdOrder(child.getSelects(), ordered)) { + return child; + } + SQLRelation outer = new SQLRelation(); + String alias = outer.newAlias(); + // A TOP-LEVEL ORDER BY / LIMIT pair is a semantic user clause: matching + // deliberately ignores its LIMIT value, and mergeLimits adopts the user's value + // by position - but a limit left INSIDE the derived table can never be reached + // by that positional merge (the frozen root would have no Limit node next to + // the user's Limit), so a matched LIMIT 200 query kept the captured 100-row + // cap. This wrapper is a PURE output relabelling (it neither filters nor + // aggregates), so hoisting the pair onto it is semantically identical and + // makes the frozen root a Limit node again. + // The ORDER BY moves WITH the LIMIT: a derived-table ORDER BY WITHOUT its + // LIMIT is only a hint the optimizer is free to drop, and the outer SELECT + // would then return an arbitrary LIMIT slice (a q02-style query replayed as + // unordered rows). Clear both BEFORE rendering the child text - toSQL() + // captures them into the string. + String sinkOrderBy = child.getOrderBy(); + String sinkLimit = child.getLimit(); + if (!sinkOrderBy.isEmpty() || !sinkLimit.isEmpty()) { + // Independently of the LIMIT: "SELECT a FROM t ORDER BY b + 1" + // (no LIMIT at all) arrives here as a Sort under the final relabel wrapper, + // and leaving the ORDER BY inside the derived table was NOT safe - the + // frozen text is re-planned at replay, and a sort that the user asked for + // but that sits one level down came back as unordered rows (Nereids drops + // an inner sort it considers redundant). The wrapper is a pure output + // relabelling, so the user clause belongs on it. Hidden sort keys the + // child's SELECT list does not export are added to it first (see + // hoistableOrderBy); a key that cannot be exported, or that a relabelled + // wrapper alias would capture, keeps the clause inside the child - its own + // query block still resolves it. ORDER BY and LIMIT + // always move TOGETHER: a limit without its order picks arbitrary rows. + String hoisted = sinkOrderBy.isEmpty() ? "" + : hoistableOrderBy(child, sinkOrderBy, ordered, true); + if (hoisted != null) { + outer.setOrderBy(hoisted); + outer.setLimit(sinkLimit); + child.setOrderBy(""); + child.setLimit(""); + } + } + outer.setFrom("(" + child.toSQL() + ") " + alias); + outer.setSelects(ordered); + return outer; + } + return visit((Plan) sink, context); + } + + /** Whether the projection ExprIds appear in the same order as ordered. */ + private static boolean sameExprIdOrder(List<Pair<ExprId, String>> projection, + List<Pair<ExprId, String>> ordered) { + if (projection.size() != ordered.size()) { + return false; + } + for (int i = 0; i < projection.size(); i++) { + if (projection.get(i).key() != ordered.get(i).key()) { + return false; + } + } + return true; + } + + /** + * Quotes an identifier for use as executable SQL text whenever it is not a plain + * unquotable identifier: a column named a-b must be emitted as `a-b`, + * otherwise the frozen projection re-parses as the subtraction a - b. A name that + * looks plain can still be a RESERVED keyword (a legal quoted column named + * from must not be emitted bare - the parser tokenizes that as the FROM + * keyword rather than an identifier), so the keyword check decides as well. + * Embedded backticks are doubled. Unquotable names stay verbatim, so ordinary + * schemas keep byte-identical frozen SQL. Also used for result-column labels (the + * expression text of an un-aliased output column such as + * round((sun_sales1 / sun_sales2), 2) is wrapped so the frozen planSql can + * carry the original column header verbatim). + */ + public static String quoteIdentifier(String name) { + if (name == null) { + return null; + } + if (name.matches("[A-Za-z_][A-Za-z0-9_]*") && NereidsParser.isValidUnquotedIdentifier(name)) { + return name; + } + return "`" + name.replace("`", "``") + "`"; + } + + /** + * Quotes every dot-separated component of a (possibly) qualified metadata name + * (catalog.db.table), so a table whose name is not a plain identifier (`my-table`) + * is still emitted as an identifier reference rather than as an expression. + */ + static String quoteQualifiedName(String name) { + if (name == null || name.isEmpty()) { + return name; + } + String[] parts = name.split("\\.", -1); + StringBuilder sb = new StringBuilder(name.length() + 4); + for (int i = 0; i < parts.length; i++) { + if (i > 0) { + sb.append('.'); + } + sb.append(quoteIdentifier(parts[i])); + } + return sb.toString(); + } + + /** + * Fully-qualified name of a table with every COMPONENT quoted separately + * (catalog.db.table). getNameWithFullQualifiers() FLATTENS a legal quoted + * component such as `t.a` into "internal.db.t.a", and splitting that on every dot + * would emit FOUR identifiers instead of the intended three-part name with `t.a` + * quoted as ONE component - a manually created frozen baseline carrying that broken + * text then fails re-analysis after a reload, and no raw fallback tree exists for a + * frozen row. The metadata components are taken structurally, never re-split. + */ + static String quoteQualifiedTableName(org.apache.doris.catalog.TableIf table) { + StringBuilder sb = new StringBuilder(); + for (String component : table.getFullQualifiers()) { + if (sb.length() > 0) { + sb.append('.'); + } + sb.append(quoteIdentifier(component)); + } + return sb.toString(); + } + + // ==================== Scan (wrapped as subquery or inline) ==================== + + /** + * PhysicalRelation (PhysicalOlapScan / PhysicalFileScan, etc.): table scan. + * + * M1 simplification: the scan output columns are registered to columnNames by their + * real column names and from is inlined as the table name (not wrapped). The + * predicate is handled by the parent PhysicalFilter. + */ + @Override + public SQLRelation visitPhysicalRelation(PhysicalRelation relation, Void context) { + if (!(relation instanceof PhysicalCatalogRelation)) { + throw new UnsupportedOperationException( + "SPMPlan2SQLBuilder does not support relation: " + relation.getClass().getSimpleName()); + } + PhysicalCatalogRelation catalogRelation = (PhysicalCatalogRelation) relation; + // A TEMPORARY table's catalog object carries the CREATOR session's internal name + // (<sessionId>_#TEMP#_<name>), not the text the user typed. Emitting it into the + // frozen planSql would persist a session-scoped identifier: another session + // running the same "FROM t" text matches the same bind key (catalog.db.t), and + // Database.getTableNullable accepts the already-marked internal name unchanged, + // so the replay would read the CREATOR's (possibly still live) temporary table + // instead of its own t. Never render it - the GLOBAL create / capture path + // rejects temporary relations before freezing (SPMPlanner), and a SESSION-scope + // freeze (where the session's own temp table is the intended target) falls back + // to the raw user text instead. + if (catalogRelation.getTable().isTemporary()) { + throw new UnsupportedOperationException( + "SPMPlan2SQLBuilder does not support temporary tables: the physical relation" + + " carries the creator session's internal name"); + } + SQLRelation sqlRelation = new SQLRelation(); + // Emit the fully qualified name (catalog.db.table) so the frozen planSql resolves + // the same table when it is replayed from a session whose current database (or + // catalog) differs from the one used at CREATE time (cross-db queries, + // information_schema, ...). Tables without a database (e.g. FunctionGenTable) + // keep the bare name. Every COMPONENT is backtick-quoted separately when it is + // not a plain identifier, so a metadata name containing operators is re-parsed + // as an identifier instead of an expression, and a legal quoted component such + // as `t.a` keeps its boundary (see quoteQualifiedTableName). Scan modifiers + // (partition selection, TABLESAMPLE, snapshot, scan parameters) follow the name + // in grammar order. + String table = catalogRelation.getTable().getDatabase() == null + ? quoteIdentifier(catalogRelation.getTable().getName()) + : quoteQualifiedTableName(catalogRelation.getTable()); + sqlRelation.setFrom(table + renderScanModifiers(relation)); + // Remember WHICH table this scan reads: two occurrences of one table are a self + // join even when their FROM texts differ (each occurrence may carry its own + // PARTITION / TABLESAMPLE pin), and the join must wrap both sides then. + sqlRelation.setRelationIdentity(table); + // Register output columns: ExprId -> real column name. Internal system columns + // (e.g. rowid columns a join may request from the scan) are execution details + // and are never registered so they cannot leak into projections / ON clauses. + for (Slot slot : relation.getOutput()) { + if (isSystemColumnName(slot.getName())) { + continue; + } + // registered as executable SQL text: a special-character column (`a-b`) + // must be backtick-quoted, otherwise a frozen projection SELECT a-b + // re-parses as the subtraction a - b and returns a different value + sqlRelation.registerRef(slot.getExprId(), quoteIdentifier(slot.getName())); + } + return sqlRelation; + } + + // ==================== table-valued functions ==================== + + /** + * PhysicalTVFRelation: a table-valued function in FROM (e.g. + * numbers('number' = '5')). The function's own SQL text is used because a TVF + * argument list is a property list, not a normal expression list; the output + * columns keep their names. + */ + @Override + public SQLRelation visitPhysicalTVFRelation(PhysicalTVFRelation tvfRelation, Void context) { + SQLRelation relation = new SQLRelation(); + relation.setFrom(renderTableValuedFunction(tvfRelation.getFunction())); + for (Slot slot : tvfRelation.getOutput()) { + relation.registerRef(slot.getExprId(), quoteIdentifier(slot.getName())); + } + return relation; + } + + /** + * Renders one TVF call safely. TableValuedFunction.computeToSql() concatenates each + * concrete key/value between apostrophes WITHOUT escaping, so a valid property such + * as an S3 object key containing an apostrophe or backslash would freeze malformed + * SQL: with a placeholder-bearing predicate the row is classified as frozen, no + * fallback tree is rebuilt, and after refresh/restart the baseline silently stops + * applying. The property map is serialized in a deterministic key order with the + * same default-mode-safe quoting as every other SPM-emitted string literal (the + * stored text is always re-parsed under MODE_DEFAULT). + */ + private static String renderTableValuedFunction(TableValuedFunction function) { + String args = new TreeMap<>(function.getTVFProperties().getMap()).entrySet().stream() + .map(kv -> quoteSqlString(kv.getKey()) + " = " + quoteSqlString(kv.getValue())) + .collect(Collectors.joining(", ")); + return quoteIdentifier(function.getName()) + "(" + args + ")"; + } + + /** PhysicalLazyMaterializeTVFScan: a TVF scan wrapped by lazy materialization. */ + @Override + public SQLRelation visitPhysicalLazyMaterializeTVFScan( + PhysicalLazyMaterializeTVFScan scan, Void context) { + return visitPhysicalTVFRelation(scan, context); + } + + /** + * Renders the scan modifiers that the frozen SQL must carry, in the order the grammar + * accepts them after the table name + * (optScanParams, materializedViewName, tableSnapshot, specifiedPartition, sample): + * + * olap scans: @paramType(...) parameters (binlog reads), the partition list + * when the scan reads a strict non-empty subset of the table partitions, and + * TABLESAMPLE; + * file scans: @paramType(...) parameters, FOR VERSION/TIME AS OF + * snapshots and TABLESAMPLE. + * + * Dropping any of them would silently change what the frozen planSql reads: a + * "FROM t PARTITION(p1)" baseline would replay over every partition, a snapshot read + * would run against the moving head of the table, and a sampled scan would return the + * full table. + * + * The partition list is emitted only when the user pinned it (a pruned scan re-derives + * its selection by replay); a frozen text without the pin simply does not match the + * pinned user query (a safe miss, never a replay over the wrong partitions). + * + * Two scan states are deliberately NOT emitted. The selected index is an optimizer + * choice (a rollup is a consistent copy, so the choice carries no semantics) that + * cannot be told apart from a user-written INDEX clause; freezing it would pin the + * optimization and stop the plain user query from matching. A user-written INDEX + * pin only survives in the user tree, so such queries never match a frozen text that + * dropped the pin. + */ + private static String renderScanModifiers(PhysicalRelation relation) { + StringBuilder modifiers = new StringBuilder(); + if (relation instanceof PhysicalOlapScan) { + PhysicalOlapScan scan = (PhysicalOlapScan) relation; + modifiers.append(renderScanParams(scan.getScanParams())); + modifiers.append(renderPartitionSelection(scan)); + modifiers.append(renderTabletSelection(scan)); + modifiers.append(renderTableSample(scan.getTableSample())); + } else if (relation instanceof PhysicalFileScan) { + PhysicalFileScan scan = (PhysicalFileScan) relation; + modifiers.append(renderScanParams(scan.getScanParams())); + modifiers.append(renderTableSnapshot(scan.getTableSnapshot())); + modifiers.append(renderTableSample(scan.getTableSample())); + } + return modifiers.toString(); + } + + /** + * OLAP partition selection: the ids of a user-written PARTITION(...) / + * TEMPORARY PARTITION(...) list are frozen as that very clause - including + * when the pin happens to cover every partition that exists right now (a cardinality + * test would drop the clause and a later ADD PARTITION would let the pinned query + * silently read the new partition) and including the temporary namespace (a temp pin + * replayed as a formal PARTITION(name) binds the wrong partition or fails to bind). + * Partition pruning also shrinks selectedPartitionIds without any user pin, + * so only the manual provenance is ever rendered; a pruned scan re-derives its + * selection by replay. Ids are sorted so the frozen text is deterministic (the + * matcher compares the selection as a set). + */ + private static String renderPartitionSelection(PhysicalOlapScan scan) { + List<Long> pinnedIds = scan.getManuallySpecifiedPartitions(); + if (pinnedIds.isEmpty()) { + return ""; + } + OlapTable table = scan.getTable(); + List<Long> sortedIds = new ArrayList<>(pinnedIds); + Collections.sort(sortedIds); + boolean temporary = table.isTemporaryPartition(sortedIds.get(0)); + StringBuilder partition = new StringBuilder(temporary ? " TEMPORARY PARTITION(" : " PARTITION("); + for (int i = 0; i < sortedIds.size(); i++) { + if (i > 0) { + partition.append(", "); + } + if (table.isTemporaryPartition(sortedIds.get(i)) != temporary) { + throw new UnsupportedOperationException( + "SPM decompile: a partition pin mixing temporary and formal partitions" + + " has no single-clause rendering"); + } + Partition partitionMeta = table.getPartition(sortedIds.get(i)); + partition.append(quoteIdentifier(partitionMeta.getName())); + } + return partition.append(')').toString(); + } + + /** + * OLAP tablet pin: TABLET(id, ...) is frozen when the user wrote it. Bucket + * pruning also fills selectedTabletIds, so only the manual provenance is + * rendered; without the clause a replayed scan would read every tablet of the + * selected partitions and could return rows the captured query excluded. Ids are + * sorted (the matcher compares the tablet list as a set). + */ + private static String renderTabletSelection(PhysicalOlapScan scan) { + List<Long> pinnedIds = scan.getManuallySpecifiedTabletIds(); + if (pinnedIds.isEmpty()) { + return ""; + } + List<Long> sortedIds = new ArrayList<>(pinnedIds); + Collections.sort(sortedIds); + StringBuilder tablet = new StringBuilder(" TABLET("); + for (int i = 0; i < sortedIds.size(); i++) { + if (i > 0) { + tablet.append(", "); + } + tablet.append(sortedIds.get(i)); + } + return tablet.append(')').toString(); + } + + /** + * TABLESAMPLE(n PERCENT | n ROWS) [REPEATABLE seed]: both olap and file scans keep the + * user's sample. Dropping it would replay over the full table and return rows the + * captured plan never sampled in. + */ + private static String renderTableSample(Optional<TableSample> tableSample) { + if (!tableSample.isPresent()) { + return ""; + } + TableSample sample = tableSample.get(); + StringBuilder modifiers = new StringBuilder(" TABLESAMPLE(") + .append(sample.sampleValue) + .append(sample.isPercent ? " PERCENT)" : " ROWS)"); + if (sample.seek >= 0) { + modifiers.append(" REPEATABLE ").append(sample.seek); + } + return modifiers.toString(); + } + + /** + * FOR VERSION AS OF / FOR TIME AS OF: a time-travel read must stay pinned to the + * captured version, otherwise the replay reads the current table contents. A numeric + * version is emitted as a literal, every other value (and all times) as a string. + */ + private static String renderTableSnapshot(Optional<TableSnapshot> tableSnapshot) { + if (!tableSnapshot.isPresent()) { + return ""; + } + TableSnapshot snapshot = tableSnapshot.get(); + if (snapshot.getType() == TableSnapshot.VersionType.TIME) { + return " FOR TIME AS OF " + quoteSqlString(snapshot.getValue()); + } + String value = snapshot.getValue(); + if (value.matches("[0-9]+")) { + return " FOR VERSION AS OF " + value; + } + return " FOR VERSION AS OF " + quoteSqlString(value); + } + + /** + * The @paramType(...) read parameters (incremental / branch / tag / options / + * snapshot / reset): the map form when the parameters carry key/value pairs, otherwise + * the bare identifier list form. Dropping them would replay a different (e.g. + * non-incremental) read than the captured one. + */ + private static String renderScanParams(Optional<TableScanParams> scanParams) { + if (!scanParams.isPresent()) { + return ""; + } + TableScanParams params = scanParams.get(); + StringBuilder modifiers = new StringBuilder(" @").append(params.getParamType()).append('('); + Map<String, String> mapParams = params.getMapParams(); + if (!mapParams.isEmpty()) { + boolean first = true; + for (Map.Entry<String, String> param : mapParams.entrySet()) { + if (!first) { + modifiers.append(", "); + } + first = false; + modifiers.append(quoteIdentifier(param.getKey())) + .append(" = ") + .append(quoteSqlString(param.getValue())); + } + } else { + List<String> listParams = params.getListParams(); + for (int i = 0; i < listParams.size(); i++) { + if (i > 0) { + modifiers.append(", "); + } + modifiers.append(quoteIdentifier(listParams.get(i))); + } + } + return modifiers.append(')').toString(); + } + + /** + * A single-quoted SQL string literal: embedded single quotes are doubled AND + * backslashes are doubled, so the value cannot terminate the literal, change the frozen + * SQL structure, or decode to a DIFFERENT value. The DEFAULT sql_mode treats a backslash + * as an escape introducer (a semantic value such as release\next would otherwise + * reparse as release + newline + ext under that mode and select another external ref), + * so the literal doubles it; the SPM re-parses pin the DEFAULT mode (see + * SPMPlanner.parseStoredSelect), which makes the round trip exact under default AND + * NO_BACKSLASH_ESCAPES sessions alike. + */ + static String quoteSqlString(String value) { + return "'" + value.replace("\\", "\\\\").replace("'", "''") + "'"; + } + + // ==================== Generate (LATERAL VIEW) ==================== + + /** + * PhysicalGenerate: LATERAL VIEW clauses over the child relation - one clause per + * generator. The parser wraps every user-written LATERAL VIEW ... into its + * own single-generator node, while MergeGenerates (for two independent + * views) folds stacked nodes into one node carrying several generators; the + * executor rolls several functions over each child row, i.e. a cartesian expansion, + * which is exactly what stacked LATERAL VIEWs express (MergeGenerates only merges + * when the upper view does not reference the lower view's output). LATERAL VIEW + * clauses attach to the child's COMPLETE query block: the child's WHERE / GROUP BY / + * HAVING / ORDER BY / LIMIT clauses must stay inside the lateral-view input (a + * derived table with LIMIT 10 limits the input, not the exploded rows). The + * generator output columns are registered on the same relation so parent operators + * can reference them; the alias is taken from the output slot's qualifier when it + * survived analysis, otherwise a per-decompile lv_N alias is used. Generate + * conjuncts (if any) are appended to the relation's WHERE. + */ + @Override + public SQLRelation visitPhysicalGenerate(PhysicalGenerate<? extends Plan> generate, Void context) { + SQLRelation childRelation = process(generate.child(0)); + List<Function> generators = generate.getGenerators(); + List<Slot> outputs = generate.getGeneratorOutput(); + if (generators.isEmpty() || generators.size() != outputs.size()) { + // post-binding plans keep exactly one output slot per generator (multi-column + // generators carry their columns in the expand alias project above) + throw new UnsupportedOperationException("SPM decompile generate: generator/output arity mismatch " + + generators.size() + "/" + outputs.size()); + } + // The LATERAL VIEW must attach to the child's COMPLETE query block. Attaching it + // to the bare FROM fragment (getFrom()) would keep the child's WHERE / GROUP BY / + // HAVING / ORDER BY / LIMIT clauses on the OUTER relation, where they apply AFTER + // the explode: a derived table with LIMIT 10 would limit the exploded rows + // instead of the lateral-view input, a GROUP BY would regroup the generator + // output, and frozen replay could return rows the captured plan filtered out. + // A FROM-less child (e.g. LATERAL VIEW over "SELECT 1 AS x") needs the wrapper + // for the same reason. + SQLRelation relation; + String baseSql; + if (childRelation.getFrom().isEmpty() || childRelation.hasOwnBlock()) { + if (childRelation.getRelationName() == null) { + childRelation.newAlias(); + } + baseSql = childRelation.toRelationSQL(); + // the wrapper carries the child's column mapping only; the child's clauses + // stay inside baseSql (moving them onto this relation would re-apply them + // after the lateral view) + relation = new SQLRelation(); + relation.getColumnNames().putAll(childRelation.getColumnNames()); + } else { + relation = childRelation; + baseSql = relation.getFrom(); + } + List<String> aliases = new ArrayList<>(outputs.size()); + // The analyzer may name a generator output with an internal column name + // ("$c$N"); such a name cannot be referenced in SQL, so give it a generated + // visible name and register the slot under that name for the parent operators. + // A generator output may also COLLIDE with a name the child already exports + // (SELECT t.x, lv.x FROM t LATERAL VIEW explode(t.arr) lv AS x is valid: the two + // slots carry different qualifiers, but the frozen SQL references a derived + // relation's columns by name) - registering both as bare x made the parent emit + // SELECT x, x ..., which fails binding as ambiguous after reload. Rename the + // generator output to a unique visible name and register THAT for the slot. + Set<String> takenNames = new HashSet<>(); + for (String existing : childRelation.getColumnNames().values()) { + if (existing != null) { + takenNames.add(existing.replace("`", "")); + } + } + List<String> columnNames = new ArrayList<>(outputs.size()); + for (Slot slot : outputs) { + String alias = slot.getQualifier().isEmpty() + ? "" : slot.getQualifier().get(slot.getQualifier().size() - 1); + aliases.add(alias.isEmpty() ? "lv_" + (lateralViewSeq++) : alias); + String name = slot.getName(); + String visible = name == null || name.startsWith("$c$") + ? "lv_col_" + (lateralViewSeq++) : name; + while (takenNames.contains(visible)) { + visible = visible + "_"; + } + takenNames.add(visible); + columnNames.add(visible); + } + StringBuilder from = new StringBuilder(baseSql); + for (int i = 0; i < generators.size(); i++) { + from.append(" LATERAL VIEW ") + .append(exprSqlBuilder.print(generators.get(i), relation)) + .append(' ') + .append(quoteIdentifier(aliases.get(i))) + .append(" AS ") + .append(quoteIdentifier(columnNames.get(i))); + } + relation.setFrom(from.toString()); + for (int i = 0; i < outputs.size(); i++) { + relation.registerRef(outputs.get(i).getExprId(), quoteIdentifier(columnNames.get(i))); + } + // When this relation is wrapped as a subquery by its parent, its SELECT list is + // the subquery output: the generator columns must be part of it, otherwise the + // parent's reference to a generated column fails to resolve ("Unknown column + // ... in table list"). An empty SELECT list means "*" and already covers them. + if (!relation.getSelects().isEmpty()) { + List<Pair<ExprId, String>> selects = new ArrayList<>(relation.getSelects()); + for (int i = 0; i < outputs.size(); i++) { + selects.add(Pair.of(outputs.get(i).getExprId(), quoteIdentifier(columnNames.get(i)))); + } + relation.setSelects(selects); + } + if (!generate.getConjuncts().isEmpty()) { + String conjuncts = generate.getConjuncts().stream() + .map(expr -> exprSqlBuilder.print(expr, relation)) + .collect(Collectors.joining(" AND ")); + relation.setWhere(relation.getWhere().isEmpty() Review Comment: [P1] Preserve LEFT JOIN UNNEST ON semantics in frozen SQL. These Generate conjuncts are the original ON predicates; PhysicalPlanTranslator passes them to TableFunctionNode.expandConjuncts, and BE emits one NULL-extended left row when every generated value fails them. Moving the predicate into WHERE after LATERAL VIEW filters that row instead. For `arr=[-1]`, `LEFT JOIN UNNEST(arr) AS u(x) ON u.x > 0` keeps the left row in the original plan, while the frozen `explode_outer` plus `WHERE u.x > 0` returns none. Render equivalent outer-join SQL or decline freezing this shape and replay its parameterized tree. -- 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]
