yuqi1129 commented on code in PR #13381:
URL: https://github.com/apache/gravitino/pull/13381#discussion_r4090672877


##########
core/src/main/java/org/apache/gravitino/connector/job/JobExecutor.java:
##########
@@ -60,15 +60,45 @@ public interface JobExecutor extends Closeable {
   String submitJob(JobTemplate jobTemplate);
 
   /**
-   * Get the status of a job by its unique identifier. The status should be 
one of the values in
-   * {@link JobHandle.Status}. The implementors should query the external job 
runner to get the job
-   * status, and map the status to the values in {@link JobHandle.Status}.
+   * Get a snapshot of the job's execution state by its unique identifier, 
including its status and
+   * when the job actually started and finished. The implementors should query 
the external job
+   * runner to get the job state, and map the status to the values in {@link 
JobHandle.Status}.
+   *
+   * <p>Gravitino pulls the job state periodically, so a job may go through 
several statuses between
+   * two pulls, for example, from {@link JobHandle.Status#QUEUED} straight to 
{@link
+   * JobHandle.Status#SUCCEEDED}. The timestamps are attributes of the job, 
not of a particular
+   * status: once the job has started, the started time must be reported in 
every later snapshot,
+   * including the terminal ones. Once reported, a timestamp should not change.
+   *
+   * <ul>
+   *   <li>The started time is null if the job hasn't started executing, or if 
it is unknown to the
+   *       job executor.
+   *   <li>The finished time is only set for a terminal status, and is null if 
it is unknown to the
+   *       job executor.
+   * </ul>
+   *
+   * <p>If the job runner can't tell when a job started or finished, leave the 
time null. Gravitino
+   * then falls back to the time it observes the job running or finished, 
which can be off by up to
+   * the job status pull interval. A job that starts and finishes between two 
pulls is never
+   * observed running, so it has no started time in that case.
+   *
+   * @param jobId The unique identifier of the job.
+   * @return The execution snapshot of the job.
+   * @throws NoSuchJobException If the job with the given identifier does not 
exist.
+   */
+  JobExecutionInfo getJobExecutionInfo(String jobId) throws NoSuchJobException;

Review Comment:
   **P1: An already-compiled custom executor can stop status polling 
permanently.** Making this method abstract makes source recompilation fail, but 
`JobExecutorFactory` loads existing executor jars by reflection without 
checking that they implement the new method. Such a class loads successfully 
and throws `AbstractMethodError` on the first `getJobExecutionInfo` call (I 
reproduced this with a plugin compiled against the old interface). The per-job 
handlers catch `Exception`/`RuntimeException`, not this `Error`, so it escapes 
`scheduleAtFixedRate` and suppresses all later polls. Please either retain a 
default bridge through `getJobStatus`, or validate the capability at startup 
and fail with an actionable error. A test using a class compiled against the 
old SPI would cover the upgrade path.



##########
core/src/main/java/org/apache/gravitino/job/JobManager.java:
##########
@@ -1239,49 +1198,236 @@ private void pullAndUpdateOwnedJobStatus(String 
metalake, JobEntity job) {
           job.name(),
           job.jobExecutionId(),
           metalake,
-          newStatus);
+          observed.status());
     } catch (Exception e) {
       // Keep the job unchanged, and retry it in the next poll.
-      newStatus = job.status();
       LOG.error(
           "Failed to pull or cancel job {} by execution id {}",
           job.name(),
           job.jobExecutionId(),
           e);
+      return;
     }
 
