u70b3 commented on code in PR #67978: URL: https://github.com/apache/doris/pull/67978#discussion_r4139924139
########## 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: done, the daemon now sleeps in slices bounded at 10s and re-reads the config at every wake, so a shortened interval takes effect within one slice (0c97204180). ########## 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: done, the dispatcher now asks the policy for all schedule-available backends and takes the first with a free possible-live slot (17568254ce). ########## 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: done, backend capacity is now counted from possible-live slot ownership via countPossibleLiveSlotsByBackend, so proven-dead RUNNING jobs stop blocking and UNKNOWN jobs still holding a slot keep counting (17568254ce). ########## 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: done, the dispatcher now resolves Env's current manager through a Supplier every round, covered by a new image-restart wiring test (b0fd3b41a8). ########## 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: done, slot counts are now aggregated in place under the read lock; no clone+sort of the full job history per tick anymore (17568254ce). -- 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]
