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


##########
core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutor.java:
##########
@@ -65,6 +91,33 @@ public class LocalJobExecutor implements JobExecutor {
 
   private static final long UNEXPIRED_TIME_IN_MS = -1L;
 
+  // Contains '.' and '-', which a metalake name can't, so it never collides 
with a metalake's
+  // staging directory.
+  private static final String OUTPUT_INDEX_DIR_NAME = ".job-output-index";

Review Comment:
   [Question] The comment justifies the directory name with metalake names not 
being able to contain `.` or `-`, but I could not find a check that enforces 
that. `MetalakeCreateRequest.validate()` 
(common/src/main/java/org/apache/gravitino/dto/requests/MetalakeCreateRequest.java:75-78)
 only requires a non-blank name, and `MetadataObjects.checkName` 
(api/src/main/java/org/apache/gravitino/MetadataObjects.java:212-217) only 
rejects `*`; there is no name pattern for metalakes in 
`NameIdentifierUtil.checkMetalake` either 
(core/src/main/java/org/apache/gravitino/utils/NameIdentifierUtil.java:453-458).
   
   A metalake named `.job-output-index` would therefore put its job staging 
directories inside the index directory. The damage looks contained - 
`cleanupOutputIndexes` globs `local-job-*.json`, and a *directory* matching 
that glob makes `Files.readAllBytes` throw, which is caught and logged rather 
than deleted - so this is about the justification rather than a data-loss path. 
Consider softening the wording to why it is safe even on a collision, or 
skipping non-regular files explicitly in the cleanup.
   
   Verified by: read all three validation sites in this run and traced what a 
colliding metalake name would do through `cleanupOutputIndexes` 
(LocalJobExecutor.java:504-536).



##########
core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutor.java:
##########
@@ -175,9 +259,13 @@ public void initialize(Map<String, String> configs) {
         jobStatusCleanupIntervalInMs,
         jobStatusCleanupIntervalInMs,
         TimeUnit.MILLISECONDS);
+    jobStatusCleanupExecutor.scheduleWithFixedDelay(
+        this::cleanupOutputIndexes,
+        OUTPUT_INDEX_CLEANUP_INTERVAL_IN_MS,
+        OUTPUT_INDEX_CLEANUP_INTERVAL_IN_MS,
+        TimeUnit.MILLISECONDS);

Review Comment:
   [Nit] The index cleanup runs on `jobStatusCleanupExecutor`, the same single 
thread as `cleanupJobStatus` (created at LocalJobExecutor.java:249-256). 
`.job-output-index` holds one file per job for as long as the staging 
directories live (`gravitino.job.stagingDirKeepTimeInMs`, 7 days by default), 
and each entry costs a `getLastModifiedTime`, a full read and possibly a 
`notExists` on shared storage (LocalJobExecutor.java:504-536), so on a busy 
deployment one scan can run long enough to postpone the in-memory status 
cleanup queued behind it. A scheduler of its own would keep the two 
independent; if the delay is considered acceptable, a line saying so would help.
   
   Verified by: both tasks are scheduled on the same single-thread executor 
here, and `cleanupOutputIndexes` iterates every index file sequentially.



##########
core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutor.java:
##########
@@ -404,29 +490,197 @@ 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 {
+          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);

Review Comment:
   [Nit] The working directory is resolved per stream, so a single 
`getJob(includeOutput=true)` reads the index file twice: `JobManager.getJob` 
calls `getJobStdout` and `getJobStderr` back to back 
(core/src/main/java/org/apache/gravitino/job/JobManager.java:468-471), and each 
goes through `locateWorkingDir`. On the shared (typically NFS) storage this 
feature targets, that is an avoidable extra round trip per request; resolving 
once and reading both files from the same `Path` would halve it.
   
   Verified by: read `JobManager.getJob` and both `getJobStdout`/`getJobStderr` 
implementations (LocalJobExecutor.java:370-378) in this run.



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