lasdf1234 commented on code in PR #13373:
URL: https://github.com/apache/gravitino/pull/13373#discussion_r4070094870
##########
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:
Nit on behavior (not blocking): a missing index returns empty output, but an
I/O error reading the index throws `RuntimeException` and can fail the whole
`getJob(includeOutput=true)` call. That seems intentional
(`testGetJobOutputFailsWhenOutputIndexCannotBeRead`), but worth calling out for
operators — a transient shared-storage glitch surfaces as an API error rather
than empty stdout. If that's the desired contract, maybe one line in
`manage-jobs-in-gravitino.md` would help.
##########
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:
Good edge-case comment. If a legacy metalake name ever collides with
`.job-output-index`, cleanup skips it rather than failing — fine. If you think
operators might hit this, a single sentence in the docs under the
index-directory section would be enough; otherwise this inline comment is
sufficient.
##########
docs/manage-jobs-in-gravitino.md:
##########
@@ -276,19 +277,35 @@ curl -X GET -H "Accept:
application/vnd.gravitino.v1+json" \
"http://localhost:8090/api/metalakes/example/jobs/runs/{job_id}?includeOutput=true&outputMaxLines=50&outputMaxBytes=8192"
```
+:::caution
+When multiple Gravitino servers share the same metadata store,
`gravitino.job.stagingDir` must be on
+storage shared by all servers (for example an NFS mount) for the local job
executor to return a
+job's output from any server. The servers may mount it at different paths.
Otherwise, only the
+server that ran a job can return its output, and the other servers return
empty output rather than
+an error.
+:::
+
+The local job executor finds a job's output through a small index file it
writes to
+`<gravitino.job.stagingDir>/.job-output-index` when the job is submitted:
+
+- Jobs submitted before Gravitino 2.0.0, or during a rolling upgrade by a
server that isn't upgraded
Review Comment:
Rolling-upgrade limitation is documented here — helpful. LGTM on making the
shared-staging-dir requirement prominent; that's the main operational
prerequisite for this fix.
##########
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:
`initialize()` now requires `stagingDir` (breaking for anyone constructing
`LocalJobExecutor` directly). `JobExecutorFactory` injects it correctly and the
new tests cover that — consider a one-line note in `LocalJobExecutorConfigs` /
server config docs that the value comes from `gravitino.job.stagingDir` and
must not be set under the executor prefix (you already test the override).
--
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]