jerryshao commented on code in PR #13373:
URL: https://github.com/apache/gravitino/pull/13373#discussion_r4070108760


##########
core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutor.java:
##########
@@ -404,29 +505,201 @@ void cleanupJobStatus() {
       jobStatus
           .entrySet()
           .removeIf(
-              entry -> {
-                boolean expired =
-                    entry.getValue().getRight() != UNEXPIRED_TIME_IN_MS
-                        && (currentTime - entry.getValue().getRight()) >= 
jobStatusKeepTimeInMs;
-                if (expired) {
-                  jobWorkingDirs.remove(entry.getKey());
-                }
-                return expired;
-              });
+              entry ->
+                  entry.getValue().getRight() != UNEXPIRED_TIME_IN_MS
+                      && (currentTime - entry.getValue().getRight()) >= 
jobStatusKeepTimeInMs);
+    }
+  }
+
+  /**
+   * Removes the output index files whose job staging directory no longer 
exists. The staging
+   * directories themselves are removed by JobManager, so this never deletes 
any job output.
+   */
+  @VisibleForTesting
+  void cleanupOutputIndexes() {
+    long now = System.currentTimeMillis();
+    // An exception escaping a scheduled task would cancel all its later runs.
+    try (DirectoryStream<Path> indexFiles =
+        Files.newDirectoryStream(
+            outputIndexDir, LOCAL_JOB_PREFIX + "*" + 
OUTPUT_INDEX_FILE_SUFFIX)) {
+      for (Path indexFile : indexFiles) {
+        // close() interrupts this thread, and every later read would then 
fail as well.
+        if (Thread.currentThread().isInterrupted()) {
+          return;
+        }
+        try {
+          // Anything but an index file, e.g. a directory of a colliding 
legacy metalake name.
+          if (!Files.isRegularFile(indexFile, LinkOption.NOFOLLOW_LINKS)) {
+            continue;
+          }
+          if (now - Files.getLastModifiedTime(indexFile).toMillis() < 
OUTPUT_INDEX_MIN_AGE_IN_MS) {
+            continue;
+          }
+          Path workingDir = parseOutputIndex(Files.readAllBytes(indexFile));
+          // An index this server can't interpret may come from a newer 
server, so keep it. Only
+          // a confirmed missing directory counts: notExists() is false on I/O 
errors too, so a
+          // storage hiccup never removes an index.
+          if (workingDir != null && Files.notExists(workingDir)) {
+            Files.deleteIfExists(indexFile);
+            LOG.debug("Removed output index {} of a deleted job staging 
directory", indexFile);
+          }
+        } catch (ClosedByInterruptException e) {
+          return;
+        } catch (NoSuchFileException e) {
+          // Removed concurrently, e.g. by another server sharing the staging 
directory.
+        } catch (IOException | RuntimeException e) {
+          LOG.warn("Failed to clean up job output index {}", indexFile, e);
+        }
+      }
+    } catch (IOException | RuntimeException e) {
+      LOG.warn("Failed to clean up job output indexes under {}", 
outputIndexDir, e);
     }
   }
 
   private List<String> getJobOutput(String jobId, String fileName, int 
maxLines, int maxBytes) {
-    File workingDir = getWorkingDir(jobId);
+    Path workingDir = locateWorkingDir(jobId);
     if (workingDir == null) {
       return ImmutableList.of();
     }
-    return readLastLines(new File(workingDir, fileName), maxLines, maxBytes);
+    return readLastLines(workingDir.resolve(fileName).toFile(), maxLines, 
maxBytes);
+  }
+
+  private void writeOutputIndex(String jobId, JobTemplate jobTemplate) {
+    // The job is already queued, so any failure here must not fail the 
submission: the caller
+    // would then clean up the staging directory of a job that still runs.
+    try {
+      Path workingDir =
+          LocalProcessBuilder.resolveWorkingDirectory(jobTemplate)
+              .toPath()
+              .toAbsolutePath()
+              .normalize();
+      if (!workingDir.startsWith(stagingRoot) || 
workingDir.equals(stagingRoot)) {
+        LOG.warn(
+            "The working directory {} of job {} is not under the job staging 
directory {}, so "
+                + "its output can't be retrieved",
+            workingDir,
+            jobId,
+            stagingRoot);
+        return;
+      }
+
+      // Relative to the staging directory, so servers mounting shared storage 
at different paths
+      // can all resolve it. '/'-separated regardless of the platform's 
separator.
+      ObjectNode index = MAPPER.createObjectNode();
+      index.put(OUTPUT_INDEX_VERSION_FIELD, OUTPUT_INDEX_VERSION);
+      index.put(
+          OUTPUT_INDEX_WORKING_DIR_FIELD, 
Joiner.on('/').join(stagingRoot.relativize(workingDir)));
+      // Recreated in case it was removed while the server is running, e.g. by 
a manual cleanup.
+      Files.createDirectories(outputIndexDir);
+      // No temp file and rename: the job id is only handed out after this 
returns, so nobody can
+      // read a partially written index, and atomic rename isn't available on 
every shared storage.
+      Files.write(
+          outputIndexFile(jobId),
+          MAPPER.writeValueAsBytes(index),
+          StandardOpenOption.CREATE_NEW,
+          StandardOpenOption.WRITE);
+    } catch (IOException | RuntimeException e) {
+      LOG.warn(
+          "Failed to write the output index of job {}, its output can't be 
retrieved", jobId, e);
+    }
+  }
+
+  @Nullable
+  private Path locateWorkingDir(String jobId) {
+    // The job id becomes a file name, so it must not be able to carry any 
path elements.
+    if (jobId == null || !JOB_ID_PATTERN.matcher(jobId).matches()) {
+      LOG.debug("Job {} is not a job of the local job executor, it has no 
output index", jobId);
+      return null;
+    }
+
+    Path indexFile = outputIndexFile(jobId);
+    byte[] content;
+    try {
+      content = Files.readAllBytes(indexFile);
+    } catch (NoSuchFileException e) {
+      warnOnce(
+          missingOutputIndexWarned,
+          "No output index found for job {} under {}, so its output can't be 
retrieved. The job "
+              + "may have been submitted before this Gravitino version, or run 
on another "
+              + "Gravitino server that doesn't share gravitino.job.stagingDir 
with this one. In a "
+              + "multi-node deployment, gravitino.job.stagingDir must be 
shared by all servers.",
+          jobId,
+          outputIndexDir);
+      return null;
+    } catch (IOException e) {
+      // Same as reading the output itself: an unexpected I/O failure must not 
be reported as
+      // "no output".
+      throw new RuntimeException("Failed to read the output index of job " + 
jobId, e);

Review Comment:
   Yes, that's the intended contract; added a line to 
manage-jobs-in-gravitino.md saying an index that exists but can't be read fails 
the request instead of returning empty output (66a58a972).
   
   _🤖 Addressed by [Claude Code](https://claude.com/claude-code)_



##########
core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutor.java:
##########
@@ -103,6 +168,31 @@ public void initialize(Map<String, String> configs) {
     this.ownedJobIdPrefix = LOCAL_JOB_PREFIX + executorId + "-";
     LOG.info("Initializing local job executor with executor id {}", 
executorId);
 
+    // Validated before any thread is started, so a failed initialization 
leaks nothing.
+    String stagingDir = configs.get(STAGING_DIR);
+    Preconditions.checkArgument(
+        StringUtils.isNotBlank(stagingDir),
+        "The job staging directory must be set for the local job executor");

Review Comment:
   Added a note to the local job executor section of 
manage-jobs-in-gravitino.md that it always uses `gravitino.job.stagingDir` and 
ignores `gravitino.jobExecutor.local.stagingDir`, and a line to the 
LocalJobExecutor javadoc that `initialize` requires `stagingDir` (66a58a972); 
the `STAGING_DIR` javadoc already says the value is set by Gravitino.
   
   _🤖 Addressed by [Claude Code](https://claude.com/claude-code)_



##########
core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutor.java:
##########
@@ -65,6 +92,34 @@ public class LocalJobExecutor implements JobExecutor {
 
   private static final long UNEXPIRED_TIME_IN_MS = -1L;
 
+  // Contains '.' and '-', which MetalakeNormalizeDispatcher rejects in 
metalake names on create and

Review Comment:
   Keeping it to the inline comment: a collision needs a metalake created 
before #3510 named exactly `.job-output-index`, and the cleanup already handles 
it without data loss, so a docs note would mostly add noise.



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