github-actions[bot] commented on code in PR #68499:
URL: https://github.com/apache/doris/pull/68499#discussion_r4095333446
##########
fe/fe-core/src/main/java/org/apache/doris/catalog/InternalSchemaInitializer.java:
##########
@@ -405,6 +405,7 @@ public static void createTbl() throws UserException {
* )
*/
createTable(getAuditLogCreateSql());
+ createTable(getSpmBaselinesCreateSql());
Review Comment:
[P1] Ensure the new SPM table on upgraded clusters
On an upgraded cluster the legacy statistics/partition/audit tables already
satisfy `created()`, so `run()` exits before this new `createTbl()` statement
is called. `created()` never checks `spm_baselines`, and there is no other
production `createTbl()` caller. `Env` then repeatedly fails to load the
missing table; `BaselineManager` stays `loaded=false` and global
CREATE/ALTER/DROP throw that the store is not ready. Please include
`spm_baselines` in the completion predicate (or ensure it independently) and
add an upgrade test where only this table is absent.
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/spm/capture/PlanCaptureManager.java:
##########
@@ -0,0 +1,404 @@
+// 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.capture;
+
+import org.apache.doris.catalog.Env;
+import org.apache.doris.common.Config;
+import org.apache.doris.common.util.MasterDaemon;
+import org.apache.doris.nereids.spm.BaselinePlan;
+import org.apache.doris.nereids.spm.BaselineSource;
+import org.apache.doris.nereids.spm.SPMPlanner;
+import org.apache.doris.nereids.spm.manager.BaselineManager;
+import org.apache.doris.qe.AutoCloseConnectContext;
+import org.apache.doris.qe.SessionVariable;
+import org.apache.doris.qe.VariableMgr;
+import org.apache.doris.statistics.util.StatisticsUtil;
+
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import com.google.common.annotations.VisibleForTesting;
+
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.atomic.AtomicLong;
+
+/**
+ * PlanCaptureManager - SPM auto capture scheduler (Phase 2, design doc 7.2.1
/ 7.2.4).
+ *
+ * A Leader-FE daemon that periodically scans the audit_log internal table and
+ * automatically creates baselines for high-value queries:
+ *
+ * - only queries executed by the Nereids planner are captured;
+ * - the capture filter (PlanCaptureFilter) enforces the multi-table /
table-exists /
+ * regex / performance-threshold rules;
+ * - the baseline is built through the Phase 1 flow
(SPMPlanner.buildBaselineFromSql:
+ * SPM-mode optimize + decompile + parameterize) with source = CAPTURE and
the actual
+ * query_time filled for candidate ordering;
+ * - duplicate (digest, planSql) baselines are skipped (BaselineManager dedup).
+ *
+ * The whole cycle is guarded by the global session variable
enable_plan_capture
+ * (default false, tunable via `SET GLOBAL enable_plan_capture = true`), and
any
+ * failure is logged and skipped so auto capture never breaks the cluster.
+ */
+public class PlanCaptureManager extends MasterDaemon {
+
+ private static final Logger LOG =
LogManager.getLogger(PlanCaptureManager.class);
+
+ private static final PlanCaptureManager INSTANCE = new
PlanCaptureManager();
+
+ private AuditLogScanner scanner = new AuditLogScanner();
+
+ /** Capture filter, refreshed from the global session variables each
cycle. */
+ private PlanCaptureFilter filter;
+
+ /** Last scan window start (epoch millis); 0 means "first run, scan one
interval". */
+ private long lastScanTimestamp = 0;
+
+ /**
+ * Re-scan overlap (millis) applied to the watermark: AuditLoader buffers
events
+ * asynchronously and writes their original event timestamp, so a row can
become
+ * visible AFTER its window has passed (it would otherwise be excluded
from every
+ * future window forever). Re-scanning a lagged/overlapping window plus
query-id
+ * deduplication makes late arrivals capturable without processing an
execution
+ * twice.
+ */
+ private static final long SCAN_WINDOW_OVERLAP_MS = 300_000L;
+
+ /** Upper bound for the processed-query-id dedup map. */
+ private static final int MAX_TRACKED_QUERY_IDS = 10000;
+
+ /** Query ids already handled in earlier (overlapping) windows. */
+ private final Map<String, Boolean> processedQueryIds = new
LinkedHashMap<>();
+
+ /**
+ * Resume cursor of a TRUNCATED scan window: (query_time, time, query_id)
of the last
+ * consumed row. Empty while no partial window is pending - a short batch
advances
+ * the watermark instead.
+ */
+ private long cursorQueryTime = 0;
+ private String cursorTime = "";
+ private String cursorQueryId = "";
+
+ /** Whether the cloud-mode warning was already logged (the gate fires
every cycle). */
+ private boolean cloudModeWarned = false;
+
+ // capture statistics (design doc 7.2.1 / 7.2.6)
+ private final AtomicLong successCount = new AtomicLong(0);
+ private final AtomicLong skipDuplicateCount = new AtomicLong(0);
+ private final AtomicLong skipSingleTableCount = new AtomicLong(0);
+ private final AtomicLong skipFilterCount = new AtomicLong(0);
+ private final AtomicLong failCount = new AtomicLong(0);
+
+ private PlanCaptureManager() {
+ super("PlanCaptureManager",
+
VariableMgr.getDefaultSessionVariable().getPlanCaptureIntervalSeconds() *
1000L);
+ this.filter = buildFilterFromGlobal();
+ }
+
+ public static PlanCaptureManager getInstance() {
+ return INSTANCE;
+ }
+
+ /**
+ * Builds a capture filter from the global session variables (so `SET
GLOBAL`
+ * changes to the thresholds / table regex take effect on the next cycle).
+ *
+ * @return a new filter
+ */
+ private static PlanCaptureFilter buildFilterFromGlobal() {
+ try {
+ SessionVariable global = VariableMgr.getDefaultSessionVariable();
+ return new PlanCaptureFilter(global.getPlanCaptureIncludePattern(),
+ global.getPlanCaptureExcludePattern(),
+ global.getPlanCaptureMinQueryTimeMs(),
+ global.getPlanCaptureMinScanRows());
+ } catch (RuntimeException e) {
+ // e.g. a legacy invalid regex in the global variable: never let
it escape the
+ // singleton constructor / the daemon cycle (leader startup calls
getInstance()
+ // before enable_plan_capture is even checked, and a
PatternSyntaxException
+ // there would terminate the FE transition)
+ LOG.error("SPM plan capture disabled: invalid capture filter
configuration", e);
+ return null;
+ }
+ }
+
+ @Override
+ protected void runAfterCatalogReady() {
+ SessionVariable global = VariableMgr.getDefaultSessionVariable();
+ // Reschedule from the cycle itself: MasterDaemon sleeps its stored
intervalMs, so
+ // only setInterval() here makes a `SET GLOBAL
plan_capture_interval_seconds`
+ // change affect future wakeups (rereading the variable in the cycle
would only
+ // change the scan window). Clamp to >= 1s so a misconfiguration
cannot spin.
+ setInterval(Math.max(1L, global.getPlanCaptureIntervalSeconds()) *
1000L);
+ // SPM baseline management (CREATE / ALTER / DROP / SHOW) explicitly
rejects cloud
+ // mode; until the full lifecycle is supported the capture daemon must
not create
+ // (or keep retrying to create) global baselines a cloud deployment
cannot show,
+ // disable or drop.
+ if (Config.isCloudMode()) {
+ if (!cloudModeWarned) {
+ cloudModeWarned = true;
+ LOG.warn("SPM plan capture is not supported in cloud mode,
skipping");
+ }
+ return;
+ }
+ if (!global.isEnablePlanCapture()) {
+ return;
+ }
+ if (!Env.getCurrentEnv().isMaster()) {
+ // auto capture runs on the Leader FE only
+ return;
+ }
+ if (Env.isCheckpointThread()) {
+ return;
+ }
+ PlanCaptureFilter newFilter = buildFilterFromGlobal();
+ if (newFilter == null) {
+ LOG.error("Plan capture filter unavailable (invalid capture
regex?),"
+ + " skipping this capture cycle");
+ return;
+ }
+ try {
+ // refresh the filter so SET GLOBAL changes take effect this cycle
+ this.filter = newFilter;
+
+ long currentTime = System.currentTimeMillis();
+ // overlap the window so audit rows loaded late (whose event time
is older
+ // than the last watermark) are still scanned; duplicates are
filtered by
+ // query id below
+ long scanStart = (lastScanTimestamp == 0)
+ ? currentTime - (long)
global.getPlanCaptureIntervalSeconds() * 1000L
+ : Math.max(0L, lastScanTimestamp - SCAN_WINDOW_OVERLAP_MS);
+ if (scanStart >= currentTime) {
+ return;
+ }
+
+ AuditLogScanner.ScanBatch batch = scanner.scan(scanStart,
currentTime,
+ global.getPlanCaptureMaxBatchSize(), cursorQueryTime,
cursorTime, cursorQueryId);
+ for (CapturedQuery candidate : batch.getCandidates()) {
+ String queryId = candidate.getQueryId();
+ if (queryId != null && !queryId.isEmpty() &&
!"NaN".equals(queryId)) {
+ if (processedQueryIds.containsKey(queryId)) {
+ continue; // already handled in an earlier overlapping
window
+ }
+ processedQueryIds.put(queryId, Boolean.TRUE);
Review Comment:
[P1] Do not permanently consume failed capture rows
This id is inserted into `processedQueryIds` before `processCandidate()`
runs, but `processCandidate()` catches every planning/catalog/persistence
exception and never removes or retries it. A transient failure therefore makes
the audit row invisible to every overlapping scan; it is recoverable only after
the 10,000-entry eviction, by which time the watermark has passed it. Please
mark ids only after successful/terminal handling or keep failed ids retryable
with bounded backoff, and add a failure-then-retry test.
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/spm/capture/AuditLogScanner.java:
##########
@@ -0,0 +1,294 @@
+// 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.capture;
+
+import org.apache.doris.qe.SessionVariable;
+import org.apache.doris.qe.VariableMgr;
+import org.apache.doris.statistics.repository.ResultRow;
+import org.apache.doris.statistics.util.StatisticsUtil;
+
+import java.time.Instant;
+import java.time.LocalDateTime;
+import java.time.ZoneId;
+import java.time.format.DateTimeFormatter;
+import java.util.ArrayList;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * AuditLogScanner - audit log query wrapper (Phase 2, design doc 7.2.3).
+ *
+ * Reads the __internal_schema.audit_log internal table through the internal
query
+ * mechanism and returns the high-value query candidates for SPM auto capture.
+ *
+ * Within a capture cycle the results are deduplicated by (catalog, db,
sql_digest): the
+ * record with the largest query_time wins, so the same query SHAPE is only
processed once
+ * per cycle - but only within one namespace. Identical unqualified SQL
executed in two
+ * databases is a DIFFERENT query for SPM (its eventual match key is
namespace-qualified),
+ * so the database / catalog must take part in the dedup key.
+ *
+ * Pagination: the batch LIMIT is applied with a stable (query_time, time,
query_id)
+ * cursor. The caller resumes from the returned cursor until a batch comes
back shorter
+ * than the limit (window exhausted); advancing the window past a truncated
batch would
+ * permanently skip every eligible row beyond the LIMIT.
+ */
+public class AuditLogScanner {
+
+ private static final DateTimeFormatter DATETIME_FORMAT =
+ DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
+
+ /** audit_log SELECT columns (order must match rowToCapturedQuery /
toBatch). */
+ private static final String SELECT_COLUMNS =
+ "`stmt`, `query_time`, `scan_rows`, `return_rows`, `sql_digest`,
`sql_hash`, `db`, `catalog`,"
+ + " `query_id`, `is_internal`, `time`";
+
+ /**
+ * Result of one audit scan: the namespace-deduplicated candidates plus
the resume
+ * cursor ((query_time, time, query_id) of the last RAW row read).
+ */
+ public static class ScanBatch {
+ private final List<CapturedQuery> candidates;
+ private final boolean windowExhausted;
+ private final long cursorQueryTime;
+ private final String cursorTime;
+ private final String cursorQueryId;
+
+ ScanBatch(List<CapturedQuery> candidates, boolean windowExhausted,
+ long cursorQueryTime, String cursorTime, String cursorQueryId)
{
+ this.candidates = candidates;
+ this.windowExhausted = windowExhausted;
+ this.cursorQueryTime = cursorQueryTime;
+ this.cursorTime = cursorTime == null ? "" : cursorTime;
+ this.cursorQueryId = cursorQueryId == null ? "" : cursorQueryId;
+ }
+
+ public List<CapturedQuery> getCandidates() {
+ return candidates;
+ }
+
+ /** Whether the batch returned fewer RAW rows than the limit (whole
window read). */
+ public boolean isWindowExhausted() {
+ return windowExhausted;
+ }
+
+ public long getCursorQueryTime() {
+ return cursorQueryTime;
+ }
+
+ public String getCursorTime() {
+ return cursorTime;
+ }
+
+ public String getCursorQueryId() {
+ return cursorQueryId;
+ }
+ }
+
+ /**
+ * Scans the audit_log table within the given time window (first page).
+ *
+ * @param startTimeMs window start (epoch millis, inclusive)
+ * @param endTimeMs window end (epoch millis, exclusive)
+ * @param maxBatchSize max number of raw rows to scan (prevents OOM)
+ * @return the scan batch (candidates + resume cursor)
+ */
+ public ScanBatch scan(long startTimeMs, long endTimeMs, int maxBatchSize) {
+ return scan(startTimeMs, endTimeMs, maxBatchSize, 0L, "", "");
+ }
+
+ /**
+ * Scans the audit_log table within the given time window, resuming after
the cursor
+ * returned by the previous batch.
+ *
+ * @param startTimeMs window start (epoch millis, inclusive)
+ * @param endTimeMs window end (epoch millis, exclusive)
+ * @param maxBatchSize max number of raw rows per batch
+ * @param cursorQueryTime query_time of the last consumed row (0 = start
from the top)
+ * @param cursorTime time (event time) of the last consumed row
+ * @param cursorQueryId query_id of the last consumed row
+ * @return the scan batch (candidates + resume cursor)
+ */
+ public ScanBatch scan(long startTimeMs, long endTimeMs, int maxBatchSize,
+ long cursorQueryTime, String cursorTime, String cursorQueryId) {
+ String start = formatTimestamp(startTimeMs);
+ String end = formatTimestamp(endTimeMs);
+
+ SessionVariable global = VariableMgr.getDefaultSessionVariable();
+ long minQueryTimeMs = global.getPlanCaptureMinQueryTimeMs();
+ long minScanRows = global.getPlanCaptureMinScanRows();
+ String sql = buildScanSql(start, end, maxBatchSize, minQueryTimeMs,
minScanRows,
+ cursorPredicate(cursorQueryTime, cursorTime, cursorQueryId));
+
+ List<ResultRow> rows = StatisticsUtil.execStatisticQuery(sql);
+ return toBatch(rows, maxBatchSize);
+ }
+
+ /**
+ * Turns one page of raw audit rows into a batch: namespace-aware dedup
plus the
+ * resume cursor. Package-visible for tests (the SQL / pagination contract
is tested
+ * against fabricated rows).
+ *
+ * @param rows the raw rows of one page
+ * @param maxBatchSize the batch limit (a shorter page exhausts the window)
+ * @return the scan batch
+ */
+ static ScanBatch toBatch(List<ResultRow> rows, int maxBatchSize) {
+ if (rows == null || rows.isEmpty()) {
+ return new ScanBatch(List.of(), true, 0L, "", "");
+ }
+ Map<String, CapturedQuery> deduped = new LinkedHashMap<>();
+ long lastQueryTime = 0L;
+ String lastTime = "";
+ String lastQueryId = "";
+ for (ResultRow row : rows) {
+ // the cursor always moves to the last RAW row read, even when
that row is
+ // unusable / filtered later: it has been consumed and must not be
scanned
+ // again by the next page
+ lastQueryTime = parseLong(row.getWithDefault(1, "0"));
+ lastTime = row.getWithDefault(10, "");
+ lastQueryId = row.getWithDefault(8, "");
+ CapturedQuery candidate = rowToCapturedQuery(row);
+ if (candidate == null || candidate.getStmt() == null ||
candidate.getStmt().isEmpty()) {
+ continue;
+ }
+ String digest = candidate.getSqlDigest();
+ if (digest == null || digest.isEmpty()) {
+ digest = candidate.getStmt();
+ }
+ // namespace-aware key: the database / catalog take part,
otherwise identical
+ // unqualified SQL from two namespaces collapses to one candidate
and the
+ // other namespace never gets a baseline (SPM namespace-qualifies
its match
+ // key, so the two executions really are different queries)
+ String key = candidate.getCatalog() + '\u0001' + candidate.getDb()
+ '\u0001' + digest;
+ deduped.merge(key, candidate, (a, b) -> b.getQueryTimeMs() >=
a.getQueryTimeMs() ? b : a);
+ }
+ return new ScanBatch(new ArrayList<>(deduped.values()), rows.size() <
maxBatchSize,
+ lastQueryTime, lastTime, lastQueryId);
+ }
+
+ /**
+ * Builds the audit_log scan SQL. Public for tests: the pushed-down
predicate shape is
+ * part of the capture contract - eligibility is query time OR scanned
rows (the same
+ * rule as PlanCaptureFilter), and internal maintenance queries are
filtered in SQL
+ * instead of relying on a hardcoded event flag.
+ *
+ * @param start window start timestamp (formatted)
+ * @param end window end timestamp (formatted)
+ * @param maxBatchSize LIMIT for the scan
+ * @param minQueryTimeMs query-time threshold
+ * @param minScanRows scan-rows threshold
+ * @return the scan SQL
+ */
+ public static String buildScanSql(String start, String end, int
maxBatchSize,
+ long minQueryTimeMs, long minScanRows) {
+ return buildScanSql(start, end, maxBatchSize, minQueryTimeMs,
minScanRows, "");
+ }
+
+ /**
+ * Builds the audit_log scan SQL with an optional resume-cursor predicate.
The ORDER
+ * BY defines the stable total order the cursor walks:
+ * (query_time DESC, time DESC, query_id DESC).
+ *
+ * @param start window start timestamp (formatted)
+ * @param end window end timestamp (formatted)
+ * @param maxBatchSize LIMIT for the scan
+ * @param minQueryTimeMs query-time threshold
+ * @param minScanRows scan-rows threshold
+ * @param cursorPredicate resume-cursor predicate (empty when starting at
the top)
+ * @return the scan SQL
+ */
+ public static String buildScanSql(String start, String end, int
maxBatchSize,
+ long minQueryTimeMs, long minScanRows, String cursorPredicate) {
+ return "SELECT " + SELECT_COLUMNS + " FROM __internal_schema.audit_log
"
+ + "WHERE `time` >= '" + start + "' AND `time` < '" + end + "' "
+ + "AND `is_query` = true "
+ + "AND `is_nereids` = true "
+ + "AND (`query_time` >= " + minQueryTimeMs
+ + " OR `scan_rows` >= " + minScanRows + ") "
+ + "AND `is_internal` = false "
+ + (cursorPredicate == null ? "" : cursorPredicate)
+ + " ORDER BY `query_time` DESC, `time` DESC, `query_id` DESC "
+ + "LIMIT " + maxBatchSize;
+ }
+
+ /**
+ * Resume-cursor predicate of the (query_time, time, query_id) total
order: strictly
+ * "after" the last consumed row, so a truncated batch continues exactly
where it
+ * stopped without re-reading or skipping rows.
+ */
+ static String cursorPredicate(long cursorQueryTime, String cursorTime,
String cursorQueryId) {
+ if (cursorQueryTime <= 0 || cursorTime == null || cursorTime.isEmpty()
Review Comment:
[P1] Preserve the resume cursor for zero/NULL query_time rows
`query_time` is nullable and scan eligibility also allows `scan_rows` alone,
so a full page can legitimately end with a row parsed as `query_time=0`. This
guard then clears the resume predicate and the next cycle starts at the first
page; query-id dedup skips repeats but rows after the page are never reached.
Track cursor presence separately (or use a sentinel outside the valid domain)
and add a zero/NULL-query-time pagination test.
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/spm/builder/SPMPlan2SQLBuilder.java:
##########
@@ -0,0 +1,2794 @@
+// 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.catalog.Column;
+import org.apache.doris.common.Pair;
+import org.apache.doris.nereids.properties.DistributionSpec;
+import org.apache.doris.nereids.properties.DistributionSpecReplicated;
+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.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.plans.AggPhase;
+import org.apache.doris.nereids.trees.plans.JoinType;
+import org.apache.doris.nereids.trees.plans.Plan;
+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.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.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.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.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 org.apache.doris.nereids.trees.plans.logical.LogicalFileScan;
+import org.apache.doris.nereids.trees.plans.physical.PhysicalFileScan;
+
+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.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.Set;
+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);
+
+ /** JOIN distribution HINT prefix constants. */
+ private static final String HINT_JOIN_BROADCAST = "BROADCAST";
+ private static final String HINT_JOIN_SHUFFLE = "SHUFFLE";
+
+ /** 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 {@code c_<seq>} 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;
+
+ /**
+ * 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) {
+ // 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;
+ 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)
renders the
+ // ASSERT_ROWS wrapper; nested ones are already handled by
toRelationSQL().
+ if (relation.isAssertRows()) {
+ return "ASSERT_ROWS (" + relation.toSQL() + ") " +
relation.getRelationAlias();
+ }
+ 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 ====================
+
+ /**
+ * 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);
+ }
+
+ /**
+ * 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).
+ 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();
+ 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
+ * {@code [A-Za-z_][A-Za-z0-9_]*} identifier: a column named {@code a-b}
must be
+ * emitted as {@code `a-b`}, otherwise the frozen projection re-parses as
the
+ * subtraction a - b. Embedded backticks are doubled. Plain 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
+ * {@code round((sun_sales1 / sun_sales2), 2)} is wrapped so the frozen
planSql can
+ * carry the original column header verbatim).
+ */
+ static String quoteIdentifier(String name) {
+ if (name == null) {
+ return null;
+ }
+ if (name.matches("[A-Za-z_][A-Za-z0-9_]*")) {
+ 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();
+ }
+
+ // ==================== 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;
+ rejectUnsupportedScan(relation);
+ 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. Each component is backtick-quoted when it is
not a plain
+ // identifier, so a metadata name containing operators is re-parsed as
an
+ // identifier instead of an expression.
+ sqlRelation.setFrom(catalogRelation.getTable().getDatabase() == null
+ ? quoteIdentifier(catalogRelation.getTable().getName())
+ :
quoteQualifiedName(catalogRelation.getTable().getNameWithFullQualifiers()));
+ // 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(tvfRelation.getFunction().toSql());
+ for (Slot slot : tvfRelation.getOutput()) {
+ relation.registerRef(slot.getExprId(),
quoteIdentifier(slot.getName()));
+ }
+ return relation;
+ }
+
+ /** PhysicalLazyMaterializeTVFScan: a TVF scan wrapped by lazy
materialization. */
+ @Override
+ public SQLRelation visitPhysicalLazyMaterializeTVFScan(
+ PhysicalLazyMaterializeTVFScan scan, Void context) {
+ return visitPhysicalTVFRelation(scan, context);
+ }
+
+ /**
+ * A scan carrying execution modifiers the decompiler cannot express in
plain SQL must
+ * never be frozen as an unrestricted catalog.db.table scan: replay would
silently run
+ * over all eligible partitions (e.g. FROM t PARTITION(p1) freezes and
later executes
+ * over every partition), use the wrong index, or drop the sampling. Fail
the decompile
+ * so CREATE falls back to the user-supplied planSql text instead.
+ *
+ * Note: selectedTabletIds is deliberately NOT a rejection criterion -
bucket pruning
+ * is derived from the query's own predicates (rule PruneOlapScanTablet),
so the
+ * replayed SQL re-derives the same selection; only sample / partition /
index
+ * selections are not reconstructible from the frozen text.
+ */
+ private static void rejectUnsupportedScan(PhysicalRelation relation) {
+ if (relation instanceof
org.apache.doris.nereids.trees.plans.physical.PhysicalOlapScan) {
+ org.apache.doris.nereids.trees.plans.physical.PhysicalOlapScan
scan =
+
(org.apache.doris.nereids.trees.plans.physical.PhysicalOlapScan) relation;
+ boolean partitionSubset = !scan.getSelectedPartitionIds().isEmpty()
+ && scan.getSelectedPartitionIds().size()
+ != scan.getTable().getPartitions().size();
+ if (scan.getSelectedIndexId() != scan.getTable().getBaseIndexId()
|| partitionSubset
+ || scan.getTableSample().isPresent()) {
+ throw new UnsupportedOperationException(
+ "SPM decompile: restricted olap scan
(index/partition/sample selection)"
+ + " is not supported yet");
+ }
+ return;
+ }
+ if (relation instanceof PhysicalFileScan) {
+ // File scans (external catalogs) carry the same class of
modifiers - partition
+ // pruning state, TABLESAMPLE, FOR VERSION AS OF snapshot state
and scan
+ // parameters - while the generic serializer emits only
catalog.db.table. A
+ // placeholder-bearing "FOR VERSION AS OF 123 ... WHERE k = 1"
baseline would
+ // replay against the current unrestricted table and return
different rows,
+ // so every non-default modifier fails the decompile: CREATE keeps
the user
+ // planSql text and the rewrite degrades to the parameterized-tree
path.
+ PhysicalFileScan scan = (PhysicalFileScan) relation;
+ boolean partitionPruned = scan.getSelectedPartitions() != null
+ && scan.getSelectedPartitions() !=
LogicalFileScan.SelectedPartitions.NOT_PRUNED;
+ if (partitionPruned || scan.getTableSample().isPresent()
+ || scan.getTableSnapshot().isPresent() ||
scan.getScanParams().isPresent()) {
+ throw new UnsupportedOperationException(
+ "SPM decompile: restricted file scan"
+ + " (partition/sample/snapshot/scan params) is
not supported yet");
+ }
+ }
+ }
+
+ // ==================== Generate (LATERAL VIEW) ====================
+
+ /**
+ * PhysicalGenerate: one LATERAL VIEW clause over the child relation. The
child
+ * relation's FROM is extended with a LATERAL VIEW clause carrying the
generator,
+ * the table alias and the column list; the generator output columns are
registered
+ * on the same relation so parent operators can reference them. The alias
is taken
+ * from the user's alias when it survived analysis (slot qualifier),
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 relation = process(generate.child(0));
+ if (generate.getGenerators().size() != 1) {
+ throw new UnsupportedOperationException("SPM decompile generate:
expected one generator, got "
+ + generate.getGenerators().size());
+ }
+ // a FROM-less child (e.g. LATERAL VIEW over "SELECT 1 AS x") has no
FROM text;
+ // force the subquery wrapping so the LATERAL VIEW has a relation to
attach to
+ String baseSql;
+ if (relation.getFrom().isEmpty()) {
+ relation.newAlias();
+ baseSql = relation.toRelationSQL();
+ } else {
+ baseSql = relation.getFrom();
Review Comment:
[P1] Preserve wrapped child clauses when adding LATERAL VIEW
When the Generate child is a wrapped relation (for example `(SELECT ...
WHERE k=1) s`), `getFrom()` returns only its raw FROM fragment. Appending
`LATERAL VIEW` here discards the child's WHERE/SELECT/GROUP/LIMIT clauses, so
frozen replay can produce rows the captured plan filtered out. Preserve the
complete child query block (force a subquery wrapper before attaching the
lateral view) and add a derived/filter-child replay test.
##########
fe/fe-core/src/main/java/org/apache/doris/qe/SessionVariable.java:
##########
@@ -2108,6 +2112,198 @@ public boolean isEnableHboNonStrictMatchingMode() {
@VarAttrDef.VarAttr(name = DISABLE_NEREIDS_RULES, needForward = true)
private String disableNereidsRules = "";
+ // ==================== SPM (SQL Plan Management) related config
====================
+ // Phase 1 SPM query rewrite switch and timeout. needForward = false: SPM
is pure FE
+ // logic and does not need to be forwarded to the BE for execution.
+ public static final String ENABLE_SPM_REWRITE = "enable_spm_rewrite";
+ public static final String SPM_REWRITE_TIMEOUT_MS =
"spm_rewrite_timeout_ms";
+
+ @VarAttrDef.VarAttr(name = ENABLE_SPM_REWRITE, needForward = false,
description =
+ "Whether to enable SPM (SQL Plan Management) query rewrite,
disabled by default for safety"
+ )
+ private boolean enableSpmRewrite = false;
+
+ @VarAttrDef.VarAttr(name = SPM_REWRITE_TIMEOUT_MS, needForward = false,
description =
+ "SPM rewrite timeout in milliseconds, fallback to normal execution
on timeout"
+ )
+ private int spmRewriteTimeoutMs = 1000;
+
+ public static final String ENABLE_SPM_FALLBACK = "enable_spm_fallback";
+
+ @VarAttrDef.VarAttr(name = ENABLE_SPM_FALLBACK, needForward = false,
description =
+ "Whether SPM falls back to the original query when the rewritten
(frozen-plan "
+ + "replay) plan fails to plan. Disabled by default so a
rewrite failure "
+ + "surfaces as an error (useful during development /
regression debugging); "
+ + "enable it for production availability so SPM never
breaks a query."
+ )
+ private boolean enableSpmFallback = false;
+
+ public boolean isEnableSpmRewrite() {
+ return enableSpmRewrite;
+ }
+
+ public void setEnableSpmRewrite(boolean enableSpmRewrite) {
+ this.enableSpmRewrite = enableSpmRewrite;
+ }
+
+ public int getSpmRewriteTimeoutMs() {
+ return spmRewriteTimeoutMs;
+ }
+
+ public void setSpmRewriteTimeoutMs(int spmRewriteTimeoutMs) {
+ this.spmRewriteTimeoutMs = spmRewriteTimeoutMs;
+ }
+
+ public boolean isEnableSpmFallback() {
+ return enableSpmFallback;
+ }
+
+ public void setEnableSpmFallback(boolean enableSpmFallback) {
+ this.enableSpmFallback = enableSpmFallback;
+ }
+
+ // ==================== SPM plan capture (Phase 2, design doc 7.2.6)
====================
+ // All are pure FE logic (needForward = false). They are registered as
session
+ // variables so they can be tuned globally with `SET GLOBAL ...` (which
updates
+ // VariableMgr.defaultSessionVariable, read by the Leader FE
PlanCaptureManager
+ // daemon). The plan_capture_* prefix distinguishes them from the Phase 1
spm_* vars.
+ public static final String ENABLE_PLAN_CAPTURE = "enable_plan_capture";
+ public static final String PLAN_CAPTURE_INTERVAL_SECONDS =
"plan_capture_interval_seconds";
+ public static final String PLAN_CAPTURE_MAX_BATCH_SIZE =
"plan_capture_max_batch_size";
+ public static final String PLAN_CAPTURE_MIN_QUERY_TIME_MS =
"plan_capture_min_query_time_ms";
+ public static final String PLAN_CAPTURE_MIN_SCAN_ROWS =
"plan_capture_min_scan_rows";
+ public static final String PLAN_CAPTURE_INCLUDE_PATTERN =
"plan_capture_include_pattern";
+ public static final String PLAN_CAPTURE_EXCLUDE_PATTERN =
"plan_capture_exclude_pattern";
+
+ @VarAttrDef.VarAttr(name = ENABLE_PLAN_CAPTURE, needForward = false,
description =
+ "Whether to enable SPM (SQL Plan Management) auto plan capture.
The Leader FE "
+ + "periodically scans the audit_log internal table and
automatically creates "
+ + "baselines for high-value queries (multi-table, slow or
heavy scan).")
+ private boolean enablePlanCapture = false;
+
+ @VarAttrDef.VarAttr(name = PLAN_CAPTURE_INTERVAL_SECONDS, needForward =
false, description =
+ "The interval (in seconds) between two SPM auto-capture cycles.
Default is 10800 (3 hours).")
+ private int planCaptureIntervalSeconds = 10800;
+
+ @VarAttrDef.VarAttr(name = PLAN_CAPTURE_MAX_BATCH_SIZE, needForward =
false, description =
+ "The max number of audit records processed in a single SPM capture
cycle.")
+ private int planCaptureMaxBatchSize = 500;
+
+ @VarAttrDef.VarAttr(name = PLAN_CAPTURE_MIN_QUERY_TIME_MS, needForward =
false, description =
+ "Queries whose execution time is below this threshold (ms) are not
captured.")
+ private long planCaptureMinQueryTimeMs = 1000;
+
+ @VarAttrDef.VarAttr(name = PLAN_CAPTURE_MIN_SCAN_ROWS, needForward =
false, description =
+ "Queries whose scan rows are below this threshold are not
captured.")
+ private long planCaptureMinScanRows = 10000;
+
+ @VarAttrDef.VarAttr(name = PLAN_CAPTURE_INCLUDE_PATTERN, needForward =
false,
+ setter = "setPlanCaptureIncludePattern", description =
+ "Only queries whose table names match this regex are captured.
Empty means all.")
+ private String planCaptureIncludePattern = "";
+
+ @VarAttrDef.VarAttr(name = PLAN_CAPTURE_EXCLUDE_PATTERN, needForward =
false,
+ setter = "setPlanCaptureExcludePattern", description =
+ "Queries with any table matching this regex are skipped.")
+ private String planCaptureExcludePattern = "";
+
+ public static final String SPM_BASELINE_REFRESH_INTERVAL_SECONDS =
+ "spm_baseline_refresh_interval_seconds";
+
+ @VarAttrDef.VarAttr(name = SPM_BASELINE_REFRESH_INTERVAL_SECONDS,
needForward = false, description =
+ "The interval (in seconds) between two SPM baseline cache refresh
cycles. Every FE "
+ + "periodically merges the shared baseline internal table
(read-only) so "
+ + "baselines created on another FE - user DDL or Leader
auto capture - become "
+ + "visible locally. Default is 60.")
+ private int spmBaselineRefreshIntervalSeconds = 60;
+
+ public boolean isEnablePlanCapture() {
+ return enablePlanCapture;
Review Comment:
[P2] Reject non-positive capture interval and batch settings
These new global setters accept zero/negative values.
`plan_capture_interval_seconds=0` makes every cycle compute `scanStart >=
currentTime` and return without scanning, while `plan_capture_max_batch_size=0`
produces `LIMIT 0`, marks the window exhausted, and advances the watermark over
eligible rows. Please reject non-positive values (or clamp before
window/pagination logic) through the SQL SET path and add range-validation
tests.
--
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]