github-actions[bot] commented on code in PR #67978:
URL: https://github.com/apache/doris/pull/67978#discussion_r4140934286
##########
fe/fe-common/src/main/java/org/apache/doris/common/Config.java:
##########
@@ -4276,4 +4276,64 @@ public void handle(Field field, String value) throws
Exception {
"Static upper bound for num_sub_vectors of Lance IVF_PQ
indexes."})
public static int lance_index_max_num_sub_vectors = 256;
+ @ConfField(mutable = true, masterOnly = true,
+ callback =
LanceIndexConfigValidator.PositiveIntConfigHandler.class,
+ description = {"Lance 索引 job 派发器(含 deadline/possible-live 扫掠与
refresh 驱动)的轮询周期(秒)。",
+ "Polling interval in seconds of the Lance index job
dispatcher "
+ + "(dispatch sweep, deadline/possible-live sweeps, and
refresh driver)."})
+ public static int lance_index_job_dispatch_interval_second = 10;
Review Comment:
[P2] Publish the five mutable numeric settings to the dispatcher thread.
ADMIN SET writes these plain static fields under ConfigBase's monitor, while
the daemon reads the interval, deadline, dispatch caps, and refresh retry
without that monitor. There is no happens-before edge, so a successful update
may leave the daemon using stale values, including a former per-BE cap. The
earlier visibility fix covered only the adjacent pause and local-file switches;
make these separate numeric fields volatile or use shared synchronization.
##########
fe/fe-core/src/main/java/org/apache/doris/service/FrontendServiceImpl.java:
##########
@@ -1145,6 +1147,23 @@ public TMasterResult finishTask(TFinishTaskRequest
request) throws TException {
return masterImpl.finishTask(request);
}
+ /**
+ * Typed result envelope of one Lance index mutation invocation. Stale or
+ * identity-mismatched reports are logged and dropped; a complete matched
+ * report is classified into the durable job state. This layer stays thin:
+ * only the master accepts reports, and everything beyond identity checking
+ * and classification lives in the report handler.
+ */
+ @Override
+ public TStatus reportLanceIndexJobResult(TLanceIndexJobReport report)
throws TException {
Review Comment:
[P1] Authenticate Lance job reports before accepting them. A user with table
SHOW can read JobId, InvocationId, BE epoch, and Revision from SHOW LANCE INDEX
JOB; immediately after markRunning, that Revision is the dispatch revision.
This RPC checks only mastership, and the FE Thrift server does not authenticate
its caller, so a client able to reach that port can submit NATIVE_OK with
CHILD_REAPED for a RUNNING job. The manager then durably marks an unexecuted
job COMMITTED and releases its slot; the genuine BE response is rejected as
late. Bind the report to the selected BE and stop exposing a report-authorizing
invocation secret through SHOW.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,699 @@
+// 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 org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+import org.apache.thrift.TApplicationException;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+import java.util.function.Supplier;
+
+/**
+ * Master-only daemon that drives the durable Lance index job records through
+ * the lifecycle after admission. Each round runs in a fixed order: converge
+ * expired RUNNING jobs to UNKNOWN, release possible-live slots whose backend
+ * process was replaced, drive the refresh a terminal job still owes, then
+ * dispatch PENDING jobs. Every durable transition goes through
+ * {@link LanceIndexJobManager} under its own lock; the daemon holds no catalog
+ * or manager lock across any call.
+ *
+ * <p>The daemon does not read the admission gate: a job that is already
durable
+ * must be driven to its terminal state, whatever the gate says now, so the
+ * thread runs unconditionally on the master and simply finds nothing to do
+ * while no jobs exist. An idle round writes no journal record.
+ *
+ * <p>Dispatch follows the durable-before-send boundary: the whole request is
+ * prepared first (so a preparation failure just leaves the job PENDING), then
+ * the markRunning edit log is written and re-read before the first byte of
+ * network I/O, and the invocation id of an attempt that lost the
compare-and-set
+ * is never reused. After a successful markRunning there is exactly one send;
+ * from that point a job converges only through a matching result callback, the
+ * deadline sweep, or the epoch sweep, never through a resend. A failure that
+ * still proves the dispatch was never enqueued (a clean pre-enqueue error
+ * status, a client-pool borrow failure, or an UNKNOWN_METHOD answer from an
+ * old backend) converges it NOT_COMMITTED through the no-enqueue channel,
+ * which releases the possible-live slot in the same durable transition;
+ * anything ambiguous after the invocation may have started converges UNKNOWN
+ * with the slot retained. The blocking time of the send loop is bounded per
+ * round (one backend RPC timeout), because this thread is also the only thread
+ * running the sweeps and the refresh driver — see
+ * {@link #dispatchPendingJobs()}.
+ *
+ * <p>The manager is resolved from the supplier once per round rather than
+ * captured at construction: {@code Env.loadLanceIndexJobManager} replaces the
+ * Env-owned manager with a brand-new object on every image load, so a cached
+ * reference would keep scanning the abandoned pre-image manager after an FE
+ * restart while replay, admission and SHOW all move on to the restored one.
+ * Every phase of one round shares the single resolved instance.
+ *
+ * <p>The sleep between rounds is sliced at {@link #MAX_SLEEP_SLICE_MS} so a
+ * shortened polling interval takes effect within one slice (see the field
+ * javadoc), and {@link Config#lance_index_job_dispatcher_paused} suspends only
+ * the dispatch phase (see {@link #dispatchPendingJobs}).
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+ private static final Logger LOG =
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+ /**
+ * Upper bound of one sleep slice, equal to the shipped default interval.
The
+ * daemon never sleeps longer than this, so a shortened
+ * {@link Config#lance_index_job_dispatch_interval_second} takes effect
within
+ * one slice instead of waiting out a previously adopted long sleep: the
+ * elapsed check in {@link #runAfterCatalogReady} is re-evaluated against
the
+ * current config at every wake. Slices bound only the sleep; rounds still
+ * honor the configured interval, because a wake whose configured interval
+ * (longer than this bound) has not elapsed since the last round skips the
+ * round. A lengthened interval takes effect at the next wake through the
same
+ * check, and an interval at or below this bound needs no check at all —
every
+ * wake runs a round, exactly one per configured period.
+ */
+ private static final long MAX_SLEEP_SLICE_MS = 10_000L;
+
+ private final Supplier<LanceIndexJobManager> jobManagerSupplier;
+
+ /** Wall time of the last executed round, or -1 before the first one. */
+ private long lastRoundMs = -1L;
+
+ public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+ this(() -> jobManager);
+ }
+
+ public LanceIndexJobDispatcher(Supplier<LanceIndexJobManager>
jobManagerSupplier) {
+ super("lance index job dispatcher", dispatchIntervalMs());
+ this.jobManagerSupplier = jobManagerSupplier;
+ }
+
+ /**
+ * Values loaded from fe.conf bypass the config validator (only ADMIN SET
runs
+ * it), so the positive invariant is re-asserted where a non-positive value
+ * would break the loop: a non-positive interval would kill this thread
inside
+ * {@code Thread.sleep} or busy-spin it, a non-positive deadline would
sweep
+ * every dispatched job UNKNOWN on the next round, and a zero cap would
stall
+ * dispatch forever. The refresh retry interval needs no such defense: a
+ * non-positive value simply disengages the throttle.
+ */
+ private static long dispatchIntervalMs() {
+ return Math.max(1, Config.lance_index_job_dispatch_interval_second) *
1000L;
+ }
+
+ private static long executeDeadlineMs(long nowMs) {
+ long second = Math.max(1L,
Config.lance_index_job_execute_deadline_second);
+ return second > (Long.MAX_VALUE - nowMs) / 1000L ? Long.MAX_VALUE :
nowMs + second * 1000L;
+ }
+
+ @Override
+ protected void runAfterCatalogReady() {
+ if (!Env.getCurrentEnv().isMaster()) {
+ return;
+ }
+ if (Env.isCheckpointThread()) {
+ return;
+ }
+ long configuredMs = dispatchIntervalMs();
+ setInterval(Math.min(configuredMs, MAX_SLEEP_SLICE_MS));
+ if (configuredMs > MAX_SLEEP_SLICE_MS && lastRoundMs >= 0 && nowMs() -
lastRoundMs < configuredMs) {
+ // A wake inside a long configured interval: the slice elapsed, the
+ // round period has not. Skipping is cheap and writes no journal
record.
+ return;
+ }
+ lastRoundMs = nowMs();
+ try {
+ runOneRound(jobManagerSupplier.get());
+ } catch (Throwable t) {
+ LOG.warn("Failed to process one round of the lance index job
dispatcher", t);
+ }
+ }
+
+ /** Clock seam for the round-period check; tests advance it instead of
sleeping. */
+ protected long nowMs() {
+ return System.currentTimeMillis();
+ }
+
+ private void runOneRound(LanceIndexJobManager jobManager) {
+ long nowMs = System.currentTimeMillis();
+ sweepExpiredRunningJobs(jobManager, nowMs);
+ sweepReplacedProcessEpochs(jobManager);
+ driveRequiredRefreshes(jobManager, nowMs);
+ dispatchPendingJobs(jobManager);
+ }
+
+ /**
+ * Deadline sweep. A RUNNING job past its wait deadline has produced no
+ * complete trusted result, so it converges to UNKNOWN through the same
+ * completeWithResult channel a callback would use. Expiry bounds the wait
+ * only: it never proves termination, so the possible-live slot, the
+ * same-name fence, and the unresolved quota all stay held.
+ */
+ private void sweepExpiredRunningJobs(LanceIndexJobManager jobManager, long
nowMs) {
+ for (LanceIndexJob job : jobManager.getExpiredRunningJobs(nowMs)) {
+ try {
+ boolean completed =
jobManager.completeWithResult(job.getJobId(),
+ dispatchRevisionOf(job), job.getInvocationId(),
job.getBeProcessEpoch(),
+ new
LanceIndexJobResult(LanceIndexJobResultCode.NO_TRUSTED_RESULT,
+ LanceIndexJobCompletionReason.NONE,
+ "execute deadline expired without a complete
trusted result", false));
+ if (completed) {
+ LOG.info("lance index job {} converged RUNNING -> UNKNOWN
on deadline expiry",
+ job.getJobId());
+ } else {
+ LOG.warn("deadline sweep skipped lance index job {}:
already converged by a callback or sweep",
+ job.getJobId());
+ }
+ } catch (Throwable t) {
+ LOG.warn("failed to sweep expired lance index job " +
job.getJobId(), t);
+ }
+ }
+ }
+
+ /**
+ * Possible-live sweep. The only slot-release proof this daemon produces is
+ * that the recorded backend process epoch no longer exists: a backend
entry
+ * reporting a different epoch proves the process that received the
dispatch
+ * was replaced. A missing backend entry or heartbeat loss proves nothing
+ * (the worker may still be running behind a partition), so such a job
keeps
+ * its slot until a stronger proof or an operator force release. An epoch
+ * change also proves nothing about the outcome, so the mutation state is
+ * never touched here.
+ */
+ private void sweepReplacedProcessEpochs(LanceIndexJobManager jobManager) {
+ for (LanceIndexJob job : jobManager.getJobsHoldingPossibleLiveSlot()) {
+ try {
+ Backend backend =
Env.getCurrentSystemInfo().getBackend(job.getBackendId());
+ if (backend == null || backend.getProcessEpoch() ==
job.getBeProcessEpoch()) {
+ continue;
+ }
+ boolean recorded =
jobManager.recordTerminationProof(job.getJobId(),
+ dispatchRevisionOf(job), job.getBackendId(),
job.getBeProcessEpoch(),
+ job.getInvocationId(),
LanceIndexTerminationProof.BE_PROCESS_EPOCH_GONE);
+ if (recorded) {
+ LOG.info("released possible-live slot of lance index job
{}: backend process epoch was replaced",
+ job.getJobId());
+ } else {
+ LOG.warn("epoch sweep skipped lance index job {}: dispatch
identity already moved",
+ job.getJobId());
+ }
+ } catch (Throwable t) {
+ LOG.warn("failed to sweep possible-live slot of lance index
job " + job.getJobId(), t);
+ }
+ }
+ }
+
+ /**
+ * The clock the FAILED-refresh throttle measures from: the dedicated
+ * refresh-failure timestamp, falling back to the generic update time only
for
+ * a legacy record replayed before the field existed. Measuring from the
+ * generic time would let an unrelated transition — a CHILD_REAPED proof
or an
+ * epoch-gone release bumping it — postpone the next retry by a full
interval
+ * while the fence and quota stay held.
+ */
+ private static long refreshThrottledSinceMs(LanceIndexJob job) {
+ return job.getRefreshFailureTimeMs() == null ? job.getUpdateTimeMs() :
job.getRefreshFailureTimeMs();
+ }
+
+ /**
+ * Refresh driver for terminal jobs with an unfinished refresh obligation.
+ * Completing the refresh is the protocol duty that releases the same-name
+ * fence and the unresolved quota once DONE; it is not a read-visibility
+ * action, because index metadata is never cached. Each job is driven
+ * through markRefreshRunning, the idempotent external-table refresh, then
+ * DONE or FAILED: a FAILED job keeps its fence and is retried, throttled
to
+ * one attempt per retry interval, while a first REQUIRED refresh is never
+ * delayed. UNKNOWN jobs never appear here; they owe no refresh.
+ */
+ private void driveRequiredRefreshes(LanceIndexJobManager jobManager, long
nowMs) {
+ for (LanceIndexJob job : jobManager.getJobsNeedingRefresh()) {
+ try {
+ if (job.getRefreshState() ==
LanceIndexJobRefreshState.RUNNING) {
+ // In flight elsewhere; the master-transfer sweep
downgrades a stale
+ // RUNNING back to REQUIRED, so a lost driver cannot
strand it.
+ continue;
+ }
+ if (job.getRefreshState() == LanceIndexJobRefreshState.FAILED
+ && nowMs - refreshThrottledSinceMs(job)
+ < Config.lance_index_job_refresh_retry_second
* 1000L) {
+ continue;
+ }
+ if (!jobManager.markRefreshRunning(job.getJobId(),
job.getRevision())) {
+ // A concurrent driver won the compare-and-set; nothing to
do here.
+ continue;
+ }
+ driveOneRefresh(jobManager, job);
+ } catch (Throwable t) {
+ LOG.warn("failed to drive the refresh of lance index job " +
job.getJobId(), t);
+ }
+ }
+ }
+
+ private void driveOneRefresh(LanceIndexJobManager jobManager,
LanceIndexJob job) {
+ long refreshRevision = job.getRevision() + 1;
+ CatalogIf catalog =
Env.getCurrentEnv().getCatalogMgr().getCatalog(job.getCatalogId());
+ if (catalog == null) {
+ // Unreachable while the unresolved-job guard blocks catalog
drops; kept as a
+ // fail-closed fallback so the job still transitions and retries
later.
+ LOG.warn("catalog of lance index job {} is gone; marking its
refresh FAILED", job.getJobId());
+ finishRefreshTransition(jobManager, job.getJobId(),
refreshRevision, false);
+ return;
+ }
+ 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. The refresh is addressed by the persisted
catalog id,
+ // never by the mutable name: a rename that hands this catalog's
old name
+ // to a different catalog between the two resolutions must not
refresh
+ // that one and bill the outcome to this job.
+
Env.getCurrentEnv().getRefreshManager().handleRefreshTable(job.getCatalogId(),
Review Comment:
[P2] Keep the refresh obligation when catalog initialization fails. This
call passes ignoreIfNotExists=true, so RefreshManager treats a null database as
successful absence; ExternalCatalog.getDbNullable also returns null after
catching makeSureInitialized errors. On a new master with a cold Lance catalog
and a transient client/namespace initialization failure, the driver therefore
marks DONE and releases the fence/quota without invalidating table caches or
writing a refresh log. Distinguish initialization failure from confirmed
absence so the job becomes FAILED and retries.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,699 @@
+// 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 org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+import org.apache.thrift.TApplicationException;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+import java.util.function.Supplier;
+
+/**
+ * Master-only daemon that drives the durable Lance index job records through
+ * the lifecycle after admission. Each round runs in a fixed order: converge
+ * expired RUNNING jobs to UNKNOWN, release possible-live slots whose backend
+ * process was replaced, drive the refresh a terminal job still owes, then
+ * dispatch PENDING jobs. Every durable transition goes through
+ * {@link LanceIndexJobManager} under its own lock; the daemon holds no catalog
+ * or manager lock across any call.
+ *
+ * <p>The daemon does not read the admission gate: a job that is already
durable
+ * must be driven to its terminal state, whatever the gate says now, so the
+ * thread runs unconditionally on the master and simply finds nothing to do
+ * while no jobs exist. An idle round writes no journal record.
+ *
+ * <p>Dispatch follows the durable-before-send boundary: the whole request is
+ * prepared first (so a preparation failure just leaves the job PENDING), then
+ * the markRunning edit log is written and re-read before the first byte of
+ * network I/O, and the invocation id of an attempt that lost the
compare-and-set
+ * is never reused. After a successful markRunning there is exactly one send;
+ * from that point a job converges only through a matching result callback, the
+ * deadline sweep, or the epoch sweep, never through a resend. A failure that
+ * still proves the dispatch was never enqueued (a clean pre-enqueue error
+ * status, a client-pool borrow failure, or an UNKNOWN_METHOD answer from an
+ * old backend) converges it NOT_COMMITTED through the no-enqueue channel,
+ * which releases the possible-live slot in the same durable transition;
+ * anything ambiguous after the invocation may have started converges UNKNOWN
+ * with the slot retained. The blocking time of the send loop is bounded per
+ * round (one backend RPC timeout), because this thread is also the only thread
+ * running the sweeps and the refresh driver — see
+ * {@link #dispatchPendingJobs()}.
+ *
+ * <p>The manager is resolved from the supplier once per round rather than
+ * captured at construction: {@code Env.loadLanceIndexJobManager} replaces the
+ * Env-owned manager with a brand-new object on every image load, so a cached
+ * reference would keep scanning the abandoned pre-image manager after an FE
+ * restart while replay, admission and SHOW all move on to the restored one.
+ * Every phase of one round shares the single resolved instance.
+ *
+ * <p>The sleep between rounds is sliced at {@link #MAX_SLEEP_SLICE_MS} so a
+ * shortened polling interval takes effect within one slice (see the field
+ * javadoc), and {@link Config#lance_index_job_dispatcher_paused} suspends only
+ * the dispatch phase (see {@link #dispatchPendingJobs}).
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+ private static final Logger LOG =
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+ /**
+ * Upper bound of one sleep slice, equal to the shipped default interval.
The
+ * daemon never sleeps longer than this, so a shortened
+ * {@link Config#lance_index_job_dispatch_interval_second} takes effect
within
+ * one slice instead of waiting out a previously adopted long sleep: the
+ * elapsed check in {@link #runAfterCatalogReady} is re-evaluated against
the
+ * current config at every wake. Slices bound only the sleep; rounds still
+ * honor the configured interval, because a wake whose configured interval
+ * (longer than this bound) has not elapsed since the last round skips the
+ * round. A lengthened interval takes effect at the next wake through the
same
+ * check, and an interval at or below this bound needs no check at all —
every
+ * wake runs a round, exactly one per configured period.
+ */
+ private static final long MAX_SLEEP_SLICE_MS = 10_000L;
+
+ private final Supplier<LanceIndexJobManager> jobManagerSupplier;
+
+ /** Wall time of the last executed round, or -1 before the first one. */
+ private long lastRoundMs = -1L;
+
+ public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+ this(() -> jobManager);
+ }
+
+ public LanceIndexJobDispatcher(Supplier<LanceIndexJobManager>
jobManagerSupplier) {
+ super("lance index job dispatcher", dispatchIntervalMs());
+ this.jobManagerSupplier = jobManagerSupplier;
+ }
+
+ /**
+ * Values loaded from fe.conf bypass the config validator (only ADMIN SET
runs
+ * it), so the positive invariant is re-asserted where a non-positive value
+ * would break the loop: a non-positive interval would kill this thread
inside
+ * {@code Thread.sleep} or busy-spin it, a non-positive deadline would
sweep
+ * every dispatched job UNKNOWN on the next round, and a zero cap would
stall
+ * dispatch forever. The refresh retry interval needs no such defense: a
+ * non-positive value simply disengages the throttle.
+ */
+ private static long dispatchIntervalMs() {
+ return Math.max(1, Config.lance_index_job_dispatch_interval_second) *
1000L;
+ }
+
+ private static long executeDeadlineMs(long nowMs) {
+ long second = Math.max(1L,
Config.lance_index_job_execute_deadline_second);
+ return second > (Long.MAX_VALUE - nowMs) / 1000L ? Long.MAX_VALUE :
nowMs + second * 1000L;
+ }
+
+ @Override
+ protected void runAfterCatalogReady() {
+ if (!Env.getCurrentEnv().isMaster()) {
+ return;
+ }
+ if (Env.isCheckpointThread()) {
+ return;
+ }
+ long configuredMs = dispatchIntervalMs();
+ setInterval(Math.min(configuredMs, MAX_SLEEP_SLICE_MS));
+ if (configuredMs > MAX_SLEEP_SLICE_MS && lastRoundMs >= 0 && nowMs() -
lastRoundMs < configuredMs) {
+ // A wake inside a long configured interval: the slice elapsed, the
+ // round period has not. Skipping is cheap and writes no journal
record.
+ return;
+ }
+ lastRoundMs = nowMs();
+ try {
+ runOneRound(jobManagerSupplier.get());
+ } catch (Throwable t) {
+ LOG.warn("Failed to process one round of the lance index job
dispatcher", t);
+ }
+ }
+
+ /** Clock seam for the round-period check; tests advance it instead of
sleeping. */
+ protected long nowMs() {
+ return System.currentTimeMillis();
+ }
+
+ private void runOneRound(LanceIndexJobManager jobManager) {
+ long nowMs = System.currentTimeMillis();
+ sweepExpiredRunningJobs(jobManager, nowMs);
+ sweepReplacedProcessEpochs(jobManager);
+ driveRequiredRefreshes(jobManager, nowMs);
+ dispatchPendingJobs(jobManager);
+ }
+
+ /**
+ * Deadline sweep. A RUNNING job past its wait deadline has produced no
+ * complete trusted result, so it converges to UNKNOWN through the same
+ * completeWithResult channel a callback would use. Expiry bounds the wait
+ * only: it never proves termination, so the possible-live slot, the
+ * same-name fence, and the unresolved quota all stay held.
+ */
+ private void sweepExpiredRunningJobs(LanceIndexJobManager jobManager, long
nowMs) {
+ for (LanceIndexJob job : jobManager.getExpiredRunningJobs(nowMs)) {
+ try {
+ boolean completed =
jobManager.completeWithResult(job.getJobId(),
+ dispatchRevisionOf(job), job.getInvocationId(),
job.getBeProcessEpoch(),
+ new
LanceIndexJobResult(LanceIndexJobResultCode.NO_TRUSTED_RESULT,
+ LanceIndexJobCompletionReason.NONE,
+ "execute deadline expired without a complete
trusted result", false));
+ if (completed) {
+ LOG.info("lance index job {} converged RUNNING -> UNKNOWN
on deadline expiry",
+ job.getJobId());
+ } else {
+ LOG.warn("deadline sweep skipped lance index job {}:
already converged by a callback or sweep",
+ job.getJobId());
+ }
+ } catch (Throwable t) {
+ LOG.warn("failed to sweep expired lance index job " +
job.getJobId(), t);
+ }
+ }
+ }
+
+ /**
+ * Possible-live sweep. The only slot-release proof this daemon produces is
+ * that the recorded backend process epoch no longer exists: a backend
entry
+ * reporting a different epoch proves the process that received the
dispatch
+ * was replaced. A missing backend entry or heartbeat loss proves nothing
+ * (the worker may still be running behind a partition), so such a job
keeps
+ * its slot until a stronger proof or an operator force release. An epoch
+ * change also proves nothing about the outcome, so the mutation state is
+ * never touched here.
+ */
+ private void sweepReplacedProcessEpochs(LanceIndexJobManager jobManager) {
+ for (LanceIndexJob job : jobManager.getJobsHoldingPossibleLiveSlot()) {
+ try {
+ Backend backend =
Env.getCurrentSystemInfo().getBackend(job.getBackendId());
+ if (backend == null || backend.getProcessEpoch() ==
job.getBeProcessEpoch()) {
+ continue;
+ }
+ boolean recorded =
jobManager.recordTerminationProof(job.getJobId(),
+ dispatchRevisionOf(job), job.getBackendId(),
job.getBeProcessEpoch(),
+ job.getInvocationId(),
LanceIndexTerminationProof.BE_PROCESS_EPOCH_GONE);
+ if (recorded) {
+ LOG.info("released possible-live slot of lance index job
{}: backend process epoch was replaced",
+ job.getJobId());
+ } else {
+ LOG.warn("epoch sweep skipped lance index job {}: dispatch
identity already moved",
+ job.getJobId());
+ }
+ } catch (Throwable t) {
+ LOG.warn("failed to sweep possible-live slot of lance index
job " + job.getJobId(), t);
+ }
+ }
+ }
+
+ /**
+ * The clock the FAILED-refresh throttle measures from: the dedicated
+ * refresh-failure timestamp, falling back to the generic update time only
for
+ * a legacy record replayed before the field existed. Measuring from the
+ * generic time would let an unrelated transition — a CHILD_REAPED proof
or an
+ * epoch-gone release bumping it — postpone the next retry by a full
interval
+ * while the fence and quota stay held.
+ */
+ private static long refreshThrottledSinceMs(LanceIndexJob job) {
+ return job.getRefreshFailureTimeMs() == null ? job.getUpdateTimeMs() :
job.getRefreshFailureTimeMs();
+ }
+
+ /**
+ * Refresh driver for terminal jobs with an unfinished refresh obligation.
+ * Completing the refresh is the protocol duty that releases the same-name
+ * fence and the unresolved quota once DONE; it is not a read-visibility
+ * action, because index metadata is never cached. Each job is driven
+ * through markRefreshRunning, the idempotent external-table refresh, then
+ * DONE or FAILED: a FAILED job keeps its fence and is retried, throttled
to
+ * one attempt per retry interval, while a first REQUIRED refresh is never
+ * delayed. UNKNOWN jobs never appear here; they owe no refresh.
+ */
+ private void driveRequiredRefreshes(LanceIndexJobManager jobManager, long
nowMs) {
+ for (LanceIndexJob job : jobManager.getJobsNeedingRefresh()) {
Review Comment:
[P2] Bound refresh work before returning to the lifecycle sweeps. This loop
processes every owed job synchronously on the sole dispatcher thread, while
handleRefreshTable can initialize a Lance catalog and list remote databases or
tables. A slow provider can consume a metadata timeout, and multiple terminal
jobs can multiply that delay before the next deadline or BE-epoch sweep and
before PENDING jobs dispatch. The existing one-RPC send budget starts after
this loop and cannot bound it. Limit refresh work per round or drive it
independently while preserving the durable transitions.
--
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]