github-actions[bot] commented on code in PR #68499: URL: https://github.com/apache/doris/pull/68499#discussion_r4230313860
########## fe/fe-core/src/main/java/org/apache/doris/nereids/spm/manager/BaselineManager.java: ########## @@ -0,0 +1,6447 @@ +// 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.spm.SPMUtils; +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.Collection; +import java.util.Collections; +import java.util.Comparator; +import java.util.HashMap; +import java.util.HashSet; +import java.util.Iterator; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +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: bindSqlHash (Long) -> baseline id list (List of Long). 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 + * 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; + } + + /** + * The newest stored update_time of the WHOLE simulated table, in epoch SECONDS + * (0 = none): the updateStatus bump reads this, so a simulator that + * keeps future-bumped update_times must expose them for the bump to apply - the + * default keeps simulators that never store future times working. + * + * @return the newest stored update_time in seconds, 0 when unavailable + */ + default long newestStoredUpdateSecond() { + return 0; + } + } + + /** + * 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) { + } + + /** + * The identity-carrying reservation: routes to the simulator's own + * storage of pendingSeqReservation. The default keeps simulators that + * model no keyed reservations working. The record fences a retry of the same key + * while its id stays unreadable and it is younger than the durable fence + * + * @param id the reserved id + * @param bindSqlDigest the baseline's bind digest + * @param planSqlHash the hash of the baseline's plan SQL + * @param reserveTimeMs the reservation instant (epoch millis) + */ + default void reserveId(long id, String bindSqlDigest, long planSqlHash, + long reserveTimeMs) { + reserveId(id); + } + + /** + * Records the DURABLE pending marker of an ambiguous create: the + * same identity-carrying append with unconfirmed = 1. The default keeps + * simulators without keyed reservations working. + * + * @param bindSqlDigest the baseline's bind digest + * @param planSqlHash the hash of the baseline's plan SQL + * @param id the id the ambiguous write consumed + * @param atMillis the marker instant (epoch millis) + */ + default void notePendingSeqState(String bindSqlDigest, long planSqlHash, long id, + long atMillis) { + } + + /** + * The latest identity record of one baseline in the simulated sequence table, or + * null: the durable unconfirmed-create fence of + * resolveDurablePendingCreate reads it. Either the explicit UNCONFIRMED + * marker of an ambiguous write, or - the plain reservation row + * appended before every baseline write (the marker wins when both exist). The + * default keeps simulators without keyed records working (no record = no fence). + * + * @param bindSqlDigest the baseline's bind digest + * @param planSqlHash the hash of the baseline's plan SQL + * @return the record, or null when none exists + */ + default SeqReservation pendingSeqReservation(String bindSqlDigest, long planSqlHash) { + return null; + } + + /** + * Retires the durable UNCONFIRMED marker of a RESOLVED ambiguous write + * the simulated equivalent of the marker DELETE. The default + * keeps simulators without keyed markers working. + * + * @param bindSqlDigest the baseline's bind digest + * @param planSqlHash the hash of the baseline's plan SQL + * @param markerId the id the marker was appended for + */ + default void retirePendingSeqState(String bindSqlDigest, long planSqlHash, + long markerId) { + } + + /** + * Records a DROP TOMBSTONE: the append-only identity of a baseline + * the DROP removed. The default keeps simulators without tombstones working (the + * load filter then finds no marker). + * + * @param id the dropped baseline id + * @param bindSqlDigest the baseline's bind digest + * @param planSqlHash the hash of the baseline's plan SQL + * @param atMillis the marker instant (epoch millis) + */ + default void appendDroppedMarker(long id, String bindSqlDigest, long planSqlHash, + long atMillis) { + } + + /** + * The recorded DROP TOMBSTONES as id|bindSqlDigest|planSqlHash keys; the + * load filter ignores rows matching one. The default keeps simulators working. + * + * @return the recorded tombstone keys (empty when none) + */ + default List<String> droppedMarkers() { + return List.of(); + } + + /** + * The newest stored update_time of the WHOLE simulated table, in epoch SECONDS + * (0 = none): see StatusProtocolStoreForTest#newestStoredUpdateSecond() + * + * @return the newest stored update_time in seconds, 0 when unavailable + */ + default long newestStoredUpdateSecond() { + return 0; + } + + 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 for the compact id high-water-mark RECORD (the value of SELECT_HWM_SQL), + * null = the internal table. A scripted value also stands in for the legacy history + * read: the seam answers the whole watermark decision (see readCompactIdWatermark). + */ + @VisibleForTesting + public static volatile java.util.function.LongSupplier hwmRecordReadForTest; + + /** + * Test seam for the scoped sequence-tail read (the value SELECT_SEQ_TAIL_SQL + * returns), null = the internal table. Only consulted with hwmRecordReadForTest + * installed - the seam pair models the two stores the watermark read consults. + */ + @VisibleForTesting + public static volatile java.util.function.LongSupplier seqTailReadForTest; + + /** + * Test seam replacing the live leadership probe of 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 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 + * 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 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 + * 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 + * refreshAfterForwardedDdl and 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; + + /** The compact id high-water-mark table (see InternalSchema#SPM_BASELINES_HWM_TBL_NAME). */ + private static final String SPM_BASELINES_HWM_TABLE = + FeConstants.INTERNAL_DB_NAME + "." + InternalSchema.SPM_BASELINES_HWM_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`, `plan_sql_digest` FROM "; + + /** + * First page of a whole-table snapshot: ordered by id so the pagination can continue + * with SELECT_PAGE_SQL from the last row read. The order must be a TOTAL + * order over the rows of ONE id - (update_time, status) break the id ties + * and the CONTENT columns break the remaining ties: two + * masters can leave two DIFFERENT rows of one id with the same stored second and the + * same status (a delayed INSERT committing after a handoff collision), and with + * ORDER BY `id` alone the engine may return such rows in ANY order in EVERY + * execution, so an OFFSET continuation landing inside the group could re-read one row + * and skip another while the row COUNT - the completeness proof - stays unchanged. + * Rows equal in EVERY column remain interchangeable (they resolve to the same + * pickDurableWinner outcome). + */ + private static final String SELECT_ALL_ORDERED_SQL = + SNAPSHOT_COLUMNS + SPM_BASELINES_TABLE + " ORDER BY `id`, `update_time`, `status`," + + " `bind_sql_digest`, `plan_sql`, `bind_sql`"; + + /** + * One continuation page of a whole-table snapshot: every row with id >= + * ${lastId}, ordered by id and SKIPPING the first ${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 SNAPSHOT_PAGE_SIZE rows - a jump past it would + * omit the rows behind the first page, possibly the newest durable status. The order + * must remain a TOTAL order within one id (see SELECT_ALL_ORDERED_SQL) or a + * tie order that differs between the two queries can make the offset skip a row of + * the interrupted group (content tie-breakers added). + */ + private static final String SELECT_PAGE_SQL = SNAPSHOT_COLUMNS + SPM_BASELINES_TABLE + + " WHERE `id` >= ${lastId} ORDER BY `id`, `update_time`, `status`," + + " `bind_sql_digest`, `plan_sql`, `bind_sql`" + + " LIMIT ${pageSize} OFFSET ${offset}"; + + /** + * Rows per snapshot page (see readPersistedSnapshot). Bounds what ONE + * internal query has to return, so a growing table can no longer make the whole + * snapshot read fail against a fixed timeout. + */ + private static final int SNAPSHOT_PAGE_SIZE = 2000; + + /** + * Fence re-reads a paginated snapshot read is allowed before it fails closed (see + * readStableSnapshot): one DDL overlapping the loop then converges on the + * retry, while a table that never stays stable must not be published as a snapshot. + */ + private static final int SNAPSHOT_STABILITY_ATTEMPTS = 3; + + /** The persistence-layer id watermark (see the class javadoc "Id source"): read + * before every id allocation. MAX over an aggregate is a light single-row query. */ + private static final String SELECT_MAX_ID_SQL = "SELECT MAX(`id`) FROM " + SPM_BASELINES_TABLE; + + /** + * The compact id high-water mark: a tiny append-only table whose rows + * carry the newest allocated id. The append-only HISTORY table + * (InternalSchema#SPM_BASELINES_SEQ_TBL_NAME) grows by one row per create + * forever, so its MAX(last_id) - the only unbounded read on the create path - + * was replaced: every allocation also records itself here (pruning the superseded + * rows), and a pre-upgrade cluster only pays the legacy full read ONCE (see + * readPersistedWatermark). + */ + private static final String SELECT_HWM_SQL = "SELECT MAX(`last_id`) FROM " + + SPM_BASELINES_HWM_TABLE + " WHERE `id` = 1"; + + /** Append one high-water-mark row (see SELECT_HWM_SQL). */ + private static final String INSERT_HWM_SQL = "INSERT INTO " + SPM_BASELINES_HWM_TABLE + + " (`id`, `last_id`, `update_time`) VALUES (1, ${lastId}, NOW())"; + + /** + * Best-effort prune of the compact high-water-mark rows (see SELECT_HWM_SQL): + * removes the rows the just-written one supersedes, so the surviving MAX read + * stays a scan of a handful of rows. Correctness never depends on it - the read takes + * the MAX - so failures are swallowed. + */ + private static final String PRUNE_HWM_SQL = "DELETE FROM " + SPM_BASELINES_HWM_TABLE + + " WHERE `last_id` < ${lastId}"; + + /** + * The scoped CONFIRMATION of SELECT_HWM_SQL (see readCompactIdWatermark): the + * highest reservation VISIBLE beyond the recorded mark. The sequence table's key + * starts with (id, last_id), so the read scans only the rows past the mark - a + * healthy cluster has none, and a stale record never hides a reservation whose + * HWM write was lost or is still unreadable. + */ + private static final String SELECT_SEQ_TAIL_SQL = "SELECT MAX(`last_id`) FROM " + + SPM_BASELINES_SEQ_TABLE + " WHERE `id` = 1 AND `last_id` > ${floor}"; + + /** + * The id high-water mark that OUTLIVES the rows (see + * InternalSchema#SPM_BASELINES_SEQ_TBL_NAME): MAX(last_id) over the append-only + * reservation rows. Kept as the ONE-TIME legacy fallback of + * readPersistedWatermark (a cluster created before the compact high-water + * mark table exists); the per-create path reads the bounded + * SELECT_HWM_SQL instead. The baselines table's own MAX(id) falls back to a + * lower value as soon as its highest row is DROPped, and an id reused for a DIFFERENT + * baseline would let a delayed DROP BASELINE PLAN IF EXISTS N retry delete the + * new baseline. + */ + private static final String SELECT_SEQ_ID_SQL = "SELECT MAX(`last_id`) FROM " + + SPM_BASELINES_SEQ_TABLE; + + /** Appends one reservation row (the id just allocated). Append-only: MAX never falls. */ + private static final String INSERT_SEQ_ID_SQL = "INSERT INTO " + SPM_BASELINES_SEQ_TABLE + + " (`id`, `last_id`, `bind_sql_digest`, `plan_sql_hash`, `reserve_time`," + + " `unconfirmed`, `dropped`)" + + " VALUES (1, ${lastId}, '${bindSqlDigest}', ${planSqlHash}, '${reserveTime}'," + + " ${unconfirmed}, 0)"; + + /** + * Appends a DROP TOMBSTONE: the identity of a baseline this FE just + * removed, with dropped = 1. A demoted master's in-flight status INSERT can + * commit AFTER the DROP deleted the row - its conditional precondition ran against + * the pre-DROP snapshot, and Doris cannot re-check it at durable commit - and the + * revived row would make the dropped baseline ACTIVE again on every loader. The + * tombstone is APPEND-ONLY and survives that commit: a load that sees a baseline row + * matching a tombstone's (id, bind_sql_digest, plan_sql_hash) treats it as deleted + * (and repairs it away). Ids are never reused (the sequence watermark), so a matching + * tombstone always describes this very incarnation. + */ + private static final String INSERT_SEQ_DROPPED_SQL = "INSERT INTO " + SPM_BASELINES_SEQ_TABLE + + " (`id`, `last_id`, `bind_sql_digest`, `plan_sql_hash`, `reserve_time`," + + " `unconfirmed`, `dropped`)" + + " VALUES (1, ${lastId}, '${bindSqlDigest}', ${planSqlHash}, '${reserveTime}'," + + " 0, 1)"; + + /** + * Reads the DROP TOMBSTONES of the GIVEN ids (see INSERT_SEQ_DROPPED_SQL, + * the append-only sequence table retains every dropped = 1 row + * forever, so an unrestricted read built a HashSet of the FULL historical drop set on + * every load / refresh - with a small active set and heavy CREATE / DROP churn the + * read grew without bound and, once it timed out, follower caches stopped + * incorporating later GLOBAL changes. The read is scoped to the ids the caller is + * actually filtering (the snapshot / point-read ids), chunked to keep each statement + * bounded. + */ + private static final String SELECT_SEQ_DROPPED_SQL = "SELECT `last_id`, `bind_sql_digest`," + + " `plan_sql_hash` FROM " + SPM_BASELINES_SEQ_TABLE + " WHERE `dropped` = 1" + + " AND `last_id` IN (${ids})"; + + /** How many ids one scoped tombstone read carries (see SELECT_SEQ_DROPPED_SQL). */ + private static final int DROPPED_MARKER_ID_CHUNK = 256; + + /** + * The LATEST identity-carrying row of one baseline: the durable half of the + * unresolved-create fence (see + * resolveDurablePendingCreate). BOTH kinds of rows fence while their write's + * outcome is unresolved: + * unconfirmed = 1: the marker of a create whose INSERT outcome was + * AMBIGUOUS; + * a PLAIN reservation (unconfirmed = 0): the row every create appends + * BEFORE its INSERT. Its separate ambiguous marker can fail / lag (that write + * is best effort), and a cross-FE retry that sees neither the baseline row nor + * the marker then allocated a SECOND id whose committed row later published a + * duplicate; the pre-INSERT reservation is the one record that always exists, + * so it fences too - but only while its row is not readable and its age is + * inside the fence bound. + * dropped = 1: a tombstone (a completed DROP or a condemned abandoned + * write) RESOLVED the identity - no fence, the key may be created again. + * The NEWEST IDENTITY wins (highest last_id): ids name the successive + * incarnations of one key, and reserve_time has only SECOND precision - a DROP of + * id N followed by a re-create as N+1 within the same stored second used to sort + * N's dropped = 1 tombstone BEFORE N+1's reservation, so resolveDurablePendingCreate + * read the key as resolved, skipped the pending fence of N+1's committed-but- + * unreadable INSERT, and allocated N+2 (both ENABLED rows could then publish). + * Within one identity the latest state wins: a tombstone (appended after the + * reservation) means resolved; the same-second order is tombstone, then marker, + * then the plain row. + */ + private static final String SELECT_PENDING_SEQ_SQL = "SELECT `last_id`, `reserve_time`," + + " `unconfirmed`, `dropped` FROM " + + SPM_BASELINES_SEQ_TABLE + " WHERE `bind_sql_digest` = '${bindSqlDigest}'" + + " AND `plan_sql_hash` = ${planSqlHash}" + + " ORDER BY `last_id` DESC, `reserve_time` DESC, `dropped` DESC," + + " `unconfirmed` DESC LIMIT 1"; + + /** + * The base of the synthetic `id` every COMPACT identity row carries (see + * INSERT_COMPACT_SEQ_SQL). The sequence table's `id` column is its DUPLICATE key and + * every historical row carries the constant 1: the identity read filtered on + * bind_sql_digest / plan_sql_hash alone, which no sort order can bound - the table is + * append-only and grows with every create / drop forever, so one identity read grew + * into a full scan of the whole history. Each state change ALSO appends a compact + * copy under this stable per-identity id, so the read seeks the key prefix and sees + * only that identity's few surviving rows (the base keeps them clear of the + * historical id = 1 rows and of the compact HWM rows of the OTHER tables). + */ + private static final long COMPACT_SEQ_ID_BASE = 2L; + + /** Appends the compact copy of one identity state (see COMPACT_SEQ_ID_BASE). */ + private static final String INSERT_COMPACT_SEQ_SQL = "INSERT INTO " + SPM_BASELINES_SEQ_TABLE + + " (`id`, `last_id`, `bind_sql_digest`, `plan_sql_hash`, `reserve_time`," + + " `unconfirmed`, `dropped`)" + + " VALUES (${compactId}, ${lastId}, '${bindSqlDigest}', ${planSqlHash}," + + " '${reserveTime}', ${unconfirmed}, ${dropped})"; + + /** + * Best-effort prune of the compact rows the just-written one supersedes + * (see COMPACT_SEQ_ID_BASE): without it one identity's slot grows with every state + * change. Correctness never depends on it - the read takes the MAX within the slot - + * so failures are swallowed like PRUNE_HWM_SQL's. A compact id is a 64-bit HASH, so + * two identities may share one slot: their rows carry their own digest / hash + * predicates and the read falls back to the history when the slot's newest row is + * another identity's. + */ + private static final String PRUNE_COMPACT_SEQ_SQL = "DELETE FROM " + + SPM_BASELINES_SEQ_TABLE + " WHERE `id` = ${compactId} AND `last_id` < ${lastId}"; + + /** + * The BOUNDED identity read (see COMPACT_SEQ_ID_BASE): the same answer as + * SELECT_PENDING_SEQ_SQL, sought through the `id` key prefix, with the identity + * predicates kept because a compact id is a hash. SELECT_PENDING_SEQ_SQL remains the + * fallback for a slot without a readable row. + */ + private static final String SELECT_COMPACT_SEQ_SQL = "SELECT `last_id`, `reserve_time`," + + " `unconfirmed`, `dropped` FROM " + + SPM_BASELINES_SEQ_TABLE + " WHERE `id` = ${compactId}" + + " AND `bind_sql_digest` = '${bindSqlDigest}'" + + " AND `plan_sql_hash` = ${planSqlHash}" + + " ORDER BY `last_id` DESC, `reserve_time` DESC, `dropped` DESC," + + " `unconfirmed` DESC LIMIT 1"; + + /** + * Retires the UNCONFIRMED marker(s) of ONE resolved ambiguous write: + * flipped by retireSeqPendingMarker. The DELETE touches only + * unconfirmed = 1 rows of that last_id - the plain reservation row appended + * before every create (and every other id ever reserved) stays, so MAX(last_id) and + * with it the id watermark never fall. + */ + private static final String DELETE_PENDING_SEQ_MARKER_SQL = "DELETE FROM " + + SPM_BASELINES_SEQ_TABLE + " WHERE `bind_sql_digest` = '${bindSqlDigest}'" + + " AND `plan_sql_hash` = ${planSqlHash} AND `unconfirmed` = 1" + + " AND `last_id` = ${lastId}"; + + /** + * The consistency fence of the paginated snapshot read (see + * readStableSnapshot): the id high-water mark, the row count and the newest + * update_time of the WHOLE table, read before AND after the page loop. An internal + * paginated snapshot issues one SELECT per page and the internal table has no + * long-lived read view, so a CREATE / ALTER / DROP committing between two pages would + * otherwise be merged into a state no single point in time ever had (the reviewer's + * example: a follower reads ENABLED low-id A on page 1, the master drops A and creates + * high-id B before page 2 - the published cache then contains BOTH, SHOW reports the + * completed DROP and matching replays A until the next refresh). MAX(id) catches every + * CREATE, COUNT(*) catches a pure DROP, MAX(update_time) catches a status flip (SECOND + * precision: a flip within the same second as the previous write remains a residual + * window, closed by the refresh daemon). + */ + private static final String SELECT_SNAPSHOT_FENCE_SQL = "SELECT MAX(`id`), COUNT(*)," + + " MAX(`update_time`) FROM " + SPM_BASELINES_TABLE; Review Comment: [P2] Bound the table-wide fence used by refresh and status changes. `readStableRawSnapshot` runs this unfiltered `MAX(id), COUNT(*), MAX(update_time)` before and after its 2,000-row page loop, each with a 10-second timeout; `readNewestStoredUpdateSecond` also runs an unfiltered `MAX(update_time)` before every real GLOBAL ENABLE/DISABLE. Pagination bounds the page SELECTs, not these statements. Segment statistics can reduce their cost, but status DELETE predicates can force row scans, and metadata work still grows with segments. On a sufficiently large history, a single timeout stops refresh/confirmed SHOW or prevents disabling a baseline. Maintain a bounded durable mutation clock/fence instead of table-wide aggregates on these paths. -- 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]
