jerryshao commented on code in PR #13373:
URL: https://github.com/apache/gravitino/pull/13373#discussion_r4068417711
##########
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:
[Question] The comment justifies the directory name with metalake names not
being able to contain `.` or `-`, but I could not find a check that enforces
that. `MetalakeCreateRequest.validate()`
(common/src/main/java/org/apache/gravitino/dto/requests/MetalakeCreateRequest.java:75-78)
only requires a non-blank name, and `MetadataObjects.checkName`
(api/src/main/java/org/apache/gravitino/MetadataObjects.java:212-217) only
rejects `*`; there is no name pattern for metalakes in
`NameIdentifierUtil.checkMetalake` either
(core/src/main/java/org/apache/gravitino/utils/NameIdentifierUtil.java:453-458).
A metalake named `.job-output-index` would therefore put its job staging
directories inside the index directory. The damage looks contained -
`cleanupOutputIndexes` globs `local-job-*.json`, and a *directory* matching
that glob makes `Files.readAllBytes` throw, which is caught and logged rather
than deleted - so this is about the justification rather than a data-loss path.
Consider softening the wording to why it is safe even on a collision, or
skipping non-regular files explicitly in the cleanup.
Verified by: read all three validation sites in this run and traced what a
colliding metalake name would do through `cleanupOutputIndexes`
(LocalJobExecutor.java:504-536).
##########
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:
[Nit] The index cleanup runs on `jobStatusCleanupExecutor`, the same single
thread as `cleanupJobStatus` (created at LocalJobExecutor.java:249-256).
`.job-output-index` holds one file per job for as long as the staging
directories live (`gravitino.job.stagingDirKeepTimeInMs`, 7 days by default),
and each entry costs a `getLastModifiedTime`, a full read and possibly a
`notExists` on shared storage (LocalJobExecutor.java:504-536), so on a busy
deployment one scan can run long enough to postpone the in-memory status
cleanup queued behind it. A scheduler of its own would keep the two
independent; if the delay is considered acceptable, a line saying so would help.
Verified by: both tasks are scheduled on the same single-thread executor
here, and `cleanupOutputIndexes` iterates every index file sequentially.
##########
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:
[Nit] The working directory is resolved per stream, so a single
`getJob(includeOutput=true)` reads the index file twice: `JobManager.getJob`
calls `getJobStdout` and `getJobStderr` back to back
(core/src/main/java/org/apache/gravitino/job/JobManager.java:468-471), and each
goes through `locateWorkingDir`. On the shared (typically NFS) storage this
feature targets, that is an avoidable extra round trip per request; resolving
once and reading both files from the same `Path` would halve it.
Verified by: read `JobManager.getJob` and both `getJobStdout`/`getJobStderr`
implementations (LocalJobExecutor.java:370-378) in this run.
--
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]