jerryshao commented on code in PR #13373:
URL: https://github.com/apache/gravitino/pull/13373#discussion_r4070108760
##########
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:
Yes, that's the intended contract; added a line to
manage-jobs-in-gravitino.md saying an index that exists but can't be read fails
the request instead of returning empty output (66a58a972).
_🤖 Addressed by [Claude Code](https://claude.com/claude-code)_
##########
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:
Added a note to the local job executor section of
manage-jobs-in-gravitino.md that it always uses `gravitino.job.stagingDir` and
ignores `gravitino.jobExecutor.local.stagingDir`, and a line to the
LocalJobExecutor javadoc that `initialize` requires `stagingDir` (66a58a972);
the `STAGING_DIR` javadoc already says the value is set by Gravitino.
_🤖 Addressed by [Claude Code](https://claude.com/claude-code)_
##########
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:
Keeping it to the inline comment: a collision needs a metalake created
before #3510 named exactly `.job-output-index`, and the cleanup already handles
it without data loss, so a docs note would mostly add noise.
--
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]