yuqi1129 commented on code in PR #13373:
URL: https://github.com/apache/gravitino/pull/13373#discussion_r4070107522
##########
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:
The containment check in `parseOutputIndex` is lexical, but `readLastLines`
follows symlinks. A local shell job can replace its `output.log` with an
absolute symlink after the process starts; when another server sharing the
staging directory handles `getJob(includeOutput=true)`, it resolves that
symlink against its own filesystem and can return a file outside staging. I
reproduced the path behavior with Java (`startsWith(stagingRoot)` true, read
from an outside file). Could we reject symlinks or validate the real output
path against the real staging root before reading, and add a cross-node symlink
test?
##########
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:
This exact class-name check misses a configured subclass of
`LocalJobExecutor` that inherits `initialize()`. Such an executor no longer
receives `stagingDir`, so the new non-blank check fails during server startup.
Could we inject the staging directory based on the instantiated executor type
(for example, `instanceof LocalJobExecutor`) and cover a subclass in
`TestJobExecutorFactory`?
--
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]