jerryshao commented on code in PR #13373:
URL: https://github.com/apache/gravitino/pull/13373#discussion_r4068880799
##########
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:
The check does exist: `MetalakeNormalizeDispatcher` validates names against
`^\w[\w]{0,63}$` on create (line 79) and rename (line 91). Since it was only
added in #3510, a legacy metalake could still collide, so the cleanup now skips
anything but regular files, and the comment says where the check lives. Fixed
in 16bbe1312.
##########
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:
Leaving this one: stdout and stderr are separate SPI calls, so sharing the
lookup would need an SPI change or a cache. The index is a few dozen bytes, so
the second read isn't worth that complexity.
##########
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:
Moved the index cleanup to a scheduler of its own in 16bbe1312.
--
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]