github-actions[bot] commented on code in PR #67978:
URL: https://github.com/apache/doris/pull/67978#discussion_r4131443651
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobManager.java:
##########
@@ -693,6 +721,75 @@ && hasFenceIdentity(job)
}
}
+ /**
+ * PENDING jobs eligible for the dispatcher, in job id order (FIFO
fairness),
+ * at most {@code limit} of them. Target identity must be complete: corrupt
+ * identity-less records are never dispatchable (their only exit is the
+ * force-release transition by job id) and are skipped here. The dispatch
+ * quad is deliberately not required: it is written by {@code markRunning},
+ * which is the step this query feeds. All matches are collected and
ordered
+ * before truncating, so a stable subset of permanently undispatchable jobs
+ * can never crowd out later ids.
+ */
+ public List<LanceIndexJob> getJobsNeedingDispatch(int limit) {
+ readLock();
+ try {
+ List<LanceIndexJob> result = new ArrayList<>();
+ for (LanceIndexJob job : jobs.values()) {
+ if (job != null && job.getMutationState() ==
LanceIndexJobMutationState.PENDING
+ && hasDispatchTarget(job)) {
+ result.add(new LanceIndexJob(job));
+ }
Review Comment:
[P1] Scan beyond ineligible pending jobs before applying the round cap. With
the default cap of 16, 16 admitted file jobs stay PENDING while local mutation
is disabled; this subList always returns those same IDs, so a later S3 job is
never visited. Count actual dispatches or advance a fair cursor after skipped
jobs.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobReportHandler.java:
##########
@@ -0,0 +1,129 @@
+// 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.thrift.TLanceIndexJobReport;
+import org.apache.doris.thrift.TLanceIndexTerminationProof;
+
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.Objects;
+
+/**
+ * Applies one typed result envelope reported by a backend to the durable job
+ * record. This is a thin shim over the manager transitions: dispatch-identity
+ * checking and result classification all live in {@link LanceIndexJobManager},
+ * so a stale or identity-mismatched report only logs a warning and changes
+ * nothing. A malformed envelope (missing result code, a code this FE does not
+ * know, or a sanitized message past the durable bound) is dropped rather than
+ * trusted; the job then converges through the dispatcher's deadline sweep.
+ * Only the typed codes are read: message text is never inspected to infer an
+ * outcome.
+ *
+ * <p>The handler runs on the report RPC thread and performs no I/O beyond the
+ * manager's own edit-log write. It starts no refresh: the metadata refresh a
+ * completed job may owe is driven by the dispatcher daemon, not here.
+ */
+public class LanceIndexJobReportHandler {
+
+ private static final Logger LOG =
LogManager.getLogger(LanceIndexJobReportHandler.class);
+
+ private final LanceIndexJobManager jobManager;
+
+ public LanceIndexJobReportHandler(LanceIndexJobManager jobManager) {
+ this.jobManager = Objects.requireNonNull(jobManager, "jobManager");
+ }
+
+ /**
+ * Handles one report: a matched report completes the job with its
+ * classified result, and a CHILD_REAPED termination proof additionally
+ * releases the possible-live slot, because reaping the exact child process
+ * proves that process ended (which still says nothing about the outcome).
+ */
+ public void handle(TLanceIndexJobReport report) {
+ if (report == null) {
+ LOG.warn("dropping null lance index job report");
+ return;
+ }
+ LanceIndexJobResult result;
+ try {
+ result = toResult(report);
+ } catch (IllegalArgumentException e) {
Review Comment:
[P2] Process an identity-matched CHILD_REAPED proof even when the result is
malformed. If a reaped child reports an overlong message or unknown result
code, toResult throws and this return drops its valid termination proof. The
deadline then leaves a false possible-live slot on a stable BE;
recordTerminationProof can validate the proof independently while the result
remains untrusted.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,522 @@
+// 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.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.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 com.google.common.collect.Maps;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+
+/**
+ * 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 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.
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+ private static final Logger LOG =
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+ private final LanceIndexJobManager jobManager;
+
+ public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+ super("lance index job dispatcher", dispatchIntervalMs());
+ this.jobManager = jobManager;
+ }
+
+ /**
+ * 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;
+ }
+ setInterval(dispatchIntervalMs());
+ try {
+ runOneRound();
+ } catch (Throwable t) {
+ LOG.warn("Failed to process one round of the lance index job
dispatcher", t);
+ }
+ }
+
+ private void runOneRound() {
+ long nowMs = System.currentTimeMillis();
+ sweepExpiredRunningJobs(nowMs);
+ sweepReplacedProcessEpochs();
+ driveRequiredRefreshes(nowMs);
+ dispatchPendingJobs();
+ }
+
+ /**
+ * 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(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() {
+ 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);
+ }
+ }
+ }
+
+ /**
+ * 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(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 - job.getUpdateTimeMs()
+ < 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(job);
+ } catch (Throwable t) {
+ LOG.warn("failed to drive the refresh of lance index job " +
job.getJobId(), t);
+ }
+ }
+ }
+
+ private void driveOneRefresh(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(job.getJobId(), refreshRevision, false);
+ return;
+ }
+ try {
+ // A half-orphan target (its db or table already dropped
externally) is a
+ // silent no-op: nothing is left to invalidate, and DONE is the
correct end
+ // state for the job.
+
Env.getCurrentEnv().getRefreshManager().handleRefreshTable(catalog.getName(),
+ 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(job.getJobId(), refreshRevision, false);
+ return;
+ }
+ finishRefreshTransition(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(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. Attempts at most
+ * {@link Config#lance_index_job_max_dispatch_per_round} fresh dispatches
per
+ * round, and never more than {@link
Config#lance_index_job_max_inflight_per_backend}
+ * in-flight jobs per backend, counted from the RUNNING snapshot 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.
+ */
+ private void dispatchPendingJobs() {
+ int maxPerRound = Math.max(1,
Config.lance_index_job_max_dispatch_per_round);
+ Map<Long, Integer> inflightByBackend = countInflightByBackend();
+ int attempted = 0;
+ for (LanceIndexJob job :
jobManager.getJobsNeedingDispatch(maxPerRound)) {
+ if (++attempted > maxPerRound) {
+ break;
+ }
+ try {
+ tryDispatch(job, inflightByBackend);
+ } catch (Throwable t) {
+ LOG.warn("failed to dispatch lance index job " +
job.getJobId(), t);
+ }
+ }
+ }
+
+ private Map<Long, Integer> countInflightByBackend() {
+ Map<Long, Integer> inflightByBackend = Maps.newHashMap();
+ for (LanceIndexJob job : jobManager.getAllJobsSnapshot()) {
+ if (job.getMutationState() == LanceIndexJobMutationState.RUNNING
&& job.getBackendId() != null) {
+ inflightByBackend.merge(job.getBackendId(), 1, Integer::sum);
+ }
+ }
+ return inflightByBackend;
+ }
+
+ /**
+ * One dispatch attempt for one PENDING job. Every early return before
+ * markRunning leaves the job PENDING for a later round. 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 void tryDispatch(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;
+ }
+ if (localDataset && Env.getCurrentEnv().getFrontends(null).size() !=
1) {
+ // Local files are only shared by a single-node deployment.
+ return;
+ }
+ SystemInfoService systemInfo = Env.getCurrentSystemInfo();
+ List<Long> backendIds = systemInfo.selectBackendIdsByPolicy(
+ new
BeSelectionPolicy.Builder().needScheduleAvailable().build(), 1);
+ if (backendIds.isEmpty()) {
+ return;
+ }
+ Backend backend = systemInfo.getBackend(backendIds.get(0));
+ if (backend == null) {
+ return;
+ }
+ if (localDataset && !isOnlyAliveBackend(systemInfo, backend.getId())) {
+ return;
+ }
+ Integer inflight = inflightByBackend.get(backend.getId());
+ if (inflight != null && inflight >= Math.max(1,
Config.lance_index_job_max_inflight_per_backend)) {
+ return;
+ }
+ String invocationId = UUID.randomUUID().toString();
+ // The process epoch is captured once, and the same value goes to the
+ // durable record and the wire: a heartbeat landing between the two
reads
+ // must not split the dispatch identity (the callback matches the
durable
+ // value, and the epoch sweep releases the slot against it).
+ long beProcessEpoch = backend.getProcessEpoch();
+ long deadlineMs = executeDeadlineMs(System.currentTimeMillis());
+ if (!jobManager.markRunning(job.getJobId(), job.getRevision(),
backend.getId(),
+ beProcessEpoch, invocationId, deadlineMs)) {
+ // The compare-and-set lost: this attempt's dispatch identity is
void and its
+ // invocation id is discarded. A fresh identity is built from
scratch next round.
+ return;
+ }
+ inflightByBackend.merge(backend.getId(), 1, Integer::sum);
+ long expectedDispatchRevision = job.getRevision() + 1;
+ LanceIndexJob fresh = jobManager.getJob(job.getJobId());
+ if (!Env.getCurrentEnv().isMaster() || fresh == null
+ || fresh.getMutationState() !=
LanceIndexJobMutationState.RUNNING
+ || fresh.getDispatchRevision() == null
+ || fresh.getDispatchRevision() != expectedDispatchRevision
+ || !invocationId.equals(fresh.getInvocationId())) {
+ // The recheck failed right before the send: no send, and no
resend either.
+ // The job is durable RUNNING, so the deadline sweep or a matching
callback
+ // converges it.
+ LOG.warn("lance index job {} did not survive the pre-send recheck;
not sending", job.getJobId());
+ return;
+ }
+ TLanceIndexJobDispatch dispatch;
+ try {
+ dispatch = buildDispatch(fresh, invocationId, deadlineMs,
beProcessEpoch,
+ resolveStorageOptions(fresh));
+ } catch (Exception e) {
+ // An FE-side resolution failure is not a trusted worker
rejection, so it must
+ // not fabricate NOT_COMMITTED. The job is already RUNNING without
a send, and
+ // the send may never happen, so converge it to UNKNOWN
fail-closed.
+ LOG.warn("failed to prepare the dispatch of lance index job {}:
{}", job.getJobId(), e.getMessage());
+ completeNoTrusted(fresh, "dispatch preparation failed before
send");
+ return;
+ }
+ TStatus status;
+ try {
+ status = sendExecuteRequest(backend, dispatch);
+ } catch (Exception e) {
+ // The request may have reached the backend, so its outcome cannot
be trusted.
+ LOG.warn("dispatch send of lance index job {} failed: {}",
job.getJobId(), e.getMessage());
+ completeNoTrusted(fresh, "dispatch send failed; the result cannot
be trusted");
+ return;
+ }
+ if (status == null || status.getStatusCode() == null) {
+ // Absence of a status is the absence of a trusted answer, not a
clean
+ // rejection; only a complete error status proves the dispatch was
not
+ // enqueued.
+ LOG.warn("dispatch send of lance index job {} returned no status",
job.getJobId());
+ completeNoTrusted(fresh, "dispatch send returned no status");
+ return;
+ }
+ if (status.getStatusCode() != TStatusCode.OK) {
Review Comment:
[P2] Release the possible-live slot on a clean pre-enqueue rejection. The
shipped BE handler always returns NOT_IMPLEMENTED_ERROR without enqueuing, yet
completePreInvocationRejected only changes the mutation state:
possibleLiveOwned remains true and SHOW reports a worker as possibly live until
that BE restarts. Record a durable no-enqueue proof for the matched dispatch
identity.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,522 @@
+// 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.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.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 com.google.common.collect.Maps;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+
+/**
+ * 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 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.
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+ private static final Logger LOG =
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+ private final LanceIndexJobManager jobManager;
+
+ public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+ super("lance index job dispatcher", dispatchIntervalMs());
+ this.jobManager = jobManager;
+ }
+
+ /**
+ * 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;
+ }
+ setInterval(dispatchIntervalMs());
+ try {
+ runOneRound();
+ } catch (Throwable t) {
+ LOG.warn("Failed to process one round of the lance index job
dispatcher", t);
+ }
+ }
+
+ private void runOneRound() {
+ long nowMs = System.currentTimeMillis();
+ sweepExpiredRunningJobs(nowMs);
+ sweepReplacedProcessEpochs();
+ driveRequiredRefreshes(nowMs);
+ dispatchPendingJobs();
+ }
+
+ /**
+ * 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(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() {
+ 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);
+ }
+ }
+ }
+
+ /**
+ * 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(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 - job.getUpdateTimeMs()
+ < 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(job);
+ } catch (Throwable t) {
+ LOG.warn("failed to drive the refresh of lance index job " +
job.getJobId(), t);
+ }
+ }
+ }
+
+ private void driveOneRefresh(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(job.getJobId(), refreshRevision, false);
+ return;
+ }
+ try {
+ // A half-orphan target (its db or table already dropped
externally) is a
+ // silent no-op: nothing is left to invalidate, and DONE is the
correct end
+ // state for the job.
+
Env.getCurrentEnv().getRefreshManager().handleRefreshTable(catalog.getName(),
+ 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(job.getJobId(), refreshRevision, false);
+ return;
+ }
+ finishRefreshTransition(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(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. Attempts at most
+ * {@link Config#lance_index_job_max_dispatch_per_round} fresh dispatches
per
+ * round, and never more than {@link
Config#lance_index_job_max_inflight_per_backend}
+ * in-flight jobs per backend, counted from the RUNNING snapshot 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.
+ */
+ private void dispatchPendingJobs() {
+ int maxPerRound = Math.max(1,
Config.lance_index_job_max_dispatch_per_round);
+ Map<Long, Integer> inflightByBackend = countInflightByBackend();
+ int attempted = 0;
+ for (LanceIndexJob job :
jobManager.getJobsNeedingDispatch(maxPerRound)) {
+ if (++attempted > maxPerRound) {
+ break;
+ }
+ try {
+ tryDispatch(job, inflightByBackend);
+ } catch (Throwable t) {
+ LOG.warn("failed to dispatch lance index job " +
job.getJobId(), t);
+ }
+ }
+ }
+
+ private Map<Long, Integer> countInflightByBackend() {
+ Map<Long, Integer> inflightByBackend = Maps.newHashMap();
+ for (LanceIndexJob job : jobManager.getAllJobsSnapshot()) {
+ if (job.getMutationState() == LanceIndexJobMutationState.RUNNING
&& job.getBackendId() != null) {
+ inflightByBackend.merge(job.getBackendId(), 1, Integer::sum);
+ }
+ }
+ return inflightByBackend;
+ }
+
+ /**
+ * One dispatch attempt for one PENDING job. Every early return before
+ * markRunning leaves the job PENDING for a later round. 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 void tryDispatch(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;
+ }
+ if (localDataset && Env.getCurrentEnv().getFrontends(null).size() !=
1) {
+ // Local files are only shared by a single-node deployment.
+ return;
+ }
+ SystemInfoService systemInfo = Env.getCurrentSystemInfo();
+ List<Long> backendIds = systemInfo.selectBackendIdsByPolicy(
+ new
BeSelectionPolicy.Builder().needScheduleAvailable().build(), 1);
+ if (backendIds.isEmpty()) {
+ return;
+ }
+ Backend backend = systemInfo.getBackend(backendIds.get(0));
+ if (backend == null) {
+ return;
+ }
+ if (localDataset && !isOnlyAliveBackend(systemInfo, backend.getId())) {
+ return;
+ }
+ Integer inflight = inflightByBackend.get(backend.getId());
+ if (inflight != null && inflight >= Math.max(1,
Config.lance_index_job_max_inflight_per_backend)) {
+ return;
+ }
+ String invocationId = UUID.randomUUID().toString();
+ // The process epoch is captured once, and the same value goes to the
+ // durable record and the wire: a heartbeat landing between the two
reads
+ // must not split the dispatch identity (the callback matches the
durable
+ // value, and the epoch sweep releases the slot against it).
+ long beProcessEpoch = backend.getProcessEpoch();
+ long deadlineMs = executeDeadlineMs(System.currentTimeMillis());
+ if (!jobManager.markRunning(job.getJobId(), job.getRevision(),
backend.getId(),
+ beProcessEpoch, invocationId, deadlineMs)) {
+ // The compare-and-set lost: this attempt's dispatch identity is
void and its
+ // invocation id is discarded. A fresh identity is built from
scratch next round.
+ return;
+ }
+ inflightByBackend.merge(backend.getId(), 1, Integer::sum);
+ long expectedDispatchRevision = job.getRevision() + 1;
+ LanceIndexJob fresh = jobManager.getJob(job.getJobId());
+ if (!Env.getCurrentEnv().isMaster() || fresh == null
+ || fresh.getMutationState() !=
LanceIndexJobMutationState.RUNNING
+ || fresh.getDispatchRevision() == null
+ || fresh.getDispatchRevision() != expectedDispatchRevision
+ || !invocationId.equals(fresh.getInvocationId())) {
+ // The recheck failed right before the send: no send, and no
resend either.
+ // The job is durable RUNNING, so the deadline sweep or a matching
callback
+ // converges it.
+ LOG.warn("lance index job {} did not survive the pre-send recheck;
not sending", job.getJobId());
+ return;
+ }
+ TLanceIndexJobDispatch dispatch;
+ try {
+ dispatch = buildDispatch(fresh, invocationId, deadlineMs,
beProcessEpoch,
+ resolveStorageOptions(fresh));
+ } catch (Exception e) {
Review Comment:
[P1] Keep a failed pre-send preparation out of UNKNOWN. During ALTER CATALOG
RENAME, CatalogMgr temporarily removes the catalog ID; a job can be marked
RUNNING, then resolveStorageOptions throws before any RPC, and this catch
retains its fence and quota indefinitely as UNKNOWN. Prepare the request before
markRunning, or durably record a proven no-send result and release the slot for
this identity.
##########
fe/fe-core/src/main/java/org/apache/doris/catalog/Env.java:
##########
@@ -860,6 +863,7 @@ public Env(boolean isCheckpointCatalog) {
this.eventProcessor = new EventProcessor(mtmvService);
this.insertOverwriteManager = new InsertOverwriteManager();
Review Comment:
[P1] Rebind the dispatcher after loading the job-manager image. The
constructor passes its initial manager into a final dispatcher field, but
loadLanceIndexJobManager later replaces Env.lanceIndexJobManager with a new
object. After restart, replay, admission, SHOW, and the master-transfer sweep
use the restored manager while every dispatcher round scans the abandoned empty
one, so no restored or newly admitted job advances. Resolve Env's current
manager each round or construct the dispatcher after image load, and cover an
image restart in the wiring test.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,522 @@
+// 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.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.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 com.google.common.collect.Maps;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+
+/**
+ * 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 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.
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+ private static final Logger LOG =
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+ private final LanceIndexJobManager jobManager;
+
+ public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+ super("lance index job dispatcher", dispatchIntervalMs());
+ this.jobManager = jobManager;
+ }
+
+ /**
+ * 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;
+ }
+ setInterval(dispatchIntervalMs());
+ try {
+ runOneRound();
+ } catch (Throwable t) {
+ LOG.warn("Failed to process one round of the lance index job
dispatcher", t);
+ }
+ }
+
+ private void runOneRound() {
+ long nowMs = System.currentTimeMillis();
+ sweepExpiredRunningJobs(nowMs);
+ sweepReplacedProcessEpochs();
+ driveRequiredRefreshes(nowMs);
+ dispatchPendingJobs();
+ }
+
+ /**
+ * 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(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() {
+ 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);
+ }
+ }
+ }
+
+ /**
+ * 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(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 - job.getUpdateTimeMs()
+ < 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(job);
+ } catch (Throwable t) {
+ LOG.warn("failed to drive the refresh of lance index job " +
job.getJobId(), t);
+ }
+ }
+ }
+
+ private void driveOneRefresh(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(job.getJobId(), refreshRevision, false);
+ return;
+ }
+ try {
+ // A half-orphan target (its db or table already dropped
externally) is a
+ // silent no-op: nothing is left to invalidate, and DONE is the
correct end
+ // state for the job.
+
Env.getCurrentEnv().getRefreshManager().handleRefreshTable(catalog.getName(),
+ 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(job.getJobId(), refreshRevision, false);
+ return;
+ }
+ finishRefreshTransition(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(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. Attempts at most
+ * {@link Config#lance_index_job_max_dispatch_per_round} fresh dispatches
per
+ * round, and never more than {@link
Config#lance_index_job_max_inflight_per_backend}
+ * in-flight jobs per backend, counted from the RUNNING snapshot 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.
+ */
+ private void dispatchPendingJobs() {
+ int maxPerRound = Math.max(1,
Config.lance_index_job_max_dispatch_per_round);
+ Map<Long, Integer> inflightByBackend = countInflightByBackend();
+ int attempted = 0;
+ for (LanceIndexJob job :
jobManager.getJobsNeedingDispatch(maxPerRound)) {
+ if (++attempted > maxPerRound) {
+ break;
+ }
+ try {
+ tryDispatch(job, inflightByBackend);
+ } catch (Throwable t) {
+ LOG.warn("failed to dispatch lance index job " +
job.getJobId(), t);
+ }
+ }
+ }
+
+ private Map<Long, Integer> countInflightByBackend() {
+ Map<Long, Integer> inflightByBackend = Maps.newHashMap();
+ for (LanceIndexJob job : jobManager.getAllJobsSnapshot()) {
+ if (job.getMutationState() == LanceIndexJobMutationState.RUNNING
&& job.getBackendId() != null) {
+ inflightByBackend.merge(job.getBackendId(), 1, Integer::sum);
+ }
+ }
+ return inflightByBackend;
+ }
+
+ /**
+ * One dispatch attempt for one PENDING job. Every early return before
+ * markRunning leaves the job PENDING for a later round. 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 void tryDispatch(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;
+ }
+ if (localDataset && Env.getCurrentEnv().getFrontends(null).size() !=
1) {
+ // Local files are only shared by a single-node deployment.
+ return;
+ }
+ SystemInfoService systemInfo = Env.getCurrentSystemInfo();
+ List<Long> backendIds = systemInfo.selectBackendIdsByPolicy(
Review Comment:
[P2] Choose a backend with a free job slot before consuming this attempt.
selectBackendIdsByPolicy(..., 1) returns a random BE; when it is at the per-BE
cap, tryDispatch returns even if another schedule-available BE is idle. With
one full and one idle BE, half the jobs are deferred for a whole polling round
on average. Filter/select across all available backends.
##########
regression-test/suites/external_table_p0/lance/test_lance_index_admission.groovy:
##########
@@ -64,13 +65,28 @@ suite("test_lance_index_admission",
"p0,external,nonConcurrent") {
assertEquals(1, quotaRows.size())
String originalGate = gateRows[0][1].toString()
String originalQuota = quotaRows[0][1].toString()
+ def dispatchIntervalRows = master_sql """ADMIN SHOW FRONTEND CONFIG LIKE
'lance_index_job_dispatch_interval_second'"""
+ assertEquals(1, dispatchIntervalRows.size())
+ String originalDispatchInterval = dispatchIntervalRows[0][1].toString()
+ // The interval pin below outlives the daemon's in-flight wait on this
shipped
+ // default, so the wait is sized correctly.
+ assertEquals("10", originalDispatchInterval)
// The main scenario admits two jobs on one table, independently of the
// cluster's original quota. The dedicated quota case temporarily lowers
it.
String suiteQuota = Math.max(2L, originalQuota.toLong()).toString()
Throwable suiteFailure = null
try {
master_sql """ADMIN SET FRONTEND CONFIG
("lance_index_job_max_unresolved_per_table" = "${suiteQuota}")"""
+ // Pin the dispatcher's polling interval to one hour so no dispatch
round can
+ // fire between admission and the PENDING assertions below: this
slice's
+ // backends answer submit_lance_index_job with a clean not-implemented
error,
+ // which converges a dispatched job to NOT_COMMITTED and would break
this
+ // suite's PENDING premise. A cycle already sleeping on the shipped
interval
Review Comment:
[P2] Synchronize with a completed dispatcher cycle before asserting PENDING.
Setting the interval to 3600 and sleeping 11 seconds does not stop a round
already running: a slow refresh before dispatch can reach the newly admitted
job, or a round blocked in a backend RPC can finish and sleep its old 10-second
interval before the next round picks up the job. The shipped BE stub then moves
it to NOT_COMMITTED and makes these PENDING assertions flaky. Use a daemon
barrier or scoped pause in both regression suites.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,522 @@
+// 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.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.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 com.google.common.collect.Maps;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+
+/**
+ * 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 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.
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+ private static final Logger LOG =
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+ private final LanceIndexJobManager jobManager;
+
+ public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+ super("lance index job dispatcher", dispatchIntervalMs());
+ this.jobManager = jobManager;
+ }
+
+ /**
+ * 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;
+ }
+ setInterval(dispatchIntervalMs());
+ try {
+ runOneRound();
+ } catch (Throwable t) {
+ LOG.warn("Failed to process one round of the lance index job
dispatcher", t);
+ }
+ }
+
+ private void runOneRound() {
+ long nowMs = System.currentTimeMillis();
+ sweepExpiredRunningJobs(nowMs);
+ sweepReplacedProcessEpochs();
+ driveRequiredRefreshes(nowMs);
+ dispatchPendingJobs();
+ }
+
+ /**
+ * 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(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() {
+ 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);
+ }
+ }
+ }
+
+ /**
+ * 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(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 - job.getUpdateTimeMs()
+ < 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(job);
+ } catch (Throwable t) {
+ LOG.warn("failed to drive the refresh of lance index job " +
job.getJobId(), t);
+ }
+ }
+ }
+
+ private void driveOneRefresh(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(job.getJobId(), refreshRevision, false);
+ return;
+ }
+ try {
+ // A half-orphan target (its db or table already dropped
externally) is a
+ // silent no-op: nothing is left to invalidate, and DONE is the
correct end
+ // state for the job.
+
Env.getCurrentEnv().getRefreshManager().handleRefreshTable(catalog.getName(),
+ 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(job.getJobId(), refreshRevision, false);
+ return;
+ }
+ finishRefreshTransition(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(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. Attempts at most
+ * {@link Config#lance_index_job_max_dispatch_per_round} fresh dispatches
per
+ * round, and never more than {@link
Config#lance_index_job_max_inflight_per_backend}
+ * in-flight jobs per backend, counted from the RUNNING snapshot 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.
+ */
+ private void dispatchPendingJobs() {
+ int maxPerRound = Math.max(1,
Config.lance_index_job_max_dispatch_per_round);
+ Map<Long, Integer> inflightByBackend = countInflightByBackend();
+ int attempted = 0;
+ for (LanceIndexJob job :
jobManager.getJobsNeedingDispatch(maxPerRound)) {
+ if (++attempted > maxPerRound) {
+ break;
+ }
+ try {
+ tryDispatch(job, inflightByBackend);
+ } catch (Throwable t) {
+ LOG.warn("failed to dispatch lance index job " +
job.getJobId(), t);
+ }
+ }
+ }
+
+ private Map<Long, Integer> countInflightByBackend() {
+ Map<Long, Integer> inflightByBackend = Maps.newHashMap();
+ for (LanceIndexJob job : jobManager.getAllJobsSnapshot()) {
+ if (job.getMutationState() == LanceIndexJobMutationState.RUNNING
&& job.getBackendId() != null) {
Review Comment:
[P1] Base backend capacity on possible-live ownership rather than RUNNING
state. After an epoch change, the sweep proves old workers dead and clears
their slots but their jobs stay RUNNING until deadline, so they block the
restarted BE for up to an hour. Conversely deadline expiry makes a potentially
live child UNKNOWN, so it stops counting and another job can exceed the
intended worker cap. Both transitions run before this count in the same round;
clear proven no-enqueue slots as well.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,522 @@
+// 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.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.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 com.google.common.collect.Maps;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+
+/**
+ * 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 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.
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+ private static final Logger LOG =
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+ private final LanceIndexJobManager jobManager;
+
+ public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+ super("lance index job dispatcher", dispatchIntervalMs());
+ this.jobManager = jobManager;
+ }
+
+ /**
+ * 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;
+ }
+ setInterval(dispatchIntervalMs());
Review Comment:
[P2] Make a shorter mutable dispatch interval take effect while the daemon
sleeps. Both changed regression suites set this to 3600 seconds, wait for the
daemon to adopt it, then restore 10 seconds; Daemon.run is still in its
one-hour Thread.sleep, so subsequent jobs can sit PENDING for nearly an hour.
Use a bounded wake cadence or wake the daemon on config change.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,522 @@
+// 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.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.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 com.google.common.collect.Maps;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+
+/**
+ * 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 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.
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+ private static final Logger LOG =
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+ private final LanceIndexJobManager jobManager;
+
+ public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+ super("lance index job dispatcher", dispatchIntervalMs());
+ this.jobManager = jobManager;
+ }
+
+ /**
+ * 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;
+ }
+ setInterval(dispatchIntervalMs());
+ try {
+ runOneRound();
+ } catch (Throwable t) {
+ LOG.warn("Failed to process one round of the lance index job
dispatcher", t);
+ }
+ }
+
+ private void runOneRound() {
+ long nowMs = System.currentTimeMillis();
+ sweepExpiredRunningJobs(nowMs);
+ sweepReplacedProcessEpochs();
+ driveRequiredRefreshes(nowMs);
+ dispatchPendingJobs();
+ }
+
+ /**
+ * 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(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() {
+ 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);
+ }
+ }
+ }
+
+ /**
+ * 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(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 - job.getUpdateTimeMs()
+ < 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(job);
+ } catch (Throwable t) {
+ LOG.warn("failed to drive the refresh of lance index job " +
job.getJobId(), t);
+ }
+ }
+ }
+
+ private void driveOneRefresh(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(job.getJobId(), refreshRevision, false);
+ return;
+ }
+ try {
+ // A half-orphan target (its db or table already dropped
externally) is a
+ // silent no-op: nothing is left to invalidate, and DONE is the
correct end
+ // state for the job.
+
Env.getCurrentEnv().getRefreshManager().handleRefreshTable(catalog.getName(),
+ 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(job.getJobId(), refreshRevision, false);
+ return;
+ }
+ finishRefreshTransition(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(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. Attempts at most
+ * {@link Config#lance_index_job_max_dispatch_per_round} fresh dispatches
per
+ * round, and never more than {@link
Config#lance_index_job_max_inflight_per_backend}
+ * in-flight jobs per backend, counted from the RUNNING snapshot 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.
+ */
+ private void dispatchPendingJobs() {
+ int maxPerRound = Math.max(1,
Config.lance_index_job_max_dispatch_per_round);
+ Map<Long, Integer> inflightByBackend = countInflightByBackend();
+ int attempted = 0;
+ for (LanceIndexJob job :
jobManager.getJobsNeedingDispatch(maxPerRound)) {
+ if (++attempted > maxPerRound) {
+ break;
+ }
+ try {
+ tryDispatch(job, inflightByBackend);
+ } catch (Throwable t) {
+ LOG.warn("failed to dispatch lance index job " +
job.getJobId(), t);
+ }
+ }
+ }
+
+ private Map<Long, Integer> countInflightByBackend() {
+ Map<Long, Integer> inflightByBackend = Maps.newHashMap();
+ for (LanceIndexJob job : jobManager.getAllJobsSnapshot()) {
+ if (job.getMutationState() == LanceIndexJobMutationState.RUNNING
&& job.getBackendId() != null) {
+ inflightByBackend.merge(job.getBackendId(), 1, Integer::sum);
+ }
+ }
+ return inflightByBackend;
+ }
+
+ /**
+ * One dispatch attempt for one PENDING job. Every early return before
+ * markRunning leaves the job PENDING for a later round. 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 void tryDispatch(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;
+ }
+ if (localDataset && Env.getCurrentEnv().getFrontends(null).size() !=
1) {
+ // Local files are only shared by a single-node deployment.
+ return;
+ }
+ SystemInfoService systemInfo = Env.getCurrentSystemInfo();
+ List<Long> backendIds = systemInfo.selectBackendIdsByPolicy(
+ new
BeSelectionPolicy.Builder().needScheduleAvailable().build(), 1);
+ if (backendIds.isEmpty()) {
+ return;
+ }
+ Backend backend = systemInfo.getBackend(backendIds.get(0));
+ if (backend == null) {
+ return;
+ }
+ if (localDataset && !isOnlyAliveBackend(systemInfo, backend.getId())) {
+ return;
+ }
+ Integer inflight = inflightByBackend.get(backend.getId());
+ if (inflight != null && inflight >= Math.max(1,
Config.lance_index_job_max_inflight_per_backend)) {
+ return;
+ }
+ String invocationId = UUID.randomUUID().toString();
+ // The process epoch is captured once, and the same value goes to the
+ // durable record and the wire: a heartbeat landing between the two
reads
+ // must not split the dispatch identity (the callback matches the
durable
+ // value, and the epoch sweep releases the slot against it).
+ long beProcessEpoch = backend.getProcessEpoch();
+ long deadlineMs = executeDeadlineMs(System.currentTimeMillis());
+ if (!jobManager.markRunning(job.getJobId(), job.getRevision(),
backend.getId(),
+ beProcessEpoch, invocationId, deadlineMs)) {
+ // The compare-and-set lost: this attempt's dispatch identity is
void and its
+ // invocation id is discarded. A fresh identity is built from
scratch next round.
+ return;
+ }
+ inflightByBackend.merge(backend.getId(), 1, Integer::sum);
+ long expectedDispatchRevision = job.getRevision() + 1;
+ LanceIndexJob fresh = jobManager.getJob(job.getJobId());
+ if (!Env.getCurrentEnv().isMaster() || fresh == null
+ || fresh.getMutationState() !=
LanceIndexJobMutationState.RUNNING
+ || fresh.getDispatchRevision() == null
+ || fresh.getDispatchRevision() != expectedDispatchRevision
+ || !invocationId.equals(fresh.getInvocationId())) {
+ // The recheck failed right before the send: no send, and no
resend either.
+ // The job is durable RUNNING, so the deadline sweep or a matching
callback
+ // converges it.
+ LOG.warn("lance index job {} did not survive the pre-send recheck;
not sending", job.getJobId());
+ return;
+ }
+ TLanceIndexJobDispatch dispatch;
+ try {
+ dispatch = buildDispatch(fresh, invocationId, deadlineMs,
beProcessEpoch,
+ resolveStorageOptions(fresh));
+ } catch (Exception e) {
+ // An FE-side resolution failure is not a trusted worker
rejection, so it must
+ // not fabricate NOT_COMMITTED. The job is already RUNNING without
a send, and
+ // the send may never happen, so converge it to UNKNOWN
fail-closed.
+ LOG.warn("failed to prepare the dispatch of lance index job {}:
{}", job.getJobId(), e.getMessage());
+ completeNoTrusted(fresh, "dispatch preparation failed before
send");
+ return;
+ }
+ TStatus status;
+ try {
+ status = sendExecuteRequest(backend, dispatch);
+ } catch (Exception e) {
Review Comment:
[P1] Separate known no-enqueue failures from ambiguous RPC failures here. An
old BE returns Thrift UNKNOWN_METHOD for this new RPC during a rolling upgrade;
a client-pool borrow can also fail while opening the socket before
submitLanceIndexJob is called. Neither case can execute a worker, yet this
catch makes the job UNKNOWN and holds its fence and quota indefinitely. Gate
old BEs and classify proven pre-call failures separately, while retaining
UNKNOWN for failures after invocation starts.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,522 @@
+// 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.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.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 com.google.common.collect.Maps;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+
+/**
+ * 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 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.
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+ private static final Logger LOG =
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+ private final LanceIndexJobManager jobManager;
+
+ public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+ super("lance index job dispatcher", dispatchIntervalMs());
+ this.jobManager = jobManager;
+ }
+
+ /**
+ * 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;
+ }
+ setInterval(dispatchIntervalMs());
+ try {
+ runOneRound();
+ } catch (Throwable t) {
+ LOG.warn("Failed to process one round of the lance index job
dispatcher", t);
+ }
+ }
+
+ private void runOneRound() {
+ long nowMs = System.currentTimeMillis();
+ sweepExpiredRunningJobs(nowMs);
+ sweepReplacedProcessEpochs();
+ driveRequiredRefreshes(nowMs);
+ dispatchPendingJobs();
+ }
+
+ /**
+ * 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(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() {
+ 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);
+ }
+ }
+ }
+
+ /**
+ * 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(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 - job.getUpdateTimeMs()
+ < 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(job);
+ } catch (Throwable t) {
+ LOG.warn("failed to drive the refresh of lance index job " +
job.getJobId(), t);
+ }
+ }
+ }
+
+ private void driveOneRefresh(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(job.getJobId(), refreshRevision, false);
+ return;
+ }
+ try {
+ // A half-orphan target (its db or table already dropped
externally) is a
+ // silent no-op: nothing is left to invalidate, and DONE is the
correct end
+ // state for the job.
+
Env.getCurrentEnv().getRefreshManager().handleRefreshTable(catalog.getName(),
+ 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(job.getJobId(), refreshRevision, false);
+ return;
+ }
+ finishRefreshTransition(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(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. Attempts at most
+ * {@link Config#lance_index_job_max_dispatch_per_round} fresh dispatches
per
+ * round, and never more than {@link
Config#lance_index_job_max_inflight_per_backend}
+ * in-flight jobs per backend, counted from the RUNNING snapshot 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.
+ */
+ private void dispatchPendingJobs() {
+ int maxPerRound = Math.max(1,
Config.lance_index_job_max_dispatch_per_round);
+ Map<Long, Integer> inflightByBackend = countInflightByBackend();
+ int attempted = 0;
+ for (LanceIndexJob job :
jobManager.getJobsNeedingDispatch(maxPerRound)) {
+ if (++attempted > maxPerRound) {
+ break;
+ }
+ try {
+ tryDispatch(job, inflightByBackend);
+ } catch (Throwable t) {
+ LOG.warn("failed to dispatch lance index job " +
job.getJobId(), t);
+ }
+ }
+ }
+
+ private Map<Long, Integer> countInflightByBackend() {
+ Map<Long, Integer> inflightByBackend = Maps.newHashMap();
+ for (LanceIndexJob job : jobManager.getAllJobsSnapshot()) {
Review Comment:
[P2] Avoid copying and sorting all historical jobs on every dispatcher tick.
getAllJobsSnapshot clones every durable job and sorts by ID under the manager
read lock, but this count only needs active backend slots. Resolved jobs are
never removed in this slice and quotas bound only unresolved jobs, so an idle
master pays unbounded O(history log history) work and allocation every 10
seconds. Query or maintain just the active slot holders without sorting
terminal history.
--
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]