u70b3 commented on code in PR #67978: URL: https://github.com/apache/doris/pull/67978#discussion_r4144170557
########## fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java: ########## @@ -0,0 +1,750 @@ +// 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.datasource.lance.job; + +import org.apache.doris.catalog.Env; +import org.apache.doris.catalog.DatabaseIf; +import org.apache.doris.catalog.TableIf; +import org.apache.doris.common.ClientPool; +import org.apache.doris.common.Config; +import org.apache.doris.common.util.MasterDaemon; +import org.apache.doris.datasource.CatalogIf; +import org.apache.doris.datasource.lance.LanceExternalCatalog; +import org.apache.doris.datasource.lance.LanceIndexDatasetCheck; +import org.apache.doris.datasource.lance.storage.LanceStorageOptions; +import org.apache.doris.persist.gson.GsonUtils; +import org.apache.doris.system.Backend; +import org.apache.doris.system.BeSelectionPolicy; +import org.apache.doris.system.SystemInfoService; +import org.apache.doris.thrift.BackendService; +import org.apache.doris.thrift.TLanceIndexJobDispatch; +import org.apache.doris.thrift.TLanceIndexMutationType; +import org.apache.doris.thrift.TNetworkAddress; +import org.apache.doris.thrift.TStatus; +import org.apache.doris.thrift.TStatusCode; + +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; +import org.apache.thrift.TApplicationException; + +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.UUID; +import java.util.function.Supplier; + +/** + * Master-only daemon that drives the durable Lance index job records through + * the lifecycle after admission. Each round runs in a fixed order: converge + * expired RUNNING jobs to UNKNOWN, release possible-live slots whose backend + * process was replaced, drive the refresh a terminal job still owes, then + * dispatch PENDING jobs. Every durable transition goes through + * {@link LanceIndexJobManager} under its own lock; the daemon holds no catalog + * or manager lock across any call. + * + * <p>The daemon does not read the admission gate: a job that is already durable + * must be driven to its terminal state, whatever the gate says now, so the + * thread runs unconditionally on the master and simply finds nothing to do + * while no jobs exist. An idle round writes no journal record. + * + * <p>Dispatch follows the durable-before-send boundary: the whole request is + * prepared first (so a preparation failure just leaves the job PENDING), then + * the markRunning edit log is written and re-read before the first byte of + * network I/O, and the invocation id of an attempt that lost the compare-and-set + * is never reused. After a successful markRunning there is exactly one send; + * from that point a job converges only through a matching result callback, the + * deadline sweep, or the epoch sweep, never through a resend. A failure that + * still proves the dispatch was never enqueued (a clean pre-enqueue error + * status, a client-pool borrow failure, or an UNKNOWN_METHOD answer from an + * old backend) converges it NOT_COMMITTED through the no-enqueue channel, + * which releases the possible-live slot in the same durable transition; + * anything ambiguous after the invocation may have started converges UNKNOWN + * with the slot retained. The blocking time of the send loop is bounded per + * round (one backend RPC timeout), because this thread is also the only thread + * running the sweeps and the refresh driver — see + * {@link #dispatchPendingJobs()}. + * + * <p>The manager is resolved from the supplier once per round rather than + * captured at construction: {@code Env.loadLanceIndexJobManager} replaces the + * Env-owned manager with a brand-new object on every image load, so a cached + * reference would keep scanning the abandoned pre-image manager after an FE + * restart while replay, admission and SHOW all move on to the restored one. + * Every phase of one round shares the single resolved instance. + * + * <p>The sleep between rounds is sliced at {@link #MAX_SLEEP_SLICE_MS} so a + * shortened polling interval takes effect within one slice (see the field + * javadoc), and {@link Config#lance_index_job_dispatcher_paused} suspends only + * the dispatch phase (see {@link #dispatchPendingJobs}). + */ +public class LanceIndexJobDispatcher extends MasterDaemon { + private static final Logger LOG = LogManager.getLogger(LanceIndexJobDispatcher.class); + + /** + * Upper bound of one sleep slice, equal to the shipped default interval. The + * daemon never sleeps longer than this, so a shortened + * {@link Config#lance_index_job_dispatch_interval_second} takes effect within + * one slice instead of waiting out a previously adopted long sleep: the + * elapsed check in {@link #runAfterCatalogReady} is re-evaluated against the + * current config at every wake. Slices bound only the sleep; rounds still + * honor the configured interval, because a wake whose configured interval + * (longer than this bound) has not elapsed since the last round skips the + * round. A lengthened interval takes effect at the next wake through the same + * check, and an interval at or below this bound needs no check at all — every + * wake runs a round, exactly one per configured period. + */ + private static final long MAX_SLEEP_SLICE_MS = 10_000L; + + private final Supplier<LanceIndexJobManager> jobManagerSupplier; + + /** Wall time of the last executed round, or -1 before the first one. */ + private long lastRoundMs = -1L; + + public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) { + this(() -> jobManager); + } + + public LanceIndexJobDispatcher(Supplier<LanceIndexJobManager> jobManagerSupplier) { + super("lance index job dispatcher", dispatchIntervalMs()); + this.jobManagerSupplier = jobManagerSupplier; + } + + /** + * Values loaded from fe.conf bypass the config validator (only ADMIN SET runs + * it), so the positive invariant is re-asserted where a non-positive value + * would break the loop: a non-positive interval would kill this thread inside + * {@code Thread.sleep} or busy-spin it, a non-positive deadline would sweep + * every dispatched job UNKNOWN on the next round, and a zero cap would stall + * dispatch forever. The refresh retry interval needs no such defense: a + * non-positive value simply disengages the throttle. + */ + private static long dispatchIntervalMs() { + return Math.max(1, Config.lance_index_job_dispatch_interval_second) * 1000L; + } + + private static long executeDeadlineMs(long nowMs) { + long second = Math.max(1L, Config.lance_index_job_execute_deadline_second); + return second > (Long.MAX_VALUE - nowMs) / 1000L ? Long.MAX_VALUE : nowMs + second * 1000L; + } + + @Override + protected void runAfterCatalogReady() { + if (!Env.getCurrentEnv().isMaster()) { + return; + } + if (Env.isCheckpointThread()) { + return; + } + long configuredMs = dispatchIntervalMs(); + setInterval(Math.min(configuredMs, MAX_SLEEP_SLICE_MS)); + if (configuredMs > MAX_SLEEP_SLICE_MS && lastRoundMs >= 0 && nowMs() - lastRoundMs < configuredMs) { + // A wake inside a long configured interval: the slice elapsed, the + // round period has not. Skipping is cheap and writes no journal record. + return; + } + lastRoundMs = nowMs(); + try { + runOneRound(jobManagerSupplier.get()); + } catch (Throwable t) { + LOG.warn("Failed to process one round of the lance index job dispatcher", t); + } + } + + /** Clock seam for the round-period check; tests advance it instead of sleeping. */ + protected long nowMs() { + return System.currentTimeMillis(); + } + + private void runOneRound(LanceIndexJobManager jobManager) { + long nowMs = System.currentTimeMillis(); + sweepExpiredRunningJobs(jobManager, nowMs); + sweepReplacedProcessEpochs(jobManager); + driveRequiredRefreshes(jobManager, nowMs); + dispatchPendingJobs(jobManager); + } + + /** + * Deadline sweep. A RUNNING job past its wait deadline has produced no + * complete trusted result, so it converges to UNKNOWN through the same + * completeWithResult channel a callback would use. Expiry bounds the wait + * only: it never proves termination, so the possible-live slot, the + * same-name fence, and the unresolved quota all stay held. + */ + private void sweepExpiredRunningJobs(LanceIndexJobManager jobManager, long nowMs) { + for (LanceIndexJob job : jobManager.getExpiredRunningJobs(nowMs)) { + try { + boolean completed = jobManager.completeWithResult(job.getJobId(), + dispatchRevisionOf(job), job.getInvocationId(), job.getBeProcessEpoch(), + new LanceIndexJobResult(LanceIndexJobResultCode.NO_TRUSTED_RESULT, + LanceIndexJobCompletionReason.NONE, + "execute deadline expired without a complete trusted result", false)); + if (completed) { + LOG.info("lance index job {} converged RUNNING -> UNKNOWN on deadline expiry", + job.getJobId()); + } else { + LOG.warn("deadline sweep skipped lance index job {}: already converged by a callback or sweep", + job.getJobId()); + } + } catch (Throwable t) { + LOG.warn("failed to sweep expired lance index job " + job.getJobId(), t); + } + } + } + + /** + * Possible-live sweep. The only slot-release proof this daemon produces is + * that the recorded backend process epoch no longer exists: a backend entry + * reporting a different epoch proves the process that received the dispatch + * was replaced. A missing backend entry or heartbeat loss proves nothing + * (the worker may still be running behind a partition), so such a job keeps + * its slot until a stronger proof or an operator force release. An epoch + * change also proves nothing about the outcome, so the mutation state is + * never touched here. + */ + private void sweepReplacedProcessEpochs(LanceIndexJobManager jobManager) { + for (LanceIndexJob job : jobManager.getJobsHoldingPossibleLiveSlot()) { + try { + Backend backend = Env.getCurrentSystemInfo().getBackend(job.getBackendId()); + if (backend == null || backend.getProcessEpoch() == job.getBeProcessEpoch()) { + continue; + } + boolean recorded = jobManager.recordTerminationProof(job.getJobId(), + dispatchRevisionOf(job), job.getBackendId(), job.getBeProcessEpoch(), + job.getInvocationId(), LanceIndexTerminationProof.BE_PROCESS_EPOCH_GONE); + if (recorded) { + LOG.info("released possible-live slot of lance index job {}: backend process epoch was replaced", + job.getJobId()); + } else { + LOG.warn("epoch sweep skipped lance index job {}: dispatch identity already moved", + job.getJobId()); + } + } catch (Throwable t) { + LOG.warn("failed to sweep possible-live slot of lance index job " + job.getJobId(), t); + } + } + } + + /** + * The clock the FAILED-refresh throttle measures from: the dedicated + * refresh-failure timestamp, falling back to the generic update time only for + * a legacy record replayed before the field existed. Measuring from the + * generic time would let an unrelated transition — a CHILD_REAPED proof or an + * epoch-gone release bumping it — postpone the next retry by a full interval + * while the fence and quota stay held. + */ + private static long refreshThrottledSinceMs(LanceIndexJob job) { + return job.getRefreshFailureTimeMs() == null ? job.getUpdateTimeMs() : job.getRefreshFailureTimeMs(); + } + + /** + * Refresh driver for terminal jobs with an unfinished refresh obligation. + * Completing the refresh is the protocol duty that releases the same-name + * fence and the unresolved quota once DONE; it is not a read-visibility + * action, because index metadata is never cached. Each job is driven + * through markRefreshRunning, the idempotent external-table refresh, then + * DONE or FAILED: a FAILED job keeps its fence and is retried, throttled to + * one attempt per retry interval, while a first REQUIRED refresh is never + * delayed. UNKNOWN jobs never appear here; they owe no refresh. + */ + private void driveRequiredRefreshes(LanceIndexJobManager jobManager, long nowMs) { + for (LanceIndexJob job : jobManager.getJobsNeedingRefresh()) { + try { + if (job.getRefreshState() == LanceIndexJobRefreshState.RUNNING) { + // In flight elsewhere; the master-transfer sweep downgrades a stale + // RUNNING back to REQUIRED, so a lost driver cannot strand it. + continue; + } + if (job.getRefreshState() == LanceIndexJobRefreshState.FAILED + && nowMs - refreshThrottledSinceMs(job) + < Config.lance_index_job_refresh_retry_second * 1000L) { + continue; + } + if (!jobManager.markRefreshRunning(job.getJobId(), job.getRevision())) { + // A concurrent driver won the compare-and-set; nothing to do here. + continue; + } + driveOneRefresh(jobManager, job); + } catch (Throwable t) { + LOG.warn("failed to drive the refresh of lance index job " + job.getJobId(), t); + } + } + } + + private void driveOneRefresh(LanceIndexJobManager jobManager, LanceIndexJob job) { + long refreshRevision = job.getRevision() + 1; + CatalogIf catalog = Env.getCurrentEnv().getCatalogMgr().getCatalog(job.getCatalogId()); + if (catalog == null) { + // Unreachable while the unresolved-job guard blocks catalog drops; kept as a + // fail-closed fallback so the job still transitions and retries later. + LOG.warn("catalog of lance index job {} is gone; marking its refresh FAILED", job.getJobId()); + finishRefreshTransition(jobManager, job.getJobId(), refreshRevision, false); + return; + } + // DONE is reserved for a refresh that verifiably did its work: a null db/table + // lookup is NOT that evidence (a cold cache or a transient remote failure also + // yields null and handleRefreshTable would silently skip the invalidation and + // the journal), so the local resolution decides how completion is proven. + DatabaseIf<? extends TableIf> db; + TableIf table; + try { + db = catalog.getDbNullable(job.getDbName()); + table = db == null ? null : db.getTableNullable(job.getTableName()); + } catch (Throwable t) { + LOG.warn("target of lance index job {} could not be resolved for its refresh; retrying later", + job.getJobId(), t); + finishRefreshTransition(jobManager, job.getJobId(), refreshRevision, false); + return; + } + if (db == null || table == null) { + // Positively verify the half-orphan through the namespace before DONE: only + // a VERIFIED_ABSENT answer is "nothing is left to invalidate". PRESENT with + // a cold local cache, or an UNRESOLVED check, marks the refresh FAILED and + // retries — never DONE without evidence. + if (catalog instanceof LanceExternalCatalog && ((LanceExternalCatalog) catalog).checkIndexJobDataset( + job.getDbName(), job.getTableName()).outcome + == LanceIndexDatasetCheck.Outcome.VERIFIED_ABSENT) { + LOG.info("refresh of lance index job {} skipped: the target is verified absent", job.getJobId()); + finishRefreshTransition(jobManager, job.getJobId(), refreshRevision, true); + return; + } + LOG.warn("refresh of lance index job {} cannot be verified this round; retrying later", job.getJobId()); + finishRefreshTransition(jobManager, job.getJobId(), refreshRevision, false); + return; + } + try { + // The target resolves: this refresh actually invalidates it. It is addressed + // by the persisted catalog id, never by the mutable name: a rename that + // hands this catalog's old name to a different catalog between the two + // resolutions must not refresh that one and bill the outcome to this job. + Env.getCurrentEnv().getRefreshManager().handleRefreshTable(job.getCatalogId(), + job.getDbName(), job.getTableName(), true); + } catch (Throwable t) { + // The typed DdlException is the expected failure; an unchecked exception out + // of the metadata path must still leave the durable refresh state, or the + // job would strand in refresh RUNNING until the next master transfer. + LOG.warn("refresh of lance index job {} failed; keeping the fence for a retry", + job.getJobId(), t); + finishRefreshTransition(jobManager, job.getJobId(), refreshRevision, false); + return; + } + finishRefreshTransition(jobManager, job.getJobId(), refreshRevision, true); + } + + /** + * Applies the DONE/FAILED transition with a bounded revision retry. A concurrent + * termination-proof write can bump the revision after markRefreshRunning succeeded, + * and silently losing that compare-and-set would leave the refresh RUNNING — a + * state only the master-transfer sweep downgrades. Re-reading the revision and + * retrying a few times converges it; a persistent loss is escalated. + */ + private void finishRefreshTransition(LanceIndexJobManager jobManager, long jobId, long expectedRevision, + boolean done) { + long revision = expectedRevision; + for (int attempt = 0; attempt < 3; attempt++) { + boolean transitioned = done ? jobManager.markRefreshDone(jobId, revision) + : jobManager.markRefreshFailed(jobId, revision); + if (transitioned) { + return; + } + LanceIndexJob fresh = jobManager.getJob(jobId); + if (fresh == null) { + break; + } + revision = fresh.getRevision(); + } + LOG.error("lance index job {} kept its refresh RUNNING: the DONE/FAILED transition kept losing the" + + " compare-and-set; the master-transfer sweep will downgrade it", jobId); + } + + /** + * PENDING dispatch. Makes at most + * {@link Config#lance_index_job_max_dispatch_per_round} fresh dispatches per + * round, and only a job this round actually made RUNNING consumes that + * budget: skipped jobs (an eligibility gate is closed, or every backend is + * at capacity) are scanned past, so a stable subset of permanently + * undispatchable jobs can never crowd out later ids. Per backend it never + * exceeds {@link Config#lance_index_job_max_inflight_per_backend} + * possible-live worker slots, counted from slot ownership (see + * {@link LanceIndexJobManager#countPossibleLiveSlotsByBackend()}) plus the + * jobs this round already made RUNNING. A job that cannot be dispatched + * keeps waiting as PENDING: there is no dispatch-exhaustion terminal state + * and no backoff beyond the daemon period. + * + * <p>{@link Config#lance_index_job_dispatcher_paused} suspends this phase + * only — the sweeps and the refresh driver keep running while it is set. + * The switch is checked at the phase entry and again before every single + * job attempt, which closes the admission race a test or operator cares + * about: anyone who sets the switch <em>before</em> admitting a job is + * guaranteed the job is never dispatched while paused. A round whose + * snapshot was taken before the admission never sees the job at all, and + * any round that can see it performs its per-job check after the + * admission, hence after the switch was set, and skips it. A skipped job + * never consumes the round's dispatch budget. + * + * <p>This daemon thread is also the only thread running the deadline sweep, + * the epoch sweep, and the refresh driver, so the blocking time of the + * dispatch attempts is bounded per round: at most one backend RPC timeout + * of it may be spent before the remaining jobs are deferred to the next + * round. A healthy round spends microseconds per attempt and never notices + * the bound; a backend whose job RPC stalls while still heartbeating can + * hold this thread for at most one in-flight timeout beyond the budget, + * instead of the per-round cap times the timeout (minutes with the + * defaults) — so the lifecycle work of the following rounds keeps its + * cadence. + */ + private void dispatchPendingJobs(LanceIndexJobManager jobManager) { + if (Config.lance_index_job_dispatcher_paused) { + return; + } + int maxPerRound = Math.max(1, Config.lance_index_job_max_dispatch_per_round); + Map<Long, Integer> inflightByBackend = jobManager.countPossibleLiveSlotsByBackend(); + // fe.conf bypasses the validator, so a non-positive timeout is clamped to + // keep at least one attempt per round. + long blockingBudgetMs = Math.max(1L, Config.backend_rpc_timeout_ms); + int dispatched = 0; + for (LanceIndexJob job : jobManager.getJobsNeedingDispatch()) { + if (Config.lance_index_job_dispatcher_paused) { + // Flipped mid-round: stop without touching the budget. + break; + } + if (dispatched >= maxPerRound) { + break; + } + if (blockingBudgetMs <= 0) { + LOG.info("lance index job dispatcher spent this round's blocking-dispatch budget;" + + " deferring lance index job {} to a later round", job.getJobId()); + break; + } + long attemptStartMs = nowMs(); + try { + if (tryDispatch(jobManager, job, inflightByBackend)) { + dispatched++; + } + } catch (Throwable t) { + LOG.warn("failed to dispatch lance index job " + job.getJobId(), t); + } + blockingBudgetMs -= nowMs() - attemptStartMs; + } + } + + /** + * One dispatch attempt for one PENDING job; returns true only when the + * attempt made the job durable RUNNING (and so consumes this round's + * dispatch budget). Every early return before markRunning leaves the job + * PENDING for a later round: the eligibility gates, the backend and + * capacity checks, and also the whole request preparation — storage-option + * resolution and the wire request build run before the durable boundary, + * so an FE-side failure there (for example a catalog id that resolves to + * nothing while ALTER CATALOG RENAME has the catalog temporarily removed) + * just retries next round instead of stranding the job UNKNOWN without a + * single byte sent. Once markRunning succeeds the job is durable RUNNING + * and this invocation id gets exactly one send attempt; after that only a + * matching callback, the deadline sweep, or the epoch sweep can converge + * the job. + */ + private boolean tryDispatch(LanceIndexJobManager jobManager, LanceIndexJob job, + Map<Long, Integer> inflightByBackend) { + boolean localDataset = isLocalFileDataset(job.getNormalizedLocator()); + if (localDataset && !Config.enable_lance_index_local_file_mutation) { + // Operator assertion is off: a local-filesystem mutation stays PENDING. + return false; + } + if (localDataset && Env.getCurrentEnv().getFrontends(null).size() != 1) { + // Local files are only shared by a single-node deployment. + return false; + } + SystemInfoService systemInfo = Env.getCurrentSystemInfo(); + // Every schedule-available worker, shuffled by the selection policy: the first + // one with a free possible-live slot takes the job, so a full backend defers + // this attempt only when every selectable backend is at the cap, never just + // because the randomly picked one is. allowOnSameHost keeps co-located + // backends visible (the policy default hides all but one per host), and + // preferComputeNode lets compute-only clusters serve Lance dispatch at all — + // the policy default filters every compute-role backend out. + List<Long> backendIds = systemInfo.selectBackendIdsByPolicy( + new BeSelectionPolicy.Builder().needScheduleAvailable().allowOnSameHost() + .preferComputeNode(true).build(), -1); Review Comment: done, selection now requests every eligible backend (assignExpectBeNum covers all registered BEs) with compute nodes first and mix nodes filling behind them, so an idle mix peer dispatches when the compute slots are full (ed204c215) ########## fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceExternalCatalog.java: ########## @@ -100,10 +107,41 @@ public String resolveCurrentIndexJobLocator(String dbName, String tableName) { try { return withClient(current -> current.resolveCurrentIndexJobLocator(dbName, tableName)); } catch (Exception e) { + LOG.warn("failed to resolve the current dataset locator of {}.{} in lance catalog {}", + dbName, tableName, getName(), e); return null; } } + /** + * Three-valued resolution of the dataset the given names point at, for callers + * that take a durable action on the verdict and must not fold "verified gone" + * and "could not tell" together. {@link #resolveCurrentIndexJobLocator} returns + * null for both on purpose (SHOW's fail-closed rule folds them); this form + * distinguishes them: a {@code TableNotFound}/{@code NamespaceNotFound} answer + * from the namespace is positive evidence of absence, while any other failure is + * logged (the sanitized client chain already masks locators and credentials) + * and reported as {@link LanceIndexDatasetCheck.Outcome#UNRESOLVED}. + * + * <p>Callers must pass the REMOTE names of the relations when they hold the + * resolved objects (admission persists local names, and the namespace is + * case-sensitive); a caller with no resolved object falls back to its local + * names, which is correct whenever the local layer could not resolve them + * either — a case-mapped remote name resolves locally, so a local miss with a + * reachable namespace means no case-insensitive match exists at all. + */ + public LanceIndexDatasetCheck checkIndexJobDataset(String dbName, String tableName) { + try { + String locator = withClient(current -> current.resolveCurrentIndexJobLocator(dbName, tableName)); + return LanceIndexDatasetCheck.present(locator); + } catch (TableNotFoundException | NamespaceNotFoundException e) { Review Comment: done, absence is now proven only through full case-insensitive namespace listings (the remote db list, then that db's table list), so a local-case spelling can no longer false-prove a cased dataset absent; any listing failure stays UNRESOLVED and retries (7b3e3426c) -- 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]