-    if (newStatus != job.status()) {
-      // Update the job entity with new status. entityStore.update() 
re-fetches the
-      // latest entity itself right before applying the updater, so the 
transition below
-      // is derived from latestJobEntity - the state as of right before the 
write - rather
-      // than the possibly-stale `job` snapshot taken by listJobs() above. A 
concurrent
-      // writer (e.g. cancelJob(), or another poll run) may have already moved 
the job to
-      // a terminal state, into CANCELLING, or recorded a real 
startedAt/finishedAt in the
-      // gap between that snapshot and this point; the updater must not 
regress any of
-      // that using the stale snapshot's view of the world.
-      JobHandle.Status finalNewStatus = newStatus;
-      updateJobEntity(
-              metalake,
-              job,
-              latestJobEntity -> toUpdatedStatusJobEntity(latestJobEntity, 
finalNewStatus))
-          .ifPresent(
-              updated ->
-                  LOG.info(
-                      "Updated the job {} with execution id {} status to {}",
-                      job.name(),
-                      job.jobExecutionId(),
-                      updated.status()));
+    // Besides the status, the job executor may report the timestamps of the 
job later than its
+    // status, so the job is also updated when only its timestamps change. The 
job is not written
+    // at all when nothing changes.
+    if (!changesJob(job, observed)) {
+      return;
     }
+
+    // Update the job entity with new status. entityStore.update() re-fetches 
the
+    // latest entity itself right before applying the updater, so the 
transition below
+    // is derived from latestJobEntity - the state as of right before the 
write - rather
+    // than the possibly-stale `job` snapshot taken by listJobs() above. A 
concurrent
+    // writer (e.g. cancelJob(), or another poll run) may have already moved 
the job to
+    // a terminal state, into CANCELLING, or recorded a real 
startedAt/finishedAt in the
+    // gap between that snapshot and this point; the updater must not regress 
any of
+    // that using the stale snapshot's view of the world.
+    JobExecutionInfo finalObserved = observed;
+    updateJobEntity(
+            metalake,
+            job,
+            latestJobEntity -> toUpdatedStatusJobEntity(latestJobEntity, 
finalObserved))
+        .ifPresent(
+            updated ->
+                LOG.info(
+                    "Updated the job {} with execution id {} status to {}",
+                    job.name(),
+                    job.jobExecutionId(),
+                    updated.status()));
   }
 
-  private JobHandle.Status cancelOwnedJob(String metalake, JobEntity job) {
+  private JobExecutionInfo cancelOwnedJob(String metalake, JobEntity job) {
     LOG.info(
         "Cancelling job {} with execution id {} under metalake {} as it is 
marked as CANCELLING",
         job.name(),
         job.jobExecutionId(),
         metalake);
     jobExecutor.cancelJob(job.jobExecutionId());
-    return jobExecutor.getJobStatus(job.jobExecutionId());
+    return jobExecutor.getJobExecutionInfo(job.jobExecutionId());
+  }
+
+  // Drops or corrects the timestamps reported by the job executor that are 
unusable, or
+  // inconsistent
+  // with the status or with each other. A faulty job executor never fails the 
status pull. The
+  // corrections are only logged at debug level, as the same report is 
corrected on every pull.
+  private JobExecutionInfo sanitizeExecutionInfo(JobEntity job, 
JobExecutionInfo observed) {
+    Instant startedAt = usableTime(job, "started", observed.startedAt());
+    Instant finishedAt = usableTime(job, "finished", observed.finishedAt());
+
+    if (startedAt != null && observed.status() == JobHandle.Status.QUEUED) {
+      LOG.debug(
+          "Job {} with execution id {} is reported as QUEUED with a started 
time {}, ignoring "
+              + "the started time",
+          job.name(),
+          job.jobExecutionId(),
+          startedAt);
+      startedAt = null;
+    }
+    if (finishedAt != null && !isFinishedStatus(observed.status())) {
+      LOG.debug(
+          "Job {} with execution id {} is reported as {} with a finished time 
{}, ignoring the "
+              + "finished time",
+          job.name(),
+          job.jobExecutionId(),
+          observed.status(),
+          finishedAt);
+      finishedAt = null;
+    }
+
+    startedAt = notBeforeQueuedAt(job, "started", startedAt);
+    finishedAt = notBeforeQueuedAt(job, "finished", finishedAt);
+
+    if (startedAt != null && finishedAt != null && 
finishedAt.isBefore(startedAt)) {
+      // The finished time is kept, as the cleanup of finished jobs relies on 
it.
+      LOG.debug(
+          "Job {} with execution id {} is reported to finish at {} before it 
started at {}, "
+              + "ignoring the started time",
+          job.name(),
+          job.jobExecutionId(),
+          finishedAt,
+          startedAt);
+      startedAt = null;
+    }
+
+    return 
observed.toBuilder().withStartedAt(startedAt).withFinishedAt(finishedAt).build();
+  }
+
+  // A reported time is stored in epoch milliseconds, so a time that doesn't 
fit is dropped. So is a
+  // time far in the future, e.g. a sentinel like Instant.MAX, which would 
otherwise keep a finished
+  // job from ever being cleaned up.
+  @Nullable
+  private static Instant usableTime(JobEntity job, String event, @Nullable 
Instant reportedTime) {
+    if (reportedTime == null) {
+      return null;
+    }
+    boolean usable;
+    try {
+      reportedTime.toEpochMilli();
+      usable = 
!reportedTime.isAfter(Instant.now().plus(MAX_REPORTED_TIME_AHEAD));
+    } catch (ArithmeticException e) {
+      usable = false;
+    }
+    if (!usable) {
+      LOG.debug(
+          "Job {} with execution id {} is reported to be {} at an unusable 
time {}, ignoring it",
+          job.name(),
+          job.jobExecutionId(),
+          event,
+          reportedTime);
+      return null;
+    }
+    return reportedTime;
+  }
+
+  // The job executor may run on another host whose clock is behind 
Gravitino's, so a reported time
+  // can be earlier than when Gravitino queued the job. Such a time is raised 
to the queued time.
+  @Nullable
+  private Instant notBeforeQueuedAt(JobEntity job, String event, @Nullable 
Instant reportedTime) {
+    Instant queuedAt = job.auditInfo() == null ? null : 
job.auditInfo().createTime();
+    if (reportedTime == null || queuedAt == null || 
!reportedTime.isBefore(queuedAt)) {
+      return reportedTime;
+    }
+    LOG.debug(
+        "Job {} with execution id {} is reported to be {} at {}, before it was 
queued at {}, "
+            + "using the queued time instead",
+        job.name(),
+        job.jobExecutionId(),
+        event,
+        reportedTime,
+        queuedAt);
+    return queuedAt;
+  }
+
+  // Whether applying the observed execution info changes the recorded status 
or timestamps of the
+  // job. It is derived by the same logic as the update, so a value that the 
update corrects is not
+  // seen as a change on every poll.
+  private boolean changesJob(JobEntity job, JobExecutionInfo observed) {
+    JobEntity updated = toUpdatedStatusJobEntity(job, observed);
+    return updated.status() != job.status()
+        || timeInMs(updated.startedAt()) != timeInMs(job.startedAt())
+        || timeInMs(updated.finishedAt()) != timeInMs(job.finishedAt());
+  }
+
+  private JobEntity toUpdatedStatusJobEntity(JobEntity latestJobEntity, 
JobExecutionInfo observed) {
+    JobHandle.Status currentStatus = latestJobEntity.status();
+    JobHandle.Status observedStatus = observed.status();
+
+    // Never regress a job out of a terminal state, and never move a 
CANCELLING job back to a
+    // non-terminal state - both would only be possible here because the 
executor status was
+    // observed against a stale snapshot of the job.
+    if (isFinishedStatus(currentStatus)

Review Comment:
   **P2: Executor timestamps reported after a terminal status cannot be 
persisted.** For an eventually consistent external runner, poll A can observe 
`SUCCEEDED` with null timestamps, causing Gravitino to store a poll-time 
`finishedAt`. On poll B the runner may return the actual `startedAt` and 
`finishedAt`, but `pullAndUpdateJobStatus` filters out terminal jobs (lines 
722-729); this early return would also discard the update. The missing start 
and inaccurate finish therefore become permanent, despite the new path allowing 
timestamps to arrive after status. Please add a bounded way to revisit terminal 
jobs awaiting timestamps, and test this two-poll sequence.



-- 
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]

Reply via email to