github-actions[bot] commented on code in PR #68499: URL: https://github.com/apache/doris/pull/68499#discussion_r4177858220
########## fe/fe-core/src/main/java/org/apache/doris/plugin/audit/AuditPublicationHorizon.java: ########## @@ -0,0 +1,358 @@ +// 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.plugin.audit; + +import org.apache.doris.catalog.Env; +import org.apache.doris.catalog.InternalSchema; +import org.apache.doris.common.FeConstants; +import org.apache.doris.qe.AuditEventProcessor; +import org.apache.doris.resource.workloadschedpolicy.WorkloadRuntimeStatusMgr; +import org.apache.doris.statistics.repository.ResultRow; +import org.apache.doris.statistics.util.StatisticsUtil; + +import com.google.common.annotations.VisibleForTesting; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.time.Instant; +import java.time.LocalDateTime; +import java.time.ZoneOffset; +import java.time.format.DateTimeFormatter; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.function.Function; +import java.util.function.Supplier; + +/** + * The cluster-wide AUDIT PUBLICATION HORIZON: the start time (epoch millis, the + * {@code time} column of {@code audit_log}) of the oldest audit event that any FE has + * accepted but not yet PUBLISHED. The SPM capture scans the shared audit table from the + * leader, so it uses this value as a progress FENCE: its next scan window must still + * start at or before it, otherwise a row an FE still owes falls behind the advanced + * watermark and is never captured. + * + * <p>Three layers make the fence complete: + * <ul> + * <li>{@link #localHorizon()} folds THIS FE's whole audit pipeline: completed queries + * still held by the {@link WorkloadRuntimeStatusMgr} (they enter the pipeline + * before any loader sees them), the {@link AuditEventProcessor} queue and its + * in-flight event (a plugin can stall while an event is dequeued), and the + * {@link AuditLoader} queue / assembled batch / not-yet-visible batch (a stream + * load can report Publish Timeout after commit). The stages are read + * UPSTREAM-FIRST (round-37 #3): every handoff enqueues the event downstream + * BEFORE the upstream stage stops covering it, so an event transferred between + * two reads can never fall in the gap - it is either still seen upstream or + * already seen downstream. Each stage keeps the event owned across its own + * handoff as well (the manager holds dequeued events until the processor call + * returns, the processor dequeues and publishes in-flight atomically - round-37 + * #1/#2).</li> + * <li>each FE REPORTS its local horizon into the shared + * {@link InternalSchema#SPM_AUDIT_HORIZON_TBL_NAME} table, so a follower's + * backlog is visible to the leader that runs the capture. INTERNAL statements + * (including the reporter's own SQL) are not part of either side: the capture + * never scans {@code is_internal = true} rows, so including them would only let + * the reporter's writes fence (and thereby re-trigger) themselves forever on an + * idle FE (round-37 #7).</li> + * <li>{@link #clusterHorizon()} is the MINIMUM over the local value and the FRESH + * rows of that table; a row the reporter stopped refreshing (its FE died or + * stopped reporting - the events are gone with it) is ignored.</li> + * </ul> + */ +public final class AuditPublicationHorizon { + + private static final Logger LOG = LogManager.getLogger(AuditPublicationHorizon.class); + + /** + * A reported row older than this is IGNORED: its FE stopped refreshing the fence + * (crashed / killed / its reporter thread is gone), so the events it still owed are + * lost with it and fencing progress forever would freeze the capture instead of + * protecting anything. Must be comfortably larger than the reporter's keepalive + * interval ({@link AuditLoader#HORIZON_KEEPALIVE_MILLIS}). + */ + public static final long ROW_STALE_MILLIS = 5 * 60 * 1000L; + + private static final String SELECT_ROWS_SQL = + "SELECT `horizon_ms`, `update_time` FROM `" + FeConstants.INTERNAL_DB_NAME + "`." + + "`" + InternalSchema.SPM_AUDIT_HORIZON_TBL_NAME + "`"; + private static final String SELECT_OWN_ROW_SQL = "SELECT `horizon_ms` FROM `" + + FeConstants.INTERNAL_DB_NAME + "`." + + "`" + InternalSchema.SPM_AUDIT_HORIZON_TBL_NAME + "` WHERE `fe_name` = '${feName}'"; + // ONE atomic statement per report: the table is a merge-on-write UNIQUE KEY(`fe_name`) + // table, so an INSERT of an existing fe_name IS the update of that FE's row - there is + // no window (a crash, or a reader between two statements) in which the row is MISSING + // while the follower still owes an old event (round-37 #9: the previous DELETE+INSERT + // committed separately and the leader could read no row in between). A zero horizon + // deletes the row instead (also one statement): a missing row and a zero row are the + // same "nothing outstanding" to every reader. + private static final String UPSERT_OWN_ROW_SQL = "INSERT INTO `" + FeConstants.INTERNAL_DB_NAME + "`." + + "`" + InternalSchema.SPM_AUDIT_HORIZON_TBL_NAME + "`" + + " (`fe_name`, `horizon_ms`, `update_time`) VALUES ('${feName}', ${horizonMs}, '${updateTime}')"; + private static final String DELETE_OWN_ROW_SQL = "DELETE FROM `" + FeConstants.INTERNAL_DB_NAME + "`." + + "`" + InternalSchema.SPM_AUDIT_HORIZON_TBL_NAME + "` WHERE `fe_name` = '${feName}'"; + private static final int IO_TIMEOUT_SECONDS = 10; + + /** + * update_time is rendered AND parsed in UTC (round-37 #4): the column is a zone-less + * DATETIME crossing FEs that may render their local wall time in different zones, so + * the previous both-sides-local rendering made a fresh row look hours old to a reader + * in another zone (discarded as stale, dropping that follower's fence). A fixed zone + * on both sides makes the freshness comparison independent of either FE's time zone. + */ + private static final DateTimeFormatter UPDATE_TIME_PATTERN = + DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"); + private static final DateTimeFormatter UPDATE_TIME_UTC_FORMATTER = + UPDATE_TIME_PATTERN.withZone(ZoneOffset.UTC); + + /** + * Test seam: the shared-table read (one row per FE). Null in production. + */ + @VisibleForTesting + static volatile Supplier<List<Object[]>> horizonRowsReaderForTest; + + /** + * Test seam: the shared-table write of this FE's row. Returns whether the write is + * CONFIRMED (the production implementation re-reads its own row). Null in production. + */ + @VisibleForTesting + static volatile Function<Long, Boolean> localHorizonWriterForTest; + + private AuditPublicationHorizon() { + } + + /** + * The oldest audit event THIS FE has accepted but not published, 0 when nothing is + * outstanding: the MINIMUM over every stage of the local pipeline (see the class + * javadoc). Cheap - no I/O - so callers may poll it. + * + * <p>The stages are read UPSTREAM-FIRST (round-37 #3): the pre-loader stages before + * the loader. A downstream stage enqueues an event BEFORE the upstream stage + * releases it, so reading upstream first means an event transferred between the two + * reads is either still seen upstream (it has not transferred yet) or already seen + * downstream (it transferred before the upstream read) - the previous + * downstream-first order could read both stages around the transfer and miss it. + */ + public static long localHorizon() { + long oldest = preLoaderHorizon(); + oldest = minPositive(oldest, AuditLoader.oldestUnpublishedEventTime()); + return oldest; + } + + /** + * The stages BEFORE the audit loader, read UPSTREAM-FIRST (round-37 #3): the runtime + * status manager (a completed query enters its list before the processor sees it, and + * its dequeued events stay fenced until the processor call returns) and then the + * processor (whose dequeued-in-flight event is published atomically with the + * dequeue). + */ + private static long preLoaderHorizon() { + long oldest = 0; + try { + Env env = Env.getCurrentEnv(); + WorkloadRuntimeStatusMgr mgr = env == null ? null : env.getWorkloadRuntimeStatusMgr(); + if (mgr != null) { + oldest = minPositive(oldest, mgr.oldestHeldAuditEventTime()); + } + } catch (Throwable t) { + LOG.debug("audit publication horizon: the workload runtime status manager is" + + " unavailable: {}", t.getMessage()); + } + try { + AuditEventProcessor processor = Env.getCurrentAuditEventProcessor(); + if (processor != null) { + oldest = minPositive(oldest, processor.oldestQueuedOrInFlightEventTime()); + } + } catch (Throwable t) { + // an FE without this component (tests, partial startup) has nothing to fence + LOG.debug("audit publication horizon: the audit event processor is unavailable: {}", + t.getMessage()); + } + return oldest; + } + + /** + * The fence the CAPTURE uses: the minimum over this FE's own pipeline and the fresh + * rows every other FE reported. Throws {@link IllegalStateException} when the shared + * table cannot be read - the caller must NOT advance without a complete fence + * (round-36 #1: an unreadable follower row is exactly the hole this guards). + */ + public static long clusterHorizon() { + long oldest = localHorizon(); + return minPositive(oldest, remoteHorizon()); + } + + /** + * The minimum horizon over the FRESH rows of the shared table (0 when none / all + * stale). A read failure propagates as a retryable {@link IllegalStateException}. + */ + private static long remoteHorizon() { + List<Object[]> rows; + Supplier<List<Object[]>> reader = horizonRowsReaderForTest; + if (reader != null) { + rows = reader.get(); + } else if (!sharedTableAvailable()) { + return 0L; // no live FE environment (unit tests / not ready): nothing reported + } else { + try { + List<ResultRow> result = StatisticsUtil.executeQuery( + SELECT_ROWS_SQL, Collections.emptyMap(), IO_TIMEOUT_SECONDS); + rows = new ArrayList<>(); + if (result != null) { + for (ResultRow row : result) { + List<String> values = row.getValues(); + if (values == null || values.size() < 2) { + continue; + } + rows.add(new Object[] {Long.parseLong(values.get(0).trim()), + parseUpdateTime(values.get(1).trim())}); + } + } + } catch (Exception e) { + throw new IllegalStateException("SPM capture cannot read the cluster audit" + + " publication horizon: " + e.getMessage(), e); + } + } + long oldest = 0; + long now = System.currentTimeMillis(); + for (Object[] row : rows) { + if (row == null || row.length < 2 || row[0] == null || row[1] == null) { + continue; + } + long horizon = (Long) row[0]; + long updatedAt = (Long) row[1]; + if (horizon <= 0) { + continue; + } + if (updatedAt <= 0 || now - updatedAt > ROW_STALE_MILLIS) { Review Comment: [P2] Keep a live follower's positive audit fence when its report expires. A follower can still hold a completed 10:00 event while keepalive upserts fail for five minutes; this branch then drops its previously confirmed positive row, and the leader checkpoints 12:00 with no fence. If the follower loads at 12:01, the next five-minute overlap starts at 11:55 and the event's 10:01 completion also misses the scanner predicate, so capture never sees it. Check FE liveness before discarding a positive row, or fail the capture cycle closed while a live reporter is overdue. ########## fe/fe-core/src/main/java/org/apache/doris/nereids/spm/manager/BaselineManager.java: ########## @@ -0,0 +1,3913 @@ +// 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.manager; + +import org.apache.doris.catalog.Env; +import org.apache.doris.catalog.InternalSchema; +import org.apache.doris.common.FeConstants; +import org.apache.doris.common.FeNameFormat; +import org.apache.doris.common.Pair; +import org.apache.doris.nereids.analyzer.UnboundRelation; +import org.apache.doris.nereids.parser.NereidsParser; +import org.apache.doris.nereids.spm.BaselinePlan; +import org.apache.doris.nereids.spm.BaselineScope; +import org.apache.doris.nereids.spm.BaselineSource; +import org.apache.doris.nereids.spm.BaselineStatus; +import org.apache.doris.nereids.spm.SPMPlanTreeSupport; +import org.apache.doris.nereids.spm.SPMPlanner; +import org.apache.doris.nereids.trees.plans.Plan; +import org.apache.doris.nereids.trees.plans.logical.LogicalPlan; +import org.apache.doris.qe.ConnectContext; +import org.apache.doris.qe.MasterOpExecutor; +import org.apache.doris.qe.QueryState; +import org.apache.doris.qe.SqlModeHelper; +import org.apache.doris.statistics.repository.ResultRow; +import org.apache.doris.statistics.util.StatisticsUtil; + +import com.google.common.annotations.VisibleForTesting; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.time.Instant; +import java.time.LocalDateTime; +import java.time.ZoneOffset; +import java.time.format.DateTimeFormatter; +import java.util.ArrayList; +import java.util.Collections; +import java.util.Comparator; +import java.util.HashMap; +import java.util.Iterator; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.locks.ReentrantReadWriteLock; +import java.util.function.Supplier; +import java.util.stream.Collectors; + +/** + * BaselineManager - baseline storage, cache and index management (M3). + * + * Corresponds to design doc section 6.6. Manages the CRUD of baselines and maintains + * two query structures: + * + * - hashIndex: {@code Map<Long, List<Long>>}, bindSqlHash -> baseline id list. Level 1 + * coarse filtering with O(1) lookup. + * - baselines: id -> BaselinePlan in-memory storage (Phase 1 MVP). Phase 2 persists it + * to the __internal_schema.spm_baselines internal table (see design doc 6.14). + * + * Candidate baseline lookup (the first two of the three-level filter): + * + * 1. Level 1: hashIndex.get(queryHash) -> candidate id list + * 2. Level 2: exact digest.equals(baseline.bindSqlDigest) matching + * 3. Sort by priority (see the comparator below) + * + * Id source (the GLOBAL id invariant): GLOBAL writes have a single writer (user DDL is + * forwarded to the master, auto capture runs on the Leader), so "read the persistence + * watermark, then allocate" is sufficient to keep the local generator from ever handing + * out an id that is already used in the shared table: before EVERY allocation + * createBaseline reads MAX(id) from the internal table (one light aggregation) and + * advances idGenerator past it; when that read fails the CREATE fails visibly with a + * retryable error and no id is allocated. This closes the failover lag hole (a new + * master may start with a generator behind the table) and the load-failure hole (a + * broken startup load leaves the generator at 1 while the table is full); it stays + * correct even when a row was inserted out of contract (e.g. a manual table write). The + * startup load / periodic refresh keep advancing the generator as a second safety net. + * Nothing that does not allocate (ALTER / DROP / matching) reads the watermark, and + * creates only come from the DDL / capture cycle, so the read is off the query hot path. + * GLOBAL ids start at 1 and stay in [1, 2^62); SESSION ids live in [2^62, 2^63) (see + * BaselineScope.ofId), so the scope of an id is always exact. + * + * Concurrency: the in-memory store (baselines / hashIndex / stateVersion / load state) is + * guarded by one read-write lock. Lookups take the read lock and never serialize on each + * other; create / drop / status / load / refresh take the write lock; internal-table I/O + * runs outside the lock wherever correctness allows (see the individual methods). + */ +public class BaselineManager { + + // ==================== test seams (never set in production) ==================== + + /** + * Test seam: routes the status-protocol durable I/O (INSERT / DELETE by status / the + * reconciliation count read) to a simulator instead of the internal table, so a unit + * test can inject faults such as "the old-row delete committed but reported + * KV_TXN_MAYBE_COMMITTED". Null in production. + */ + @VisibleForTesting + interface StatusProtocolStoreForTest { + void insert(BaselinePlan plan); + + void deleteByIdAndStatus(long id, BaselineStatus status); + + int countByIdAndStatus(long id, BaselineStatus status); + + /** + * The CONDITIONAL insert half of a status flip: mirrors the durable + * {@code INSERT ... SELECT ... WHERE id / status} statement - the new-status row + * must NOT be written once the previous-status row is gone (the caller then + * refuses the flip, which is how the DROP-while-ALTER-stalled conflict surfaces). + * The default keeps the unconditional simulators working: their scenarios always + * keep the previous row. + * + * @return whether the new-status row was written + */ + default boolean insertIfPreviousPresent(BaselinePlan plan, BaselineStatus previousStatus) { + insert(plan); + return true; + } + } + + /** + * Test seam for the create-time id allocator / collision protocol: routes the + * watermark read, the INSERT, the by-id collision probe, the identity delete AND its + * ambiguous-commit reconciliation read to a simulator, so a unit test can inject a + * COMPETING master's row between the INSERT and the probe (the latch-driven handoff + * scenario) or an unconfirmable delete. Null in production. + */ + @VisibleForTesting + interface IdAllocatorStoreForTest { + long watermark(); + + default long seqWatermark() { + return 0; + } + + default void reserveId(long id) { + } + + void insert(BaselinePlan plan); + + List<BaselinePlan> readById(long id); + + void deleteByIdentity(BaselinePlan plan); + } + + @VisibleForTesting + public static volatile StatusProtocolStoreForTest statusProtocolStoreForTest; + + @VisibleForTesting + public static volatile IdAllocatorStoreForTest idAllocatorStoreForTest; + + /** + * Test seam replacing the live leadership probe of {@link #assertLeaderForWrite} + * (null in production). The store simulators bypass the live fence by design, so + * without this seam a unit test cannot interleave a master handoff with an in-flight + * write (the insert / delete halves of a status flip). + */ + @VisibleForTesting + public static volatile java.util.function.BooleanSupplier leaderProbeForTest; + + /** + * Test seam for the read-back visibility confirmation of a reported-successful write + * (see {@link #confirmInsertVisible}): one call is ONE probe attempt, true = the row + * (insert) or its status row is READABLE, false = not yet visible. A test + * decrements an invisible window here to simulate the COMMITTED-but-not-yet-published + * state the real store exposes. Null in production. + */ + @VisibleForTesting + interface DurableVisibilityProbeForTest { + boolean isReadable(long id, BaselineStatus status); + + /** + * The confirmation of a row JUST WRITTEN additionally checks the ATTEMPTED + * STORED SECOND: the requested status ALONE is weak evidence (a previously + * failed old-row delete can leave a STALE row of that very status behind, see + * {@link #observedInsertRowIsOurs}). The default delegates to the two-argument + * form so a simulator that models only the visibility window of one row keeps + * its semantics. + * + * @param id the baseline id + * @param status the status the write attempted + * @param updateTime the attempted row's update time (stored seconds) + * @return whether THAT row is readable + */ + default boolean isReadable(long id, BaselineStatus status, long updateTime) { + return isReadable(id, status); + } + } + + @VisibleForTesting + public static volatile DurableVisibilityProbeForTest durableVisibilityProbeForTest; + + /** + * Test seam replacing the snapshot READ of the load path (loadFromInternalTable / + * the promotion reload): lets a unit test return a controlled snapshot and, together + * with {@link #snapshotReadStartedHookForTest}, invalidate the store WHILE a load is + * still inside its read - the stale snapshot must then be discarded instead of + * republished. Null in production. + */ + @VisibleForTesting + public static volatile Supplier<Map<Long, BaselinePlan>> snapshotReaderForTest; + + /** + * Test seam counting the background load threads that were actually STARTED by + * {@link #scheduleAsyncLoad} (one per load-slot claim). A query burst must coalesce + * onto the in-flight load instead of starting one thread per caller, which this + * counter makes observable. Null in production. + */ + @VisibleForTesting + public static volatile java.util.concurrent.atomic.AtomicInteger asyncLoadSpawnCountForTest; + + /** + * Test seam replacing the journal synchronization of + * {@link #refreshAfterForwardedDdl} and {@link #confirmGlobalRowsForShow} (null in + * production): the real sync asks the master for its max journal id and waits + * locally, which a unit test cannot do. A test whose snapshot reader returns a + * PRE-DDL snapshot until this seam ran proves the sync happens BEFORE the snapshot + * read. + */ + @VisibleForTesting + public static volatile Runnable forwardedDdlSyncForTest; + + /** + * Test seam invoked by a load right after it captured its generation and BEFORE the + * snapshot read: a test blocks here, invalidates the store (the promotion window) + * and lets the load continue - the now-stale snapshot must be discarded. + */ + @VisibleForTesting + public static volatile Runnable snapshotReadStartedHookForTest; + + private static final Logger LOG = LogManager.getLogger(BaselineManager.class); + + /** Singleton. */ + private static final BaselineManager INSTANCE = new BaselineManager(); + + // ==================== internal-table persistence (__internal_schema.spm_baselines) ========== + + /** Fully qualified internal table (FeConstants.INTERNAL_DB_NAME == "__internal_schema"). */ + private static final String SPM_BASELINES_TABLE = + FeConstants.INTERNAL_DB_NAME + "." + InternalSchema.SPM_BASELINES_TBL_NAME; + + /** The append-only id reservation table (see InternalSchema#SPM_BASELINES_SEQ_TBL_NAME). */ + private static final String SPM_BASELINES_SEQ_TABLE = + FeConstants.INTERNAL_DB_NAME + "." + InternalSchema.SPM_BASELINES_SEQ_TBL_NAME; + + /** Column order follows InternalSchema.SPM_BASELINES_SCHEMA (unpaged selection). */ + private static final String SNAPSHOT_COLUMNS = + "SELECT `id`, `bind_sql`, `bind_sql_digest`," + + " `bind_sql_hash`, `plan_sql`, `query_id`, `cost`, `query_time_ms`, `source`," + + " `status`, `create_time`, `update_time`, `sql_mode`, `plan_sql_mode`," + + " `plan_frozen`, `schema_fingerprint` FROM "; + + /** + * First page of a whole-table snapshot: ordered by id so the pagination can continue + * with {@link #SELECT_PAGE_SQL} from the last row read (and so the duplicate-id + * resolution sees a stable order). + */ + private static final String SELECT_ALL_ORDERED_SQL = + SNAPSHOT_COLUMNS + SPM_BASELINES_TABLE + " ORDER BY `id`"; + + /** + * One continuation page of a whole-table snapshot: every row with {@code id >= + * ${lastId}}, ordered by id and SKIPPING the first {@code ${offset}} rows of that + * range. The offset is what keeps an id group larger than one page readable: a + * repeated opposite-status ALTER failure leaves one more row under the id every time, + * so the group can outgrow {@link #SNAPSHOT_PAGE_SIZE} rows - a jump past it would + * omit the rows behind the first page, possibly the newest durable status. + */ + private static final String SELECT_PAGE_SQL = SNAPSHOT_COLUMNS + SPM_BASELINES_TABLE Review Comment: [P2] Give equal-id rows a stable order across snapshot pages. This `OFFSET` is applied to a separate SELECT ordered only by `id`. With 1,999 lower-id rows plus old and newer status rows for id 2000, page one can end on the old row; if the next SELECT reverses those tied rows, `OFFSET 1` reads the old row again and skips the newer one. The read count still equals the unchanged table's `COUNT(*)`, so refresh can publish the wrong status. Order by a unique row identity or page from a stable snapshot, and test tie reordering between page queries. ########## fe/fe-core/src/main/java/org/apache/doris/nereids/spm/SPMOptimizer.java: ########## @@ -0,0 +1,593 @@ +// 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; + +import org.apache.doris.common.AnalysisException; +import org.apache.doris.common.UserException; +import org.apache.doris.nereids.NereidsPlanner; +import org.apache.doris.nereids.StatementContext; +import org.apache.doris.nereids.cost.Cost; +import org.apache.doris.nereids.memo.GroupExpression; +import org.apache.doris.nereids.parser.NereidsParser; +import org.apache.doris.nereids.properties.PhysicalProperties; +import org.apache.doris.nereids.properties.SelectHint; +import org.apache.doris.nereids.properties.SelectHintSetVar; +import org.apache.doris.nereids.rules.RuleType; +import org.apache.doris.nereids.trees.plans.Plan; +import org.apache.doris.nereids.trees.plans.commands.Command; +import org.apache.doris.nereids.trees.plans.commands.ExplainCommand.ExplainLevel; +import org.apache.doris.nereids.trees.plans.logical.LogicalPlan; +import org.apache.doris.nereids.trees.plans.logical.LogicalSelectHint; +import org.apache.doris.nereids.trees.plans.physical.PhysicalPlan; +import org.apache.doris.qe.ConnectContext; +import org.apache.doris.qe.OriginStatement; +import org.apache.doris.qe.SessionVariable; + +import com.google.common.collect.ImmutableList; +import com.google.common.collect.ImmutableSet; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.BitSet; +import java.util.HashSet; +import java.util.LinkedHashSet; +import java.util.List; +import java.util.Locale; +import java.util.Set; +import java.util.stream.Collectors; +import java.util.stream.Stream; + +/** + * SPMOptimizer - baseline-dedicated optimizer (M3). + * + * Corresponds to design doc section 6.7. When a baseline is created, a dedicated + * optimizer is used that deliberately disables many "state-sensitive" optimization + * rules (MV rewrite, table pruning, UKFK JOIN pruning, equivalence derivation, + * structural rewrites, ...), so the baseline plan only depends on SQL semantics and the + * general cost model - because a baseline is a cross-time promise that must remain + * reproducible and semantically correct even after data, statistics or MVs change. + * + * Rule exclusion mechanism (WHITELIST mode): the session variable enable_nereids_rules + * carries a comma-separated rule WHITELIST; when non-empty, the engine only applies the + * listed rules (the statement-level rule mask forbids every rule outside the list, see + * StatementContext#getOrCacheDisableRules). SPMOptimizer temporarily replaces that + * variable with the SPM whitelist - every RuleType except the excluded set below, see + * buildSpmEnabledRules - while a baseline is created, and restores the original value + * after optimization completes (the original value is cached before calling and + * restored in a finally block). disable_nereids_rules is NOT touched: the user's own + * disable list keeps applying on top of the whitelist. + * + * The excluded set is: + * - an explicit list of state-sensitive RuleType names (categories 1, 3, 4, 5, 6), and + * - the whole RuleTypeClass.MATERIALIZE_VIEW rule family (category 2) enumerated + * programmatically via RuleType.isMaterializedViewRule(), so new MV rules are covered + * automatically. + */ +public class SPMOptimizer { + + /** + * RuleType names excluded from the SPM whitelist (aligned with design doc 5.4). + * All names must exist in RuleType (the whitelist is parsed with + * RuleType.valueOf, which throws on unknown names). + */ + public static final List<String> SPM_EXCLUDED_RULE_NAMES = List.of( + // ===== category 1: data-state sensitive ===== + // simple aggregate to constant (wrong result on an empty table) + "REWRITE_SIMPLE_AGG_TO_CONSTANT", + // skewed JOIN salt splitting (skew pattern changes over time) + "SALT_JOIN", + // grouping sets decomposition (topology change, cardinality dependent) + "DECOMPOSE_REPEAT", + + // ===== category 2: MV rewrite (all) ===== + // Excluded via the whole RuleTypeClass.MATERIALIZE_VIEW family (see + // getMaterializedViewRuleNames()), not listed here. + + // ===== category 3: table pruning + UKFK ===== + "ELIMINATE_JOIN_BY_UK", // UK constraint JOIN elimination + "ELIMINATE_JOIN_BY_FK", // FK constraint JOIN elimination + "ELIMINATE_GROUP_BY_KEY", // UKFK GROUP BY key elimination + "ELIMINATE_GROUP_BY_KEY_BY_UNIFORM", // uniform-distribution GROUP BY key elimination + // ORDER BY key elimination by a declared UNIQUE key: ORDER BY a, b LIMIT 1 can + // collapse to ORDER BY a LIMIT 1, and the frozen SQL keeps the reduced order. + // Dropping the UNIQUE declaration (or its backing constraint state) changes no + // value the fingerprint hashes, so the replay may pick another row with the + // same smallest a - exclude the rule like the other UKFK ones. + "ELIMINATE_ORDER_BY_KEY", + // PK/FK-derived aggregate push down below the (FK) join: the rewritten + // topology stops being correct once the constraint state changes. The rule + // derives its rewrite from canEliminateByFk, so dropping the constraints and + // adding duplicate keys on the former primary side makes the original + // aggregate above the multiplying join return one doubled group while the + // frozen pre-aggregate replay returns duplicate undoubled rows - a wrong + // result, not just a missed rewrite. Audit (mutable-constraint consumers in + // the whitelist): the remaining canEliminateByFk / canEliminateByUk consumers + // are EliminateJoinByFK / EliminateJoinByUK (both excluded above) and the MV + // comparator family (whole RuleTypeClass excluded); this rule was the only + // unexcluded one. + "PUSH_DOWN_AGG_THROUGH_JOIN_ON_PKFK", + // uniqueness-dependent rewrites: each fires only when DataTrait proves the + // relevant slots UNIQUE (and NOT NULL), and that proof may come from a + // DECLARED UNIQUE constraint - which schemaFingerprint does not capture (it + // hashes only the table id + base columns), so dropping the declaration + // after freezing would leave the rewrite in place: + // - SIMPLIFY_WINDOW_EXPRESSION / AGG_SCALAR_SUBQUERY_TO_WINDOW_FUNCTION: + // Window(SUM(v) PARTITION BY k) -> Project(v) -> Scan(t) once t declares + // UNIQUE(k); after dropping the declaration and adding (k,10),(k,20) the + // original window returns 30,30 while the frozen replay returns 10,20; + // - PUSH_DOWN_TOP_N_DISTINCT_THROUGH_JOIN / _PROJECT_JOIN: a hard TopN + // pushed through the join / project retains rows justified by the + // declared uniqueness. + "SIMPLIFY_WINDOW_EXPRESSION", + "AGG_SCALAR_SUBQUERY_TO_WINDOW_FUNCTION", + "PUSH_DOWN_TOP_N_DISTINCT_THROUGH_JOIN", + "PUSH_DOWN_TOP_N_DISTINCT_THROUGH_PROJECT_JOIN", + + // ===== category 4: equivalence derivation ===== + "INFER_PREDICATES", // predicate derivation + "INFER_FILTER_NOT_NULL", // Filter NOT NULL derivation + "INFER_JOIN_NOT_NULL", // JOIN NOT NULL derivation + "CONSTANT_PROPAGATION", // constant propagation + + // ===== category 5: structural rewrite ===== + "EXTRACT_SINGLE_TABLE_EXPRESSION_FROM_DISJUNCTION", // split OR into single table + "OR_EXPANSION", // OR expansion into UNION + "PUSH_DOWN_FILTER_THROUGH_WINDOW", // window predicate push down + "ELIMINATE_AGG_CASE_WHEN", // aggregate CASE WHEN elimination + "ELIMINATE_OUTER_JOIN", // outer join elimination + "ELIMINATE_LIMIT", // LIMIT elimination + // ELIMINATE_LIMIT is registered in the same rule class (EliminateLimit) as + // ELIMINATE_LIMIT_ON_ONE_ROW_RELATION; excluding only the former let + // "SELECT 1 LIMIT 1" freeze the one-row child WITHOUT its LIMIT, while the + // sibling digest of "SELECT 1 LIMIT 0" still matched it (top-level LIMIT + // values are deliberately ignored during matching) - the replay then + // returned one row instead of none. + "ELIMINATE_LIMIT_ON_ONE_ROW_RELATION", + "ELIMINATE_AGGREGATE", // aggregate elimination + // two-phase LIMIT split GLOBAL(l, o) -> LOCAL(l + o, 0): the freeze keeps + // only the UPPER limit topology, the decompiled SQL then carries both + // phases as two query blocks, and the rewrite-time LIMIT merge (which only + // reaches the block that has a user-tree counterpart) can never grow the + // inner one - a captured order-free LIMIT 10 kept returning 10 rows when a + // matching user query asked for LIMIT 20. The split is an execution + // detail; freezing the single-phase limit keeps LIMIT ... OFFSET + // semantics (see SPMPlan2SQLBuilder.visitPhysicalLimit for the + // collapse of already-frozen pairs). + "SPLIT_LIMIT", + + // ===== category 6: external sources / empty relations (data dependent) ===== + "PUSH_FILTER_INTO_SCHEMA_SCAN", // schema table predicate push down + "ELIMINATE_JOIN_ON_EMPTYRELATION", // empty-relation operator elimination + "ELIMINATE_FILTER_ON_EMPTYRELATION", + "ELIMINATE_AGG_ON_EMPTYRELATION", + "ELIMINATE_PROJECT_ON_EMPTYRELATION", + "ELIMINATE_UNION_ON_EMPTYRELATION", + "ELIMINATE_TOPN_ON_EMPTYRELATION", + "ELIMINATE_SORT_ON_EMPTYRELATION", + "ELIMINATE_INTERSECTION_ON_EMPTYRELATION", + "ELIMINATE_EXCEPT_ON_EMPTYRELATION", + "ELIMINATE_LIMIT_ON_EMPTY_RELATION", + "PRUNE_EMPTY_PARTITION" + ); + + /** + * SET_VAR keys that carry the SPM safety overrides (see optimize()): a plan-SQL hint + * must never clear or weaken them for the nested statement, otherwise the frozen plan + * could be produced with state-sensitive rewrites re-enabled (e.g. + * SET_VAR(enable_nereids_rules='') re-enables FK join elimination, so a join could be + * frozen away and return different rows once constraint state changes) or with TopN + * lazy materialization / CTE inlining undoing the plan-serialization guards. + */ + private static final Set<String> SPM_LOCKED_SET_VAR_KEYS = ImmutableSet.of( + SessionVariable.ENABLE_NEREIDS_RULES, + "topn_lazy_materialization_threshold", + SessionVariable.ENABLE_CTE_MATERIALIZE, + SessionVariable.INLINE_CTE_REFERENCED_THRESHOLD, + SessionVariable.CTE_INLINE_MODE); + + private SPMOptimizer() { + } + + /** + * Installs the REPLAY-side rule mask on a replaying statement's context: every + * MATERIALIZED_VIEW rewrite is forbidden while the frozen plan is re-planned. + * + * <p>The frozen plan was produced under the SPM whitelist (see {@link #optimize}) - + * every MV rewrite excluded - so its fingerprint pins the SOURCE tables. The replay + * is normally planned with the session's full rule set, and an MV that became + * eligible AFTER the baseline was created (the reported case: an async MTMV over the + * very query, built and refreshed later) then substitutes the MV's storage table for + * the source table. The post-plan fingerprint guard + * ({@code SPMPlanner#verifyReplayMetadata}, which pins the frozen plan's tables) + * rejects its own replay and, with the default {@code enable_spm_fallback=false}, the + * SELECT fails although the source table has not changed. + * + * <p>The mask is deliberately NARROWER than the CREATE-side whitelist: it forbids + * exactly the rules that SUBSTITUTE one table for another. The remaining + * SPM-excluded rules are optimization opportunities that a re-plan may legally take + * (they either keep the frozen tables - salt join, structural rewrites, predicate + * inference - or their absence would break semantics: ELIMINATE_LIMIT turns the + * CALLER's LIMIT 0 into an empty relation, and a full CREATE mask forbids it, which + * replayed {@code LIMIT 0} as one row). Rules that REMOVE a table (UK / FK join + * elimination, aggregate-to-constant) keep their established fail-closed behavior: + * the fingerprint mismatch surfaces through the caller's fallback policy. + * + * @param statementContext the REPLAYING statement's context; its rule cache must not + * have been computed yet (the callers install the mask before + * the planner runs) + */ + public static void installSpmReplayRuleMask(StatementContext statementContext) { + statementContext.setSpmExcludedRules(materializedViewRuleMask()); + } + + /** + * The replay mask: one bit per materialized-view rewrite rule + * ({@link RuleType#isMaterializedViewRule()}, so newly added MV rules are covered + * automatically). + * + * @return the forbidden-rule mask of the whole MV family + */ + private static BitSet materializedViewRuleMask() { + BitSet mask = new BitSet(); + for (RuleType ruleType : RuleType.values()) { + if (ruleType.isMaterializedViewRule()) { + mask.set(ruleType.ordinal()); + } + } + return mask; + } + + /** + * All materialized view rewrite rule names (RuleTypeClass.MATERIALIZE_VIEW), + * enumerated programmatically so newly added MV rules are excluded from the SPM + * whitelist automatically. + * + * @return the immutable list of MV rewrite RuleType names + */ + public static List<String> getMaterializedViewRuleNames() { + return ImmutableList.copyOf(Arrays.stream(RuleType.values()) + .filter(RuleType::isMaterializedViewRule) + .map(Enum::name) + .collect(Collectors.toList())); + } + + /** + * The full set of RuleType names excluded from the SPM whitelist: the explicit + * state-sensitive list plus the whole MV rewrite family (deduplicated). + * + * @return the immutable list of all SPM-excluded rule names + */ + public static List<String> getSpmExcludedRuleNames() { + return ImmutableList.copyOf(Stream.concat( + SPM_EXCLUDED_RULE_NAMES.stream(), + getMaterializedViewRuleNames().stream()) + .distinct() + .collect(Collectors.toList())); + } + + /** + * Builds the SPM rule whitelist written into enable_nereids_rules while a baseline is + * created: every RuleType name except the SPM-excluded set, intersected with the + * caller's existing whitelist when one is configured (a session whitelist is + * preserved, never widened). The engine turns a non-empty enable_nereids_rules into + * the statement-level forbidden-rule mask (StatementContext#getOrCacheDisableRules), + * so the excluded rules cannot apply during baseline creation. + * + * @param originalEnabled the previous raw enable_nereids_rules value (may be null / + * empty) + * @return the comma-separated whitelist (RuleType declaration order) + * @throws AnalysisException when the original value names an unknown rule, or when + * its intersection with the SPM whitelist is empty (the + * session whitelists only rules that SPM excludes) + */ + public static String buildSpmEnabledRules(String originalEnabled) throws AnalysisException { + Set<String> excluded = new LinkedHashSet<>(getSpmExcludedRuleNames()); + Set<String> original = new LinkedHashSet<>(); + if (originalEnabled != null && !originalEnabled.isEmpty()) { + for (String ruleName : originalEnabled.split(",")) { + String normalized = ruleName.trim().toUpperCase(Locale.ROOT); + if (normalized.isEmpty()) { + continue; + } + try { + RuleType.valueOf(normalized); + } catch (IllegalArgumentException e) { + throw new AnalysisException( + "Unknown rule in enable_nereids_rules: " + normalized); + } + original.add(normalized); + } + } + List<String> enabled = new ArrayList<>(); + for (RuleType ruleType : RuleType.values()) { + String name = ruleType.name(); + if (excluded.contains(name)) { + continue; + } + if (!original.isEmpty() && !original.contains(name)) { + continue; + } + enabled.add(name); + } + if (enabled.isEmpty()) { + throw new AnalysisException("enable_nereids_rules only whitelists rules that SPM excludes: " + + originalEnabled); + } + return String.join(",", enabled); + } + + /** + * Optimizes a planSql in SPM mode: parse -> analyze -> rewrite -> CBO optimize with + * the SPM rule whitelist installed in enable_nereids_rules, then return the best + * physical plan and its estimated cost. + * + * The original enable_nereids_rules value is cached before the call and restored in + * a finally block (disable_nereids_rules is not touched). A fresh StatementContext is + * used so the per-statement forbidden-rule cache + * (CascadesContext.getAndCacheDisableRules) reflects the whitelist. + * + * @param ctx the connect context (provides session variables and catalog) + * @param planSql the plan SQL to optimize (a SELECT statement, may contain SET_VAR + * hints) + * @return the optimization result (best physical plan + estimated cost) + * @throws UserException when the SQL cannot be parsed or planned + */ + public static OptimizeResult optimize(ConnectContext ctx, String planSql) throws UserException { + Plan parsed = new NereidsParser().parseSingle(planSql); + // A SELECT statement is parsed as a logical plan; DDL/DML parse to a Command. + if (!(parsed instanceof LogicalPlan) || parsed instanceof Command) { + throw new AnalysisException("SPM only supports SELECT statements: " + planSql); + } + return optimize(ctx, (LogicalPlan) parsed, planSql); + } + + /** + * Optimizes an already-parsed (possibly parameterized) plan tree in SPM mode: + * analyze -> rewrite -> CBO optimize with the SPM rule whitelist installed in + * enable_nereids_rules, then return the best physical plan and its estimated cost. + * + * This is the entry used by CREATE BASELINE when the plan tree is first + * parameterized (literals replaced by SpmConstVar / SpmConstList placeholders) and + * then optimized: the placeholders travel through the optimizer and survive into + * the physical plan, so the decompiled frozen planSql keeps the placeholder ids + * (matching the SR model where the frozen SQL is re-parsed and user values are + * substituted by id at rewrite time). + * + * @param ctx the connect context (provides session variables and catalog) + * @param logicalPlan the (possibly parameterized) SELECT plan tree + * @param originSql the original SQL text (used for the statement context / + * error messages) + * @return the optimization result (best physical plan + estimated cost) + * @throws UserException when the plan cannot be analyzed / planned + */ + public static OptimizeResult optimize(ConnectContext ctx, LogicalPlan logicalPlan, String originSql) + throws UserException { + // Plan-SQL hints are applied DURING analysis (after the SPM overrides below are + // installed): SelectHintSetVar writes the SAME SessionVariable object, so a hint + // such as SET_VAR(enable_nereids_rules='') would silently clear the whitelist and + // the TopN / CTE guards for the nested statement. Reject such hints up front - + // CREATE then keeps the user's planSql text instead of freezing an unguarded plan. + checkProtectedSetVarHints(logicalPlan); + StatementContext statementContext = new StatementContext(ctx, + new OriginStatement(originSql, 0)); + NereidsPlanner planner = new NereidsPlanner(statementContext); + + SessionVariable sessionVar = ctx.getSessionVariable(); + String originalEnabled = sessionVar.getEnableNereidsRulesStr(); + StatementContext originalCtx = ctx.getStatementContext(); + int originalTopnLazyThreshold = SessionVariable.getTopNLazyMaterializationThreshold(); + boolean originalCteMaterialize = sessionVar.enableCTEMaterialize; + int originalInlineCteThreshold = sessionVar.inlineCTEReferencedThreshold; + int originalCteInlineMode = sessionVar.cteInlineMode; + try { + // WHITELIST mode: only the rules SPM allows may apply while the baseline plan + // is produced. The mask is installed on THIS nested statement's context - NOT + // by redefining the public enable_nereids_rules variable, which would change + // its established behavior for every ordinary statement in the session (a + // session value like enable_nereids_rules='ELIMINATE_GROUP_BY_KEY_BY_UNIFORM' + // would forbid every binding / implementation rule not named there). The + // user's own disable_nereids_rules keeps applying on top of the mask. + statementContext.setSpmExcludedRules(buildSpmExcludedRuleMask(originalEnabled)); + // TopN lazy materialization is an execution detail (post-process): it prunes + // base-table columns from the physical plan and re-reads them later by rowid + // (PhysicalLazyMaterialize). Such pruned columns would be missing from the + // decompiled frozen planSql ("Unknown column in table list" on replay), so it + // is disabled while the baseline plan is produced - the frozen SQL must carry + // the full column set; a re-plan at rewrite time may still apply lazy + // materialization itself. + sessionVar.setTopNLazyMaterializationThreshold(-1); + // The WITH structure of the user query must survive into the frozen planSql: + // the decompiler turns PhysicalCTEAnchor / PhysicalCTEProducer / + // PhysicalCTEConsumer into a real WITH clause (one shared definition, + // referenced by alias), like StarRocks. Three settings are overridden while + // the baseline plan is produced: + // - enable_cte_materialize / inline_cte_referenced_threshold: by default the + // engine inlines a CTE with a single consumer (a common TPCDS shape), which + // deletes the anchor from the optimized plan before the decompiler can see it; + // - cte_inline_mode = -1: mode 0 (the default) builds an alternative fully + // inlined plan and uses it when consumer filters can eliminate union branches + // of the CTE body, which silently drops WITH for such queries (TPCDS q04/ + // q11/q74); SPM needs the anchored plan, and a re-plan at rewrite time still + // applies the user's own cte_inline_mode. + // CTEInline still inlines the CTEs a recursive CTE requires inlined + // (StatementContext mustInlineCTEs), so WITH RECURSIVE planning is unaffected. + sessionVar.enableCTEMaterialize = true; + sessionVar.inlineCTEReferencedThreshold = 0; + sessionVar.cteInlineMode = -1; + // Install the FRESH statement context for the whole nested plan: the table + // collector caches the SELECT tables into the ConnectContext's CURRENT + // statement context, and collectAndLockTable() then locks THAT object. When + // the caller (StmtExecutor for a CREATE BASELINE statement) already installed + // the outer command's context, leaving it active made collectRelation cache + // into the outer object while lock() ran on the fresh (empty) one: the nested + // plan was optimized / decompiled WITHOUT its metadata-stability locks and + // could race concurrent ALTER / DROP. The exact original (possibly null) is + // restored in the finally below. + ctx.setStatementContext(statementContext); + // planWithLock runs preprocess (SET_VAR hint) -> analyze -> rewrite -> + // optimize -> postProcess; distribution planning is not needed for the + // decompiler (the physical plan already carries distribution specs). + // The root physical plan is the RETURN value of planWithLock: the planner's + // physicalPlan field is only assigned through the lockCallback used by + // plan(), so planner.getPhysicalPlan() would be null here. + Plan resultPlan = planner.planWithLock(logicalPlan, + PhysicalProperties.ANY, ExplainLevel.NONE); + if (!(resultPlan instanceof PhysicalPlan)) { + throw new AnalysisException("SPM failed to plan SQL: " + originSql); + } + // Belt-and-braces for any hint path the pre-scan does not model: the protected + // variables must still carry the SPM values after planning, otherwise the plan + // may have been produced without a guard - fail the freeze instead of + // publishing it. + verifySpmOverridesIntact(sessionVar); + return new OptimizeResult((PhysicalPlan) resultPlan, + extractCost((PhysicalPlan) resultPlan)); + } finally { + sessionVar.setTopNLazyMaterializationThreshold(originalTopnLazyThreshold); + sessionVar.enableCTEMaterialize = originalCteMaterialize; + sessionVar.inlineCTEReferencedThreshold = originalInlineCteThreshold; + sessionVar.cteInlineMode = originalCteInlineMode; + ctx.setStatementContext(originalCtx); Review Comment: [P1] Close the nested statement context after baseline planning. An Iceberg bind can load its snapshot through the fresh context installed here; that context's connector scope owns a table-cache lease or tracked table, released only by `StatementContext.close()`. This `finally` restores the outer context without closing the nested one, so normal statement teardown closes only the outer context and every such CREATE/refresh can leave an Iceberg resource pinned. Close the nested context in this `finally` and cover successful and failing nested plans with a scoped closeable. -- 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]
