github-actions[bot] commented on code in PR #68499: URL: https://github.com/apache/doris/pull/68499#discussion_r4176998371
########## regression-test/suites/spm/test_spm_review_round27.groovy: ########## @@ -0,0 +1,180 @@ +// 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. + +suite("test_spm_review_round27", "spm") { + + // Twenty-seventh review round: replay side stability and the SHOW / CREATE surfaces. + // + // - #1: a baseline is produced under the SPM rule whitelist (every MV rewrite + // excluded), so its fingerprint freezes the SOURCE table. The replay used to be + // planned by the ordinary planner with the session's full rule set: an async MTMV + // that became eligible AFTER the baseline was created substituted its storage table + // for the source table, the post-plan fingerprint guard rejected its own replay and + // the SELECT failed (default enable_spm_fallback=false) although the source table + // had not changed. The replay now installs the same whitelist mask. + // - #2: scan selectors are compared per table in STATEMENT order: the ambiguous swap + // of PARTITION pins between self-join occurrences is rejected at CREATE. + // - #4: SHOW BASELINE PLANS WHERE 1 = 1 used to become pattern = '1' and searched + // SQL / status / source text for 1: it now raises the advertised analysis error. + // #3 (authoritative GLOBAL rows on a follower) needs two FEs: it is covered by + // BaselineManagerConcurrencyTest. + + sql """set enable_spm_rewrite = true""" + sql """set enable_spm_fallback = false""" + + def ownBaselines = { + sql("""SHOW BASELINE PLANS""").findAll { it[1].toString().contains("spm_r27_") } + } + def dropOwnBaselines = { + ownBaselines().each { row -> + sql """DROP BASELINE PLAN ${row[0]}""" + } + } + dropOwnBaselines() + + def explainOf = { String query -> sql("""EXPLAIN ${query}""").toString() } + def createBaseline = { String bind, String plan -> + (sql('CREATE GLOBAL BASELINE PLAN "' + bind + '" WITH "' + plan + '"')[0][0] as Long) + } + + // ==================== #4: no supported column -> analysis error ==================== + // (each statement must stay on ONE line: a test{} block treats every line as its own + // statement) + test { + sql 'SHOW BASELINE PLANS WHERE 1 = 1' + exception "only supports" + } + test { + sql 'SHOW BASELINE PLANS WHERE id + 0 = 1' + exception "only supports" + } + // the supported shapes keep working + assertEquals(0, sql("""SHOW BASELINE PLANS WHERE id = -1""").size()) + assertTrue(sql("""SHOW BASELINE PLANS LIKE '%'""").size() >= 0) + + // ==================== #1: an MTMV appears after the baseline ==================== + sql """DROP TABLE IF EXISTS spm_r27_mv_t""" + sql """ + CREATE TABLE spm_r27_mv_t (k INT, v INT) + DISTRIBUTED BY HASH(k) BUCKETS 1 + PROPERTIES("replication_num" = "1") + """ + sql """INSERT INTO spm_r27_mv_t VALUES (1, 1), (1, 2), (2, 3), (3, 4)""" + + sql """DROP MATERIALIZED VIEW IF EXISTS spm_r27_mv""" + // the baseline is frozen while NO eligible MV exists: the fingerprint pins the source + // table and the frozen plan reads it + String agg = "SELECT k, SUM(v) FROM spm_r27_mv_t GROUP BY k" + long aggId = createBaseline(agg, agg) + assertTrue(explainOf(agg).contains("SPM baseline hit: id=${aggId}"), + "the aggregate baseline must be hit: " + explainOf(agg)) + assertEquals("[[1, 3], [2, 3], [3, 4]]", sql(agg).sort().toString(), Review Comment: [P3] Record fixed SQL results with generated regression output. This aggregate check and the fixed row checks at lines 126, 151 and 176 use hardcoded `assertEquals` values, but the suite has no `.out` file. AGENTS.md Testing Standard 6 requires determined query results to use ordered `qt_` cases and runner-generated output, so these replay results can be inspected consistently. The fixed `group_concat` results in its suite need the same treatment. ########## fe/fe-core/src/main/java/org/apache/doris/plugin/audit/AuditLoader.java: ########## @@ -300,9 +382,14 @@ public synchronized void loadIfNecessary(boolean force) { } private void resetBatch(long currentTime) { - this.auditLogBuffer = new StringBuilder(); - this.lastLoadTimeAuditLog = currentTime; - this.auditLogNum = 0; + synchronized (this) { Review Comment: [P2] Keep the audit horizon until a committed batch is visible. A stream load can return `Publish Timeout` after commit while its rows remain unreadable, but this unconditional reset clears `batchOldestEventTime` and the response above is never checked. For a 10:00 event delayed in the loader until an 11:00 load, capture at 11:05 can advance with only the fixed overlap; when the row publishes later it is behind the watermark. Retain a fence for such batches until publication is confirmed, and cover a Publish Timeout followed by delayed visibility. ########## fe/fe-core/src/main/java/org/apache/doris/plugin/audit/AuditLoader.java: ########## @@ -143,6 +171,60 @@ public void exec(AuditEvent event) { private synchronized void assembleAudit(AuditEvent event) { fillLogBuffer(event, auditLogBuffer); ++auditLogNum; + long eventTime = event.timestamp; + if (eventTime > 0 && (batchOldestEventTime == 0 || eventTime < batchOldestEventTime)) { + batchOldestEventTime = eventTime; + } + } + + /** + * Start time (epoch millis, the {@code time} column of {@code audit_log}) of the OLDEST + * event this FE's builtin loader has accepted but not published yet - queued events plus + * the assembled, not yet flushed batch. Returns 0 when the loader is not running (or has + * nothing outstanding), i.e. when there is no known publication delay to retain. + * + * <p>The SPM capture uses this as a progress fence: its next scan window must still start + * at or before this instant, otherwise a row the local loader still owes (e.g. a query the + * {@code query_audit_log_timeout_ms} hold released late, or one sitting behind a slow + * stream load in the {@link #auditEventQueue}) would fall behind the advanced watermark Review Comment: [P2] Include upstream audit events in this publication horizon. Completed queries first enter `WorkloadRuntimeStatusMgr` and `AuditEventProcessor`; this method sees an event only after the processor's single worker calls `AuditLoader.exec`. An 11:40 event can be queued there, or dequeued but stalled in the earlier AuditLogBuilder plugin, while the loader is empty at the 12:00 capture scan. The checkpoint advances and publication at 12:02 leaves the row before the next fixed overlap at 11:55. Track queued and in-flight pre-loader events or use a durable publication watermark; test a stalled processor with an empty loader. ########## fe/fe-core/src/main/java/org/apache/doris/nereids/spm/capture/PlanCaptureManager.java: ########## @@ -0,0 +1,2256 @@ +// 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.FeConstants; +import org.apache.doris.common.UserException; +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.SPMUtils; +import org.apache.doris.nereids.spm.manager.BaselineManager; +import org.apache.doris.plugin.audit.AuditLoader; +import org.apache.doris.qe.AutoCloseConnectContext; +import org.apache.doris.qe.ConnectContext; +import org.apache.doris.qe.GlobalVariable; +import org.apache.doris.qe.SessionVariable; +import org.apache.doris.qe.SqlModeHelper; +import org.apache.doris.qe.VariableMgr; +import org.apache.doris.statistics.repository.ResultRow; +import org.apache.doris.statistics.util.StatisticsUtil; + +import com.google.common.annotations.VisibleForTesting; +import com.google.gson.Gson; +import com.google.gson.reflect.TypeToken; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.atomic.AtomicLong; +import java.util.function.LongSupplier; +import java.util.function.Supplier; + +/** + * 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 { + + /** + * Test seam replacing the live leadership probe of {@link #persistCheckpoint} (null in + * production). A capture cycle runs on the master, but an in-flight cycle can reach + * its checkpoint write AFTER a handoff (the daemon checks isMaster only at the cycle + * start), which a unit test cannot interleave otherwise. + */ + @VisibleForTesting + public static volatile java.util.function.BooleanSupplier checkpointLeadershipProbeForTest; + + /** + * Statement timeout (seconds) of the checkpoint read / write. The default + * StatisticsUtil overloads assign the ANALYZE timeout (43,200 seconds), so a stalled + * internal-table read or write could hold the single capture cycle for hours and + * delay every later capture / retry. Both operations are latency-sensitive: fail + * fast, keep the cycle consistent, retry next cycle. + */ + static final int CHECKPOINT_IO_TIMEOUT_SECONDS = 10; + + /** Bounded read-back attempts confirming the first reservation is VISIBLE. */ + private static final int CHECKPOINT_VISIBILITY_ATTEMPTS = 5; + + /** Delay between the reservation visibility reads (millis). */ + private static final long CHECKPOINT_VISIBILITY_RETRY_MILLIS = 200L; + + private static final Logger LOG = LogManager.getLogger(PlanCaptureManager.class); + + private static final PlanCaptureManager INSTANCE = new PlanCaptureManager(); + + /** + * 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; + + /** + * Audit pages ONE wakeup may consume. The daemon interval (default 3h) bounds how + * often the backlog is drained, so consuming a single page per wakeup left a window + * truncated at the page limit needing one extra interval per page - a window holding + * more than `plan_capture_max_batch_size` eligible rows per interval could never catch + * up. The drain stays BOUNDED so one cycle cannot run unboundedly long (each page is + * one bounded query plus its checkpoint write). + */ + private static final int MAX_PAGES_PER_CYCLE = 50; + + /** + * Wakeup delay used while a window is still pending after a cycle: the backlog drains + * promptly instead of one page per `plan_capture_interval_seconds`. + */ + private static final long PENDING_WINDOW_RESUME_INTERVAL_MS = 5_000L; + + /** + * In-memory budget of the queued retry candidates, in statement characters (the + * dominant part of a queued entry; the statement text is what the leader FE holds). + * One drain of {@link #MAX_PAGES_PER_CYCLE} pages can enqueue up to a page budget of + * failures per page, so a transient external-metadata outage (every capture fails + * with "table metadata unavailable") would otherwise retain tens of thousands of full + * statements - hundreds of megabytes - before the first cohort reaches its third + * attempt. When the budget is reached the DRAIN pauses: the pending window keeps its + * bounds and its cursor (the unconsumed rows stay reachable by the keyset scan) while + * the replay burns the queue down, and the daemon resumes promptly (see + * pendingWindowNeedsPromptResume). + */ + private static final long MAX_QUEUED_FAILURE_CHARS = 64L * 1024 * 1024; + + /** + * Test seam overriding {@link #MAX_QUEUED_FAILURE_CHARS} (null = the production + * budget): a unit test cannot queue tens of megabytes of statements just to reach it. + * Written through {@link #setQueuedFailureBudgetForTest}. + */ + private static volatile Long queuedFailureBudgetForTest; + + /** + * Queued failures replayed in ONE cycle (see {@link #replayQueuedFailures}): replanning + * a whole outage-sized queue every wakeup would keep the FE busy for the length of the + * outage itself. + */ + private static final int MAX_RETRY_REPLAY_PER_CYCLE = 1000; + + /** + * Bounded retries for a FAILED capture: the query id stays retryable for later + * overlapping scans until it either succeeds or reaches this attempt count. Marking + * the id before processing would make a transient failure permanent - the + * overlapping scans would skip the row and the watermark passes it long before the + * dedup map evicts the entry. + */ + private static final int MAX_CAPTURE_ATTEMPTS = 3; + + /** Durable checkpoint key: the internal table holds exactly one row. */ + private static final long CHECKPOINT_ID = 1L; + + /** Upper bound for the retry entries written into the checkpoint row (row size). */ + private static final int MAX_PERSISTED_RETRIES = 64; + + /** Table of the durable capture checkpoint (see InternalSchema). */ + private static final String CHECKPOINT_TABLE = + "`__internal_schema`.`spm_capture_checkpoint`"; + + private static final String CHECKPOINT_SELECT_SQL = + "SELECT `last_scan_timestamp`, `pending_window_start`, `pending_window_end`," + + " `cursor_query_time`, `cursor_time`, `cursor_query_id`," + + " `failed_attempts`, `retry_queue`, `cursor_tail`," + + " `min_query_time_ms`, `min_scan_rows`, `include_pattern`," + + " `exclude_pattern`, `scan_zone` FROM " + CHECKPOINT_TABLE + + " WHERE `id` = " + CHECKPOINT_ID + " ORDER BY `update_time` DESC LIMIT 1"; + + /** + * One UPSERT statement: the table is UNIQUE-key(id) with merge-on-write, so inserting + * the row again REPLACES it atomically. The previous delete-then-insert pair was two + * separately committed statements: a crash / leadership loss / timeout / failed + * INSERT after the DELETE left NO row for the next leader, which then derived a fresh + * window and permanently skipped the deleted pending window's unconsumed tail. + * + * The target columns are listed EXPLICITLY. The VALUES order below follows + * {@link org.apache.doris.catalog.InternalSchema#SPM_CAPTURE_CHECKPOINT_SCHEMA}, but + * the PHYSICAL order of an upgraded table can differ: the upgrade of a pre-existing + * table APPENDS the columns it adds ({@code InternalSchemaInitializer# + * upgradeSpmCaptureCheckpointSchema}), which used to place cursor_tail after + * update_time. A positional INSERT then shifts every value behind the first + * out-of-position column - the tail JSON was written into failed_attempts, the retry + * JSON into update_time and NOW() into cursor_tail - and the checkpoint write failed + * / persisted garbage. Address the columns by NAME instead: the write must stay + * correct on every physical layout, exactly like the (by-name) CHECKPOINT_SELECT_SQL + * read. + */ + private static final String CHECKPOINT_INSERT_SQL = + "INSERT INTO " + CHECKPOINT_TABLE + + " (`id`, `last_scan_timestamp`, `pending_window_start`, `pending_window_end`," + + " `cursor_query_time`, `cursor_time`, `cursor_query_id`, `cursor_tail`," + + " `failed_attempts`, `retry_queue`, `min_query_time_ms`, `min_scan_rows`," + + " `include_pattern`, `exclude_pattern`, `scan_zone`, `update_time`)" + + " VALUES (" + CHECKPOINT_ID + ", ${lastScan}, ${pendingStart}, ${pendingEnd}," + + " ${cursorQueryTime}, '${cursorTime}', '${cursorQueryId}', '${cursorTail}'," + + " '${failedAttempts}', '${retryQueue}', ${minQueryTimeMs}, ${minScanRows}," + + " '${includePattern}', '${excludePattern}', '${scanZone}', NOW())"; + + 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; + + /** + * The session time_zone (zone ID) the most recent scan PASS rendered its window bounds + * in - the zone the audited rows' {@code time} columns are stored in. Empty = never + * scanned (a fresh process follows the global zone). audit_log keeps the WRITER's + * local rendering, so a global time_zone change makes the already published rows + * invisible to bounds rendered in the new zone: while this differs from the current + * global zone, the next pass re-renders the window in this zone first and only then in + * the new one (see {@link #resolveScanPassZone} and the exhaustion branch of + * {@link #runCaptureCycle}). Persisted with the checkpoint so a takeover resumes the + * same rendering. + */ + private String lastScanZone = ""; + + /** + * Pending scan window of a TRUNCATED cycle: the (start, end) pair the resume cursor + * below belongs to. While set, every cycle keeps scanning the SAME window - the end + * must stay fixed until the window is fully consumed, because the next + * interval-derived window would start around this window's end and leave every row the + * cursor has not reached yet permanently out of scope. + */ + private long pendingWindowStart = 0; + private long pendingWindowEnd = 0; + + /** + * The FILTER SNAPSHOT the pending window was opened with (null while none is pending). + * The audit SQL and the in-memory {@link PlanCaptureFilter#shouldCapture} stage must + * judge one window's rows by the SAME thresholds: a `SET GLOBAL + * plan_capture_min_query_time_ms` between two pages of the same window otherwise made + * the SQL return rows the stale filter rejected terminally (they were consumed, + * never captured) or pushed already-passed rows behind the cursor where a LOWERED + * threshold could no longer reach them. The whole filter is pinned, so a pattern + * change applies from the next window on. + */ + private PlanCaptureFilter pendingWindowFilter; + + /** + * Set when a cycle could not finish its pending window (page budget / failed + * checkpoint write): the daemon then reschedules the next cycle promptly instead of + * waiting the full `plan_capture_interval_seconds` (default 3h), which would grow the + * backlog by one interval's worth of eligible rows per consumed page. + */ + private volatile boolean pendingWindowNeedsPromptResume; + + /** Query ids already handled in earlier (overlapping) windows. */ + private final Map<String, Boolean> processedQueryIds = new LinkedHashMap<>(); + + /** + * Failure attempts per query id (bounded retry, see MAX_CAPTURE_ATTEMPTS). An id is + * removed here when it succeeds or is given up on; the map is capped like the + * processed-id map so a long-running failure burst cannot grow unbounded. + */ + private final Map<String, Integer> failedCaptureAttempts = new LinkedHashMap<>(); + + /** + * Candidates whose capture failed and that still have retry budget: keyset pagination + * advances the scan cursor past their audit rows and the window overlap only re-reads + * recent rows, so they are REPLAYED one attempt per cycle from here. Bounded like the + * other query-id maps; entries leave on success, on give-up, or when the id turns + * terminal elsewhere. + */ + private final Map<String, CapturedQuery> failedCaptureQueue = new LinkedHashMap<>(); + + /** + * Statement characters currently retained by {@link #failedCaptureQueue} (see + * {@link #MAX_QUEUED_FAILURE_CHARS}). Guarded by the capture daemon thread (all + * mutations happen inside a cycle) plus the test seams. + */ + private long queuedFailureChars = 0; + + /** + * Pre-page checkpoint state of the page a queued retry was FIRST seen on: the durable + * checkpoint must never move past an entry that the persisted JSON drops + * ({@link #MAX_PERSISTED_RETRIES}) - keyset pagination has already moved beyond its + * audit row, so only a cursor BEFORE that row can reach it after a restart / handoff. + * First-wins (pages only move forward) and evicted together with the queue. + */ + private final Map<String, RetryAnchor> failedCaptureAnchors = new LinkedHashMap<>(); + + /** One pre-page checkpoint state (see {@link #failedCaptureAnchors}). */ + private static final class RetryAnchor { + final long lastScanTimestamp; + final long windowStart; + final long windowEnd; + final long cursorQueryTime; + final String cursorTime; + final String cursorQueryId; + final String cursorTail; + + /** + * The zone THIS page's bounds / cursor were RENDERED in. The anchor describes a + * position of the audit stream, and that position is only reachable when it is + * rendered in the same zone again (see {@link #resolveScanPassZone}): persisting + * the CURRENT cycle's zone beside an earlier window's bounds made a takeover + * scan those bounds in a zone the rows were never written under, and the omitted + * retries (the ones the truncated queue could not carry) stayed unreachable. + */ + final String scanZone; + + /** The filter snapshot the page was judged by (see pageStartFilter). */ + final PlanCaptureFilter filter; + + RetryAnchor(long lastScanTimestamp, long windowStart, long windowEnd, + long cursorQueryTime, String cursorTime, String cursorQueryId, + String cursorTail, String scanZone, PlanCaptureFilter filter) { + this.lastScanTimestamp = lastScanTimestamp; + this.windowStart = windowStart; + this.windowEnd = windowEnd; + this.cursorQueryTime = cursorQueryTime; + this.cursorTime = cursorTime; + this.cursorQueryId = cursorQueryId; + this.cursorTail = cursorTail; + this.scanZone = scanZone; + this.filter = filter; + } + } + + /** + * Resume cursor of a TRUNCATED scan window: the FULL ORDER BY key tuple of the last + * consumed row -- (time, query_time, query_id) plus the encoded tail (client_ip, + * sql_hash, scan_rows, return_rows, statement hash) that uniquely separates audit + * rows sharing the first three keys. CURSOR_ABSENT while no partial window is + * pending - a short batch advances the watermark instead. Zero and NULL query_time + * are VALID cursors (see AuditLogScanner.CURSOR_QUERY_TIME_NULL). + */ + private long cursorQueryTime = AuditLogScanner.CURSOR_ABSENT; + private String cursorTime = ""; + private String cursorQueryId = ""; + private String cursorTail = ""; + + /** + * The pre-page state (watermark + window bounds + cursor) the CURRENT cycle's scan + * started from. It is the durable fallback persisted when the retry state is + * truncated by {@link #MAX_PERSISTED_RETRIES} (see persistCheckpoint): the durable + * cursor must never move past retries the checkpoint can no longer carry, otherwise + * a restart / leader handoff neither replays them from the queue nor re-reads their + * audit rows (the keyset cursor is beyond them and they can age outside the + * five-minute overlap), silently losing those captures. + */ + private long pageStartLastScanTimestamp = 0; + private long pageStartWindowStart = 0; + private long pageStartWindowEnd = 0; + private long pageStartCursorQueryTime = AuditLogScanner.CURSOR_ABSENT; + private String pageStartCursorTime = ""; + private String pageStartCursorQueryId = ""; + private String pageStartCursorTail = ""; + + /** + * The zone the CURRENT page's bounds / cursor were rendered in (the pass zone of the + * cycle that opened the page, see {@link #resolveScanPassZone}). It travels with + * {@link #currentPageAnchor()} and is persisted whenever the retry state rewinds the + * durable cursor to a page anchor: a takeover must re-render those bounds in the + * SAME zone, otherwise the audit rows written under the anchor's rendering are + * invisible to the re-scan (see {@link RetryAnchor#scanZone}). + */ + private String pageStartZoneId = ""; + + /** + * The FILTER SNAPSHOT the CURRENT page was scanned with (the same value handed to + * {@link AuditLogScanner#scan}). It travels with {@link #currentPageAnchor()} and is + * persisted whenever the retry state rewinds the durable cursor to an anchor: the + * rows of that page were admitted (or filtered) by THESE thresholds / patterns, so a + * takeover must re-scan the rewound range with the same eligibility - judging the + * re-scan by a configuration that changed in between could terminally filter the + * omitted oldest failure before its retry is even reachable. + */ + private PlanCaptureFilter pageStartFilter; + + /** Whether the durable checkpoint was already consulted in this process. */ + private boolean checkpointLoaded = false; + + /** + * Whether THIS process has ever seen a durable checkpoint row - read it from the store + * or written by this process. While it is false, the window a cycle consumes exists + * only in memory: runCaptureCycle records that window BEFORE scanning (see the + * initial reservation there), so a takeover can still resume it. + */ + private boolean durableCheckpointObserved = false; + + /** + * Checkpoint read / write seams. Production talks to the internal table through + * StatisticsUtil; tests replace them to simulate a failing first read and to observe + * the exact statements a persist issues. + */ + private Supplier<List<ResultRow>> checkpointReader = () -> StatisticsUtil.executeQuery( + CHECKPOINT_SELECT_SQL, Collections.emptyMap(), CHECKPOINT_IO_TIMEOUT_SECONDS); + + /** One checkpoint write statement. */ + @VisibleForTesting + public interface CheckpointWriter { + void write(String sql, Map<String, String> params) throws Exception; + } + + private CheckpointWriter checkpointWriter = (sql, params) -> StatisticsUtil.execUpdate( + sql, params, CHECKPOINT_IO_TIMEOUT_SECONDS); + + /** + * Whether a scripted checkpoint read / write seam is installed (tests only). The + * leadership fence of {@link #persistCheckpoint} guards LIVE writes: a scripted store + * stands in for the internal table, exactly like the simulator stores in + * BaselineManager.assertLeaderForWrite. + */ + private boolean checkpointSeamsForTest = false; + + /** + * Whether the cloud-mode warning was already logged (the gate fires every cycle). + */ + private boolean cloudModeWarned = false; + + /** + * The start of the window the FIRST cycle would have consumed when its checkpoint read + * FAILED (0 = none). A failed read records no window and the daemon retries promptly; + * without this floor the first cycle that succeeds on an EMPTY store would derive its + * own [now-interval, now) and permanently skip the rows of the first attempted window + * (every later window starts even later). Cleared as soon as a durable row is read or + * the reserved window becomes durable. + */ + private long firstAttemptedWindowStart = 0; + + /** + * The LOCAL audit loader's publication horizon: the start time (epoch millis) of the + * oldest audit event this FE's loader has accepted but not published yet (0 = nothing + * outstanding). Production reads the live {@link AuditLoader}; tests replace it. + */ + private LongSupplier auditQueueHorizon = AuditLoader::oldestUnpublishedEventTime; + + // 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", + Math.max(1, VariableMgr.getDefaultSessionVariable().getPlanCaptureIntervalSeconds()) + * 1000L); + this.filter = buildFilterFromGlobal(); + } + + /** The pre-page state of the CURRENT page (anchors entries queued by this page). */ + private RetryAnchor currentPageAnchor() { + return new RetryAnchor(pageStartLastScanTimestamp, pageStartWindowStart, + pageStartWindowEnd, pageStartCursorQueryTime, pageStartCursorTime, + pageStartCursorQueryId, pageStartCursorTail, pageStartZoneId, pageStartFilter); + } + + 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; + } + runCaptureCycle(global, newFilter); + if (pendingWindowNeedsPromptResume) { + // A window this cycle could not finish (page budget reached or its progress + // not durable) must NOT wait another full interval: the next wakeup continues + // exactly where this one stopped. The configured interval would add one + // interval's worth of eligible rows per consumed page, so a window holding + // more than one page per interval would never drain. + setInterval(Math.min( + Math.max(1L, global.getPlanCaptureIntervalSeconds()) * 1000L, + PENDING_WINDOW_RESUME_INTERVAL_MS)); + } + } + + /** + * One capture cycle body: everything after the runtime guards (cloud mode, enable + * flag, leader / checkpoint-thread checks, filter refresh). Split out so unit tests + * can drive a FULL cycle - checkpoint read through window derivation, scan, state + * advance and persist - without the process-global guards (master / cloud / enable) + * a test environment cannot satisfy. + * + * @param global the global session variables of this cycle + * @param newFilter the filter refreshed for this cycle + */ + @VisibleForTesting + public void runCaptureCycle(SessionVariable global, PlanCaptureFilter newFilter) { + try { + // A restarted / newly promoted leader must NOT start from a fresh + // interval-derived window: a truncated window from the previous leader is + // checkpointed here, and skipping it would permanently exclude its unconsumed + // tail (the overlap only reaches rows younger than the NEW watermark). + // A FAILED read returns false and the cycle aborts BEFORE deriving or + // persisting anything: writing a freshly derived window while the previous + // leader's unconsumed tail is still unreadable would overwrite its only + // record (the write path shares the same internal table the read failed on). + if (!loadCheckpointIfNeeded()) { + // The read failed and recorded nothing, so this process still has no window. + // Remember the window this cycle WOULD have consumed and retry promptly: the + // internal-schema initializer is asynchronous, and a later cycle deriving its + // OWN [now-interval, now) would permanently skip every eligible short row of + // this first attempted window (no later overlap reaches behind a NEW window's + // start). Nothing is written here - an unreadable checkpoint must never be + // replaced by a freshly derived one. + long attemptedStart = System.currentTimeMillis() + - Math.max(1L, global.getPlanCaptureIntervalSeconds()) * 1000L; + if (firstAttemptedWindowStart == 0 || attemptedStart < firstAttemptedWindowStart) { + firstAttemptedWindowStart = attemptedStart; + } + pendingWindowNeedsPromptResume = true; + LOG.warn("Plan capture cycle skipped: durable checkpoint not confirmed"); + return; + } + + // A window that is still PENDING keeps the filter snapshot it was opened with: + // its rows behind the cursor were already judged by those thresholds, and the + // audit SQL must use exactly the same values (see AuditLogScanner#scan). The + // snapshot is chosen AFTER the checkpoint load, so a TAKEOVER continues the + // restored window with the restored thresholds in its very first cycle. A NEW + // window follows the filter refreshed for this cycle, so `SET GLOBAL + // plan_capture_min_query_time_ms` takes effect from the next window on. + PlanCaptureFilter cycleFilter = pendingWindowFilter != null + ? pendingWindowFilter : newFilter; + this.filter = cycleFilter; + pendingWindowNeedsPromptResume = false; + + long currentTime = System.currentTimeMillis(); + // a non-positive interval / batch size can never be written through SQL SET + // (see SessionVariable), but clamp defensively: an interval of 0 would make + // every window empty and a batch size of 0 would return LIMIT 0, mark the + // window exhausted and advance the watermark over every eligible row + long intervalMs = Math.max(1L, global.getPlanCaptureIntervalSeconds()) * 1000L; + int batchSize = Math.max(1, global.getPlanCaptureMaxBatchSize()); + // overlap the window so audit rows loaded late (published after their event + // time has passed) are still scanned; the overlap follows the audit loader's + // configured batch interval so rows written at the tail of a loader batch - + // whose event time predates the new watermark - are not lost. Duplicates are + // filtered by query id below. + long[] window = resolveScanWindow(lastScanTimestamp, pendingWindowStart, pendingWindowEnd, + currentTime, intervalMs, + scanWindowOverlapMs(GlobalVariable.auditPluginMaxBatchInternalSec, + auditQueueHorizon.getAsLong(), currentTime), Review Comment: [P2] Include follower audit queues in the capture horizon. Capture runs only on the master, but this supplier reads that FE's process-local AuditLoader; follower FEs also emit audit events to their own loaders and write the shared audit table. If a follower's 11:50 event remains queued until after the master completes its 12:00 window, the master's empty local queue allows progress past 11:50 and the later row falls outside the fixed overlap. The earlier queue-delay thread prompted this local fix, but it cannot see another FE's backlog. Fence progress with a cluster-wide publication horizon, and test a delayed follower load. -- 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]
