jerryshao commented on code in PR #13373:
URL: https://github.com/apache/gravitino/pull/13373#discussion_r4070170696
##########
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);
Review Comment:
Good catch, fixed in aead588a7: output is only read from a working directory
whose real path equals the real staging root plus the indexed relative path (so
a symlinked working directory, or a symlinked directory above it such as one
pointing at another metalake's job, is rejected), and output.log/error.log are
opened with NOFOLLOW_LINKS. Added
TestLocalJobExecutor#testGetJobOutputDoesNotFollowSymlinks (file, working dir
and parent dir symlinks, read from a second executor) and
TestJobManagerMultiNode#testGetJobOutputDoesNotFollowSymlinkFromAnotherNode,
where the job replaces its output.log with an absolute symlink and node B reads
it.
_🤖 Addressed by [Claude Code](https://claude.com/claude-code)_
##########
core/src/main/java/org/apache/gravitino/job/JobExecutorFactory.java:
##########
@@ -60,7 +61,13 @@ public static JobExecutor create(Config config) {
jobExecutorName);
Map<String, String> configs =
- config.getConfigsWithPrefix(JOB_EXECUTOR_CONF_PREFIX + jobExecutorName
+ ".");
+ Maps.newHashMap(
+ config.getConfigsWithPrefix(JOB_EXECUTOR_CONF_PREFIX +
jobExecutorName + "."));
+ if (LocalJobExecutor.class.getCanonicalName().equals(clzName)) {
Review Comment:
Done in aead588a7: the staging directory is now injected when the
instantiated executor is a LocalJobExecutor (instanceof), and
TestJobExecutorFactory#testLocalJobExecutorSubclassUsesJobStagingDir covers a
subclass that inherits initialize().
_🤖 Addressed by [Claude Code](https://claude.com/claude-code)_
--
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]