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


##########
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:
   Nit on behavior (not blocking): a missing index returns empty output, but an 
I/O error reading the index throws `RuntimeException` and can fail the whole 
`getJob(includeOutput=true)` call. That seems intentional 
(`testGetJobOutputFailsWhenOutputIndexCannotBeRead`), but worth calling out for 
operators — a transient shared-storage glitch surfaces as an API error rather 
than empty stdout. If that's the desired contract, maybe one line in 
`manage-jobs-in-gravitino.md` would help.



##########
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:
   Good edge-case comment. If a legacy metalake name ever collides with 
`.job-output-index`, cleanup skips it rather than failing — fine. If you think 
operators might hit this, a single sentence in the docs under the 
index-directory section would be enough; otherwise this inline comment is 
sufficient.



##########
docs/manage-jobs-in-gravitino.md:
##########
@@ -276,19 +277,35 @@ curl -X GET -H "Accept: 
application/vnd.gravitino.v1+json" \
   
"http://localhost:8090/api/metalakes/example/jobs/runs/{job_id}?includeOutput=true&outputMaxLines=50&outputMaxBytes=8192";
 ```
 
+:::caution
+When multiple Gravitino servers share the same metadata store, 
`gravitino.job.stagingDir` must be on
+storage shared by all servers (for example an NFS mount) for the local job 
executor to return a
+job's output from any server. The servers may mount it at different paths. 
Otherwise, only the
+server that ran a job can return its output, and the other servers return 
empty output rather than
+an error.
+:::
+
+The local job executor finds a job's output through a small index file it 
writes to
+`<gravitino.job.stagingDir>/.job-output-index` when the job is submitted:
+
+- Jobs submitted before Gravitino 2.0.0, or during a rolling upgrade by a 
server that isn't upgraded

Review Comment:
   Rolling-upgrade limitation is documented here — helpful. LGTM on making the 
shared-staging-dir requirement prominent; that's the main operational 
prerequisite for this fix.



##########
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:
   `initialize()` now requires `stagingDir` (breaking for anyone constructing 
`LocalJobExecutor` directly). `JobExecutorFactory` injects it correctly and the 
new tests cover that — consider a one-line note in `LocalJobExecutorConfigs` / 
server config docs that the value comes from `gravitino.job.stagingDir` and 
must not be set under the executor prefix (you already test the override).



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