This is an automated email from the ASF dual-hosted git repository.
yuqi1129 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/main by this push:
new 1c70cdcc6a [#13372] fix(core): Make job staging directories
independent of template and metalake names (#13425)
1c70cdcc6a is described below
commit 1c70cdcc6acfea1b816684864b1784d34421b80a
Author: Jerry Shao <[email protected]>
AuthorDate: Wed Sep 23 21:52:09 2026 +0800
[#13372] fix(core): Make job staging directories independent of template
and metalake names (#13425)
### What changes were proposed in this pull request?
A job's staging directory is now
`<gravitino.job.stagingDir>/job-runs/job-<id>`, derived from the job id
alone, instead of `<stagingDir>/<metalake>/<template>/job-<id>`.
Both `cleanUpStagingDirs()` and `deleteJobTemplate()` delete a job's
directory through one helper. It first tries the new location. If that
doesn't exist, the job was submitted by an earlier version, and its
directory is looked up by id under `<stagingDir>/*/*/job-<id>` instead
of being rebuilt from the current names. This also removes directories
that were already leaked by an earlier rename. Symbolic links and hidden
entries such as `.job-output-index` are skipped, and an unreadable
directory (e.g. `lost+found`) doesn't stop the lookup.
Also updates the staging paths in the table maintenance docs, and fixes
the `JobIT` check for a rejected job's staging directory, which passed
vacuously with the new layout.
### Why are the changes needed?
The cleanup rebuilt the staging path from the current template and
metalake names. After a template or metalake rename, the path didn't
exist, nothing was deleted, and the job entity was removed anyway, so
the directory (executables, scripts, jars, output logs) leaked forever.
The metalake rename case is verified by a new test.
Fix: #13372
### Does this PR introduce _any_ user-facing change?
The layout of the job staging directory changes as described above.
Directories of jobs submitted before the upgrade are still cleaned up.
No API or configuration changes.
During a rolling upgrade with a shared staging directory, a server
running the old version can't find the directories of jobs submitted by
upgraded servers, so the jobs it cleans up leave their directories
behind.
### How was this patch tested?
- New tests in `TestJobManagerMultiNode` (real JDBC backend and
`LocalJobExecutor`) for template rename + expiry, repeated template
renames + `deleteJobTemplate`, and metalake rename + expiry. All three
fail on main.
- New unit tests in `TestJobManager` for the legacy directory lookup,
including symlinks, hidden directories, same-name files and unreadable
directories.
- `./gradlew :core:test -PskipITs --tests 'org.apache.gravitino.job.*'`,
and `JobIT` in embedded mode.
🤖 Generated with [Claude Code](https://claude.com/claude-code)
---------
Co-authored-by: Claude Opus 5 <[email protected]>
---
.../gravitino/client/integration/test/JobIT.java | 16 +-
.../java/org/apache/gravitino/job/JobManager.java | 140 +++++++++---
.../org/apache/gravitino/job/TestJobManager.java | 237 +++++++++++++++++++--
.../gravitino/job/TestJobManagerMultiNode.java | 90 +++++++-
.../gravitino/job/local/TestLocalJobExecutor.java | 13 +-
.../optimizer-cli-reference.md | 2 +-
.../optimizer-troubleshooting.md | 2 +-
7 files changed, 441 insertions(+), 59 deletions(-)
diff --git
a/clients/client-java/src/test/java/org/apache/gravitino/client/integration/test/JobIT.java
b/clients/client-java/src/test/java/org/apache/gravitino/client/integration/test/JobIT.java
index 08f38a9d24..1442354417 100644
---
a/clients/client-java/src/test/java/org/apache/gravitino/client/integration/test/JobIT.java
+++
b/clients/client-java/src/test/java/org/apache/gravitino/client/integration/test/JobIT.java
@@ -19,6 +19,7 @@
package org.apache.gravitino.client.integration.test;
import com.google.common.collect.ImmutableMap;
+import com.google.common.collect.ImmutableSet;
import com.google.common.collect.Lists;
import java.io.File;
import java.io.IOException;
@@ -380,6 +381,8 @@ public class JobIT extends BaseIT {
.withClassName("org.apache.gravitino.test.SparkJob")
.build();
Assertions.assertDoesNotThrow(() ->
metalake.registerJobTemplate(template));
+ File jobRunsDir = new File(testStagingDir, "job-runs");
+ Set<String> jobStagingDirsBefore = listFileNames(jobRunsDir);
// The run request is rejected with the reason instead of being queued and
failing later.
IllegalArgumentException e =
@@ -396,9 +399,11 @@ public class JobIT extends BaseIT {
// No job is created, and the staging directory of the rejected job is
removed.
Assertions.assertTrue(metalake.listJobs(template.name()).isEmpty());
- String[] jobStagingDirs =
- new File(testStagingDir, METALAKE_NAME + File.separator +
template.name()).list();
- Assertions.assertTrue(jobStagingDirs == null || jobStagingDirs.length ==
0);
+ // Directories of earlier jobs may expire meanwhile, but no new one may be
left behind.
+ Set<String> jobStagingDirsAfter = listFileNames(jobRunsDir);
+ Assertions.assertTrue(
+ jobStagingDirsBefore.containsAll(jobStagingDirsAfter),
+ "Staging directories left behind: " + jobStagingDirsAfter);
// Shell jobs are not affected by the missing Spark installation.
JobTemplate shellTemplate =
builder.withName("test_run_shell_without_spark_submit").build();
@@ -671,6 +676,11 @@ public class JobIT extends BaseIT {
return job;
}
+ private static Set<String> listFileNames(File dir) {
+ String[] names = dir.list();
+ return names == null ? Collections.emptySet() : ImmutableSet.copyOf(names);
+ }
+
private String generateTestEntryScript() {
String content =
"#!/bin/bash\n"
diff --git a/core/src/main/java/org/apache/gravitino/job/JobManager.java
b/core/src/main/java/org/apache/gravitino/job/JobManager.java
index 7e67602a43..a2f2069d08 100644
--- a/core/src/main/java/org/apache/gravitino/job/JobManager.java
+++ b/core/src/main/java/org/apache/gravitino/job/JobManager.java
@@ -27,7 +27,11 @@ import com.google.common.base.Preconditions;
import java.io.File;
import java.io.IOException;
import java.net.URI;
+import java.nio.file.DirectoryIteratorException;
+import java.nio.file.DirectoryStream;
import java.nio.file.Files;
+import java.nio.file.LinkOption;
+import java.nio.file.Path;
import java.time.Instant;
import java.util.ArrayList;
import java.util.Arrays;
@@ -84,14 +88,20 @@ public class JobManager implements JobOperationDispatcher {
private static final Pattern PLACEHOLDER_PATTERN =
Pattern.compile("\\{\\{([\\w.-]+)\\}\\}");
- private static final String JOB_STAGING_DIR =
- File.separator
- + "%s"
- + File.separator
- + "%s"
- + File.separator
- + JobHandle.JOB_ID_PREFIX
- + "%s";
+ // A job's staging directory is <stagingDir>/job-runs/job-<id>, derived from
the job id alone:
+ // template and metalake names can change after the job runs, so a path
built from them can't be
+ // rebuilt at cleanup time. The name contains '-', which metalake names
can't, so it never
+ // collides with a metalake directory of the legacy
<stagingDir>/<metalake>/<template>/job-<id>
+ // layout.
+ private static final String JOB_RUNS_DIR_NAME = "job-runs";
+
+ private static final Pattern JOB_DIR_PATTERN =
+ Pattern.compile(Pattern.quote(JobHandle.JOB_ID_PREFIX) + "\\d+");
+
+ // Bounds how deep a legacy staging directory is looked up below a job
template's first path
+ // element. A template name nested deeper than this keeps its directory,
which is better than
+ // walking an unrelated directory tree an operator put in the staging
directory.
+ private static final int LEGACY_JOB_STAGING_DIR_MAX_DEPTH = 16;
private static final long JOB_STAGING_DIR_CLEANUP_MIN_TIME_IN_MS = 600 *
1000L; // 10 minute
@@ -312,16 +322,13 @@ public class JobManager implements JobOperationDispatcher
{
}
// Only remove directories belonging to the observed jobs. A same-name
template can be
- // recreated after the metadata transaction commits, so its parent
directory is not ours to
- // delete.
+ // recreated after the metadata transaction commits, so a legacy template
directory is not ours
+ // to delete.
for (JobEntity job : jobs) {
- String jobStagingPath =
- stagingDir.getAbsolutePath()
- + String.format(JOB_STAGING_DIR, metalake,
job.jobTemplateName(), job.id());
try {
- FileUtils.deleteDirectory(new File(jobStagingPath));
+ deleteJobStagingDir(job);
} catch (IOException e) {
- LOG.error("Failed to delete job staging directory: {}",
jobStagingPath, e);
+ LOG.error("Failed to delete the staging directory of job {}",
job.name(), e);
}
}
@@ -496,10 +503,7 @@ public class JobManager implements JobOperationDispatcher {
// Create staging directory.
long jobId = idGenerator.nextId();
- String jobStagingPath =
- stagingDir.getAbsolutePath()
- + String.format(JOB_STAGING_DIR, metalake, jobTemplateName, jobId);
- File jobStagingDir = new File(jobStagingPath);
+ File jobStagingDir = jobStagingDir(jobId);
try {
Files.createDirectories(jobStagingDir.toPath());
} catch (IOException e) {
@@ -821,15 +825,7 @@ public class JobManager implements JobOperationDispatcher {
try {
entityStore.delete(
NameIdentifierUtil.ofJob(metalake, job.name()),
Entity.EntityType.JOB);
-
- String jobStagingPath =
- stagingDir.getAbsolutePath()
- + String.format(JOB_STAGING_DIR, metalake,
job.jobTemplateName(), job.id());
- File jobStagingDir = new File(jobStagingPath);
- if (jobStagingDir.exists()) {
- FileUtils.deleteDirectory(jobStagingDir);
- LOG.info("Deleted job staging directory {} for job {}",
jobStagingPath, job.name());
- }
+ deleteJobStagingDir(job);
} catch (OptimisticLockException e) {
// Keep the files when deletion loses its CAS. The next cleanup
run re-reads the
// job and checks retention eligibility again; this batch can
process other jobs.
@@ -1108,6 +1104,94 @@ public class JobManager implements
JobOperationDispatcher {
.build();
}
+ @VisibleForTesting
+ File jobStagingDir(long jobId) {
+ return new File(new File(stagingDir, JOB_RUNS_DIR_NAME),
JobHandle.JOB_ID_PREFIX + jobId);
+ }
+
+ /**
+ * Finds the staging directory of a job submitted by an earlier Gravitino
version, which used the
+ * {@code <stagingDir>/<metalake>/<template>/job-<id>} layout. It's looked
up by the job id rather
+ * than rebuilt from the current names, because the template or metalake may
have been renamed
+ * since the job ran. The job id is unique, so there is at most one match in
practice. Called only
+ * for a job without a directory in the current layout, which after an
upgrade are the jobs of the
+ * earlier version until they expire.
+ */
+ @VisibleForTesting
+ List<File> findLegacyJobStagingDirs(long jobId) throws IOException {
+ String jobDirName = JobHandle.JOB_ID_PREFIX + jobId;
+ List<File> legacyJobStagingDirs = new ArrayList<>();
+ try (DirectoryStream<Path> metalakeDirs =
Files.newDirectoryStream(stagingDir.toPath())) {
+ for (Path metalakeDir : metalakeDirs) {
+ String name = metalakeDir.getFileName().toString();
+ // Metalake names can't start with '.', so hidden entries, e.g. the
job output index
+ // directory, are not metalake directories. Symbolic links are never
followed, so nothing
+ // outside the staging directory can be deleted.
+ if (name.equals(JOB_RUNS_DIR_NAME)
+ || name.startsWith(".")
+ || !Files.isDirectory(metalakeDir, LinkOption.NOFOLLOW_LINKS)) {
+ continue;
+ }
+
+ // An unreadable directory, e.g. "lost+found" when the staging
directory is the root of a
+ // file system, must not prevent finding the job under the other
directories.
+ try (DirectoryStream<Path> templateDirs =
Files.newDirectoryStream(metalakeDir)) {
+ for (Path templateDir : templateDirs) {
+ if (Files.isDirectory(templateDir, LinkOption.NOFOLLOW_LINKS)) {
+ collectLegacyJobStagingDirs(
+ templateDir, jobDirName, LEGACY_JOB_STAGING_DIR_MAX_DEPTH,
legacyJobStagingDirs);
+ }
+ }
+ } catch (IOException | DirectoryIteratorException e) {
+ LOG.warn(
+ "Failed to look up the legacy staging directory of job {} under
{}, skip it",
+ jobDirName,
+ metalakeDir,
+ e);
+ }
+ }
+ }
+ return legacyJobStagingDirs;
+ }
+
+ // The legacy layout put the template name in the path as is, and a template
name may contain
+ // '/', so the job directory can be nested deeper than
<metalake>/<template>/job-<id>, e.g. under
+ // a template named "team/etl". The staging directory of another job is
never descended into: its
+ // contents are that job's own files.
+ private static void collectLegacyJobStagingDirs(
+ Path dir, String jobDirName, int remainingDepth, List<File>
legacyJobStagingDirs)
+ throws IOException {
+ try (DirectoryStream<Path> children = Files.newDirectoryStream(dir)) {
+ for (Path child : children) {
+ if (!Files.isDirectory(child, LinkOption.NOFOLLOW_LINKS)) {
+ continue;
+ }
+
+ String name = child.getFileName().toString();
+ if (name.equals(jobDirName)) {
+ legacyJobStagingDirs.add(child.toFile());
+ } else if (remainingDepth > 0 &&
!JOB_DIR_PATTERN.matcher(name).matches()) {
+ collectLegacyJobStagingDirs(child, jobDirName, remainingDepth - 1,
legacyJobStagingDirs);
+ }
+ }
+ }
+ }
+
+ private void deleteJobStagingDir(JobEntity job) throws IOException {
+ File jobStagingDir = jobStagingDir(job.id());
+ if (jobStagingDir.exists()) {
+ FileUtils.deleteDirectory(jobStagingDir);
+ LOG.info("Deleted job staging directory {} for job {}", jobStagingDir,
job.name());
+ return;
+ }
+
+ for (File legacyJobStagingDir : findLegacyJobStagingDirs(job.id())) {
+ FileUtils.deleteDirectory(legacyJobStagingDir);
+ LOG.info(
+ "Deleted legacy job staging directory {} for job {}",
legacyJobStagingDir, job.name());
+ }
+ }
+
private void deleteStagingDirOfUnsubmittedJob(File jobStagingDir, long
jobId) {
// The job is not tracked by any job entity, so the periodic cleanup will
never remove its
// staging directory. A cleanup failure must not mask the original
submission failure.
diff --git a/core/src/test/java/org/apache/gravitino/job/TestJobManager.java
b/core/src/test/java/org/apache/gravitino/job/TestJobManager.java
index 1091b38fec..d19c274c02 100644
--- a/core/src/test/java/org/apache/gravitino/job/TestJobManager.java
+++ b/core/src/test/java/org/apache/gravitino/job/TestJobManager.java
@@ -33,6 +33,7 @@ import static org.mockito.Mockito.when;
import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableMap;
+import com.google.common.collect.ImmutableSet;
import com.google.common.collect.Lists;
import com.sun.net.httpserver.HttpServer;
import java.io.File;
@@ -53,6 +54,7 @@ import java.util.UUID;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.function.Function;
+import java.util.stream.Collectors;
import javax.annotation.Nullable;
import org.apache.commons.io.FileUtils;
import org.apache.commons.lang3.ArrayUtils;
@@ -405,10 +407,7 @@ public class TestJobManager {
.delete(
NameIdentifierUtil.ofJobTemplate(metalake, "shell_job"),
Entity.EntityType.JOB_TEMPLATE);
- File directory =
- new File(
- testStagingDir,
- metalake + File.separator + "shell_job" + File.separator +
finishedJob.name());
+ File directory = jobManager.jobStagingDir(finishedJob.id());
Assertions.assertTrue(directory.mkdirs() || directory.isDirectory());
File artifact = new File(directory, "artifact");
Assertions.assertTrue(artifact.createNewFile());
@@ -428,9 +427,7 @@ public class TestJobManager {
@Test
public void testDeletePreservesReplacementStaging() throws IOException {
doReturn(Collections.emptyList()).when(jobManager).listJobs(metalake,
Optional.of("shell_job"));
- File replacementDir =
- new File(
- testStagingDir, metalake + File.separator + "shell_job" +
File.separator + "job_999");
+ File replacementDir = jobManager.jobStagingDir(999L);
File replacementArtifact = new File(replacementDir, "new-job-artifact");
when(entityStore.delete(
NameIdentifierUtil.ofJobTemplate(metalake, "shell_job"),
@@ -455,10 +452,7 @@ public class TestJobManager {
doReturn(Collections.singletonList(finishedJob))
.when(jobManager)
.listJobs(metalake, Optional.of("shell_job"));
- File directory =
- new File(
- testStagingDir,
- metalake + File.separator + "shell_job" + File.separator +
finishedJob.name());
+ File directory = jobManager.jobStagingDir(finishedJob.id());
Assertions.assertTrue(directory.mkdirs());
File artifact = new File(directory, "artifact");
Assertions.assertTrue(artifact.createNewFile());
@@ -824,10 +818,9 @@ public class TestJobManager {
// No job entity is registered and the staging directory of the rejected
job is removed.
verify(entityStore, never()).put(any(JobEntity.class), anyBoolean());
- File templateStagingDir =
- new File(testStagingDir, metalake + File.separator +
shellJobTemplate.name());
- String[] jobStagingDirs = templateStagingDir.list();
- Assertions.assertTrue(jobStagingDirs == null || jobStagingDirs.length ==
0);
+ File jobRunsDir = jobManager.jobStagingDir(0L).getParentFile();
+ Assertions.assertTrue(jobRunsDir.isDirectory(), "The job staging directory
was never created");
+ Assertions.assertArrayEquals(new String[0], jobRunsDir.list());
}
@Test
@@ -1639,7 +1632,7 @@ public class TestJobManager {
for (JobEntity job : ImmutableList.of(queuedJob, startedJob,
cancellingJob, activeJob)) {
stubEntityStoreUpdateToApply(job, job);
}
- File jobStagingDir = new File(testStagingDir, metalake + "/shell_job/" +
startedJob.name());
+ File jobStagingDir = jobManager.jobStagingDir(startedJob.id());
Assertions.assertTrue(jobStagingDir.mkdirs());
long beforeCleanUp = System.currentTimeMillis();
@@ -1744,8 +1737,8 @@ public class TestJobManager {
.thenThrow(new OptimisticLockException("job changed"))
.thenReturn(true);
when(entityStore.delete(otherIdent,
Entity.EntityType.JOB)).thenReturn(true);
- File conflictedDir = new File(testStagingDir, metalake + "/shell_job/" +
conflicted.name());
- File otherDir = new File(testStagingDir, metalake + "/shell_job/" +
other.name());
+ File conflictedDir = jobManager.jobStagingDir(conflicted.id());
+ File otherDir = jobManager.jobStagingDir(other.id());
Assertions.assertTrue(conflictedDir.mkdirs());
Assertions.assertTrue(otherDir.mkdirs());
File artifact = new File(conflictedDir, "artifact");
@@ -1760,6 +1753,214 @@ public class TestJobManager {
verify(entityStore, times(1)).delete(otherIdent, Entity.EntityType.JOB);
}
+ @Test
+ public void testCleanUpStagingDirsDeletesLegacyStagingDirOfRenamedTemplate()
throws IOException {
+ // The job ran from template "shell_job" before an upgrade, in the legacy
layout, and the
+ // template was renamed afterwards, so the job now reports the new name.
+ JobEntity job = expiredJob();
+ JobEntity renamedJob =
+ JobEntity.builder()
+ .withId(job.id())
+ .withJobExecutionId(job.jobExecutionId())
+ .withNamespace(job.namespace())
+ .withJobTemplateName("renamed_shell_job")
+ .withStartedAt(job.startedAt())
+ .withFinishedAt(job.finishedAt())
+ .withStatus(job.status())
+ .withAuditInfo(job.auditInfo())
+ .build();
+ mockListActiveJobs(renamedJob);
+ when(entityStore.delete(NameIdentifierUtil.ofJob(metalake, job.name()),
Entity.EntityType.JOB))
+ .thenReturn(true);
+ File legacyDir = new File(testStagingDir, metalake + "/shell_job/" +
job.name());
+ Assertions.assertTrue(legacyDir.mkdirs());
+ Assertions.assertTrue(new File(legacyDir, "output.log").createNewFile());
+ File otherJobDir = new File(testStagingDir, metalake + "/shell_job/job-" +
(job.id() + 1));
+ Assertions.assertTrue(otherJobDir.mkdirs());
+
+ Assertions.assertDoesNotThrow(() -> jobManager.cleanUpStagingDirs());
+
+ Assertions.assertFalse(legacyDir.exists());
+ Assertions.assertTrue(otherJobDir.isDirectory());
+ }
+
+ @Test
+ public void testDeleteJobTemplateDeletesLegacyStagingDirOfRenamedMetalake()
throws IOException {
+ // The job ran in the legacy layout under the metalake's former name.
+ JobEntity job = expiredJob();
+
doReturn(Collections.singletonList(job)).when(jobManager).listJobs(metalake,
Optional.of("t"));
+ doReturn(true)
+ .when(entityStore)
+ .delete(NameIdentifierUtil.ofJobTemplate(metalake, "t"),
Entity.EntityType.JOB_TEMPLATE);
+ File legacyDir = new File(testStagingDir, "old_metalake_name/t/" +
job.name());
+ Assertions.assertTrue(legacyDir.mkdirs());
+
+ Assertions.assertTrue(jobManager.deleteJobTemplate(metalake, "t"));
+
+ Assertions.assertFalse(legacyDir.exists());
+ // Only the job's own directory is deleted, never the template or metalake
directory.
+ Assertions.assertTrue(legacyDir.getParentFile().isDirectory());
+ }
+
+ @Test
+ public void testDeleteJobStagingDirPrefersCurrentLayout() throws IOException
{
+ JobEntity job = expiredJob();
+
doReturn(Collections.singletonList(job)).when(jobManager).listJobs(metalake,
Optional.of("t"));
+ doReturn(true)
+ .when(entityStore)
+ .delete(NameIdentifierUtil.ofJobTemplate(metalake, "t"),
Entity.EntityType.JOB_TEMPLATE);
+ File jobStagingDir = jobManager.jobStagingDir(job.id());
+ Assertions.assertTrue(jobStagingDir.mkdirs());
+
+ Assertions.assertTrue(jobManager.deleteJobTemplate(metalake, "t"));
+
+ Assertions.assertFalse(jobStagingDir.exists());
+ verify(jobManager, never()).findLegacyJobStagingDirs(job.id());
+ }
+
+ @Test
+ public void testFindLegacyJobStagingDirs() throws IOException {
+ long jobId = idGenerator.nextId();
+ String jobDirName = JobHandle.JOB_ID_PREFIX + jobId;
+ File stagingDir = new File(testStagingDir).getAbsoluteFile();
+ File legacyDir = new File(stagingDir, metalake + "/shell_job/" +
jobDirName);
+ Assertions.assertTrue(legacyDir.mkdirs());
+
+ // Not the job's directory: a file of the same name, a directory at the
wrong depth, a hidden
+ // top-level directory and the directory of the current layout.
+ File sameNameFile = new File(stagingDir, metalake + "/other_job/" +
jobDirName);
+ Assertions.assertTrue(sameNameFile.getParentFile().mkdirs());
+ Assertions.assertTrue(sameNameFile.createNewFile());
+ Assertions.assertTrue(new File(stagingDir, metalake + "/" +
jobDirName).mkdirs());
+ Assertions.assertTrue(new File(stagingDir, ".job-output-index/x/" +
jobDirName).mkdirs());
+ Assertions.assertTrue(jobManager.jobStagingDir(jobId).mkdirs());
+ Assertions.assertTrue(
+ new File(jobManager.jobStagingDir(jobId).getParentFile(), "x/" +
jobDirName).mkdirs());
+ Assertions.assertTrue(new File(stagingDir, "a_file").createNewFile());
+
+ // Symbolic links are never followed, so nothing outside the staging
directory is found.
+ File outsideDir =
Files.createTempDirectory("gravitino-test-outside-staging").toFile();
+ try {
+ Assertions.assertTrue(new File(outsideDir, "t/" + jobDirName).mkdirs());
+ Files.createSymbolicLink(
+ new File(stagingDir, "linked_metalake").toPath(),
outsideDir.toPath());
+ Assertions.assertTrue(new File(stagingDir, "metalake_2").mkdirs());
+ Files.createSymbolicLink(
+ new File(stagingDir, "metalake_2/linked_template").toPath(),
+ new File(outsideDir, "t").toPath());
+ File linkedJobTemplateDir = new File(stagingDir, metalake +
"/linked_job");
+ Assertions.assertTrue(linkedJobTemplateDir.mkdirs());
+ Files.createSymbolicLink(
+ new File(linkedJobTemplateDir, jobDirName).toPath(),
+ new File(outsideDir, "t/" + jobDirName).toPath());
+ } catch (IOException e) {
+ FileUtils.deleteDirectory(outsideDir);
+ throw e;
+ }
+
+ try {
+ Assertions.assertEquals(
+ ImmutableList.of(legacyDir.getAbsolutePath()),
+ jobManager.findLegacyJobStagingDirs(jobId).stream()
+ .map(File::getAbsolutePath)
+ .collect(Collectors.toList()));
+ Assertions.assertTrue(
+ jobManager.findLegacyJobStagingDirs(idGenerator.nextId()).isEmpty(),
+ "No directory of an unknown job");
+ } finally {
+ FileUtils.deleteDirectory(outsideDir);
+ }
+ }
+
+ @Test
+ public void
testCleanUpStagingDirsDeletesLegacyStagingDirOfNestedTemplateName()
+ throws IOException {
+ // Template names may contain '/', which the legacy layout put in the path
as is, so the job
+ // directory is nested deeper. The template was never renamed: this is an
upgraded job.
+ JobEntity job = expiredJob();
+ JobEntity nestedTemplateJob =
+ JobEntity.builder()
+ .withId(job.id())
+ .withJobExecutionId(job.jobExecutionId())
+ .withNamespace(job.namespace())
+ .withJobTemplateName("team/etl")
+ .withStartedAt(job.startedAt())
+ .withFinishedAt(job.finishedAt())
+ .withStatus(job.status())
+ .withAuditInfo(job.auditInfo())
+ .build();
+ mockListActiveJobs(nestedTemplateJob);
+ when(entityStore.delete(NameIdentifierUtil.ofJob(metalake, job.name()),
Entity.EntityType.JOB))
+ .thenReturn(true);
+ File legacyDir = new File(testStagingDir, metalake + "/team/etl/" +
job.name());
+ Assertions.assertTrue(legacyDir.mkdirs());
+
+ Assertions.assertDoesNotThrow(() -> jobManager.cleanUpStagingDirs());
+
+ Assertions.assertFalse(legacyDir.exists());
+ }
+
+ @Test
+ public void testFindLegacyJobStagingDirsOfNestedTemplateNames() throws
IOException {
+ long jobId = idGenerator.nextId();
+ String jobDirName = JobHandle.JOB_ID_PREFIX + jobId;
+ File stagingDir = new File(testStagingDir).getAbsoluteFile();
+ File nestedDir = new File(stagingDir, metalake + "/team/etl/" +
jobDirName);
+ Assertions.assertTrue(nestedDir.mkdirs());
+ File deeplyNestedDir = new File(stagingDir, metalake + "/a/b/c/d/" +
jobDirName);
+ Assertions.assertTrue(deeplyNestedDir.mkdirs());
+
+ // The staging directory of another job is not descended into: a directory
of its own files
+ // named like this job is not this job's staging directory.
+ File otherJobDir =
+ new File(
+ stagingDir,
+ metalake
+ + "/shell_job/"
+ + JobHandle.JOB_ID_PREFIX
+ + (jobId + 1)
+ + File.separator
+ + jobDirName);
+ Assertions.assertTrue(otherJobDir.mkdirs());
+
+ // A template name nested deeper than the lookup goes keeps its directory.
+ StringBuilder tooDeepPath = new StringBuilder(metalake);
+ for (int i = 0; i < 20; i++) {
+ tooDeepPath.append(File.separator).append("d").append(i);
+ }
+ File tooDeepDir = new File(stagingDir, tooDeepPath + File.separator +
jobDirName);
+ Assertions.assertTrue(tooDeepDir.mkdirs());
+
+ Assertions.assertEquals(
+ ImmutableSet.of(nestedDir.getAbsolutePath(),
deeplyNestedDir.getAbsolutePath()),
+ jobManager.findLegacyJobStagingDirs(jobId).stream()
+ .map(File::getAbsolutePath)
+ .collect(Collectors.toSet()));
+ }
+
+ @Test
+ public void testFindLegacyJobStagingDirsSkipsUnreadableDirectory() throws
IOException {
+ long jobId = idGenerator.nextId();
+ File stagingDir = new File(testStagingDir);
+ File legacyDir =
+ new File(stagingDir, metalake + "/shell_job/" +
JobHandle.JOB_ID_PREFIX + jobId);
+ Assertions.assertTrue(legacyDir.mkdirs());
+ // E.g. "lost+found" when the staging directory is the root of a file
system.
+ File unreadableDir = new File(stagingDir, "lost+found");
+ Assertions.assertTrue(unreadableDir.mkdirs());
+ Assertions.assertTrue(unreadableDir.setReadable(false, false));
+
+ try {
+ Assertions.assertEquals(
+ ImmutableList.of(legacyDir.getAbsolutePath()),
+ jobManager.findLegacyJobStagingDirs(jobId).stream()
+ .map(File::getAbsolutePath)
+ .collect(Collectors.toList()));
+ } finally {
+ Assertions.assertTrue(unreadableDir.setReadable(true, false));
+ }
+ }
+
@Test
public void testCleanUpStagingDirs() throws IOException,
InterruptedException {
JobEntity job = newJobEntity("shell_job", JobHandle.Status.STARTED);
diff --git
a/core/src/test/java/org/apache/gravitino/job/TestJobManagerMultiNode.java
b/core/src/test/java/org/apache/gravitino/job/TestJobManagerMultiNode.java
index 9f830bdb4c..f2f50cda2d 100644
--- a/core/src/test/java/org/apache/gravitino/job/TestJobManagerMultiNode.java
+++ b/core/src/test/java/org/apache/gravitino/job/TestJobManagerMultiNode.java
@@ -24,10 +24,14 @@ import com.google.common.collect.Lists;
import java.io.File;
import java.io.IOException;
import java.nio.file.Files;
+import java.nio.file.Path;
import java.time.Instant;
import java.util.Collections;
import java.util.EnumSet;
+import java.util.List;
import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+import java.util.stream.Stream;
import org.apache.commons.io.FileUtils;
import org.apache.commons.lang3.reflect.FieldUtils;
import org.apache.gravitino.Config;
@@ -42,6 +46,7 @@ import org.apache.gravitino.job.local.LocalJobExecutor;
import org.apache.gravitino.job.local.LocalJobExecutorConfigs;
import org.apache.gravitino.lock.LockManager;
import org.apache.gravitino.meta.AuditInfo;
+import org.apache.gravitino.meta.BaseMetalake;
import org.apache.gravitino.meta.JobEntity;
import org.apache.gravitino.meta.JobTemplateEntity;
import org.apache.gravitino.storage.RandomIdGenerator;
@@ -227,7 +232,7 @@ public class TestJobManagerMultiNode extends
TestJDBCBackend {
String newName = ECHO_TEMPLATE + "_renamed";
nodeA.alterJobTemplate(METALAKE, ECHO_TEMPLATE,
JobTemplateChange.rename(newName));
- // The job reports the new template name, while its staging directory
keeps the old one.
+ // The job reports the new template name, which its staging directory
doesn't depend on.
JobEntity jobWithOutput = nodeB.getJob(METALAKE, job.name(), true);
Assertions.assertEquals(newName, jobWithOutput.jobTemplateName());
Assertions.assertEquals(ImmutableList.of("hello b"),
jobWithOutput.stdout());
@@ -235,7 +240,7 @@ public class TestJobManagerMultiNode extends
TestJDBCBackend {
@TestTemplate
public void testGetJobOutputOfTemplateNamedWithSpecialCharacters() throws
IOException {
- // Template names are not restricted, and name a directory level of the
job staging directory.
+ // Template names are not restricted, and must not affect running the job
or reading its output.
for (String templateName : ImmutableList.of("etl job \"v2\" 䏿–‡ #1",
"team/etl")) {
backend.insert(
newScriptJobTemplateEntity(
@@ -300,6 +305,87 @@ public class TestJobManagerMultiNode extends
TestJDBCBackend {
return job;
}
+ @TestTemplate
+ public void testStagingDirIsCleanedUpAfterTemplateRenamed() throws
IOException {
+ JobEntity job = runFinishedJobOnNodeA();
+ Assertions.assertTrue(jobStagingDir(job).isDirectory());
+
+ // Before the fix, the cleanup rebuilt the staging path from the new
template name, missed the
+ // directory and deleted only the job entity, leaking the directory.
+ nodeA.alterJobTemplate(METALAKE, TEMPLATE,
JobTemplateChange.rename("renamed_sleep_job"));
+ Assertions.assertEquals("renamed_sleep_job",
getJob(job.name()).jobTemplateName());
+ moveJobTimestampsBack(job.name());
+ nodeB.cleanUpStagingDirs();
+
+ Assertions.assertFalse(jobExists(job.name()));
+ assertNoStagingDirLeft(job);
+ }
+
+ @TestTemplate
+ public void testStagingDirIsDeletedWithRenamedTemplate() throws IOException {
+ JobEntity job = runFinishedJobOnNodeA();
+
+ nodeA.alterJobTemplate(METALAKE, TEMPLATE,
JobTemplateChange.rename("renamed_sleep_job"));
+ nodeA.alterJobTemplate(
+ METALAKE, "renamed_sleep_job",
JobTemplateChange.rename("renamed_twice_sleep_job"));
+ Assertions.assertTrue(nodeB.deleteJobTemplate(METALAKE,
"renamed_twice_sleep_job"));
+
+ Assertions.assertFalse(jobExists(job.name()));
+ assertNoStagingDirLeft(job);
+ }
+
+ @TestTemplate
+ public void testStagingDirIsCleanedUpAfterMetalakeRenamed() throws
IOException {
+ JobEntity job = runFinishedJobOnNodeA();
+ moveJobTimestampsBack(job.name());
+
+ String newMetalake = METALAKE + "_renamed";
+ entityStore.update(
+ NameIdentifierUtil.ofMetalake(METALAKE),
+ BaseMetalake.class,
+ Entity.EntityType.METALAKE,
+ metalake ->
+ BaseMetalake.builder()
+ .withId(metalake.id())
+ .withName(newMetalake)
+ .withComment(metalake.comment())
+ .withProperties(metalake.properties())
+ .withAuditInfo(metalake.auditInfo())
+ .withVersion(metalake.getVersion())
+ .build());
+ nodeB.cleanUpStagingDirs();
+
+ Assertions.assertThrows(
+ NoSuchJobException.class, () -> nodeB.getJob(newMetalake, job.name(),
false));
+ assertNoStagingDirLeft(job);
+ }
+
+ private JobEntity runFinishedJobOnNodeA() throws IOException {
+ JobEntity job = nodeA.runJob(METALAKE, TEMPLATE,
ImmutableMap.of("seconds", "0"));
+ Awaitility.await()
+ .atMost(1, TimeUnit.MINUTES)
+ .until(() -> executorA.getJobStatus(job.jobExecutionId()) ==
JobHandle.Status.SUCCEEDED);
+ nodeA.pullAndUpdateJobStatus();
+ Assertions.assertEquals(JobHandle.Status.SUCCEEDED,
getJob(job.name()).status());
+ return job;
+ }
+
+ private File jobStagingDir(JobEntity job) {
+ return new File(testDir, "staging/job-runs/" + job.name());
+ }
+
+ // Checks the whole staging directory rather than the expected path, so a
directory left in any
+ // layout is caught.
+ private void assertNoStagingDirLeft(JobEntity job) throws IOException {
+ try (Stream<Path> paths = Files.walk(new File(testDir,
"staging").toPath())) {
+ List<Path> left =
+ paths
+ .filter(path -> path.getFileName().toString().equals(job.name()))
+ .collect(Collectors.toList());
+ Assertions.assertTrue(left.isEmpty(), "Staging directory left behind: "
+ left);
+ }
+ }
+
private JobEntity runLongJobOnNodeA() throws IOException {
JobEntity job = nodeA.runJob(METALAKE, TEMPLATE,
ImmutableMap.of("seconds", "600"));
Awaitility.await()
diff --git
a/core/src/test/java/org/apache/gravitino/job/local/TestLocalJobExecutor.java
b/core/src/test/java/org/apache/gravitino/job/local/TestLocalJobExecutor.java
index 6c5441dd79..ee2d0dae25 100644
---
a/core/src/test/java/org/apache/gravitino/job/local/TestLocalJobExecutor.java
+++
b/core/src/test/java/org/apache/gravitino/job/local/TestLocalJobExecutor.java
@@ -785,22 +785,23 @@ public class TestLocalJobExecutor {
@Test
public void testSubmitJobWritesOutputIndexRelativeToStagingDir() throws
IOException {
- // Laid out like the {metalake}/{template}/job-{id} staging directory of
JobManager.
- File jobDir = new File(workingDir, "metalake/template/job-1");
+ // Laid out like the job-runs/job-{id} staging directory of JobManager.
+ File jobDir = new File(workingDir, "job-runs/job-1");
Assertions.assertTrue(jobDir.mkdirs());
String jobId = runSucceededJob(jobDir);
JsonNode index =
JsonUtils.anyFieldMapper().readTree(outputIndexFile(jobId));
Assertions.assertEquals(1, index.get("version").intValue());
Assertions.assertEquals(
- workingDir.getName() + "/metalake/template/job-1",
index.get("workingDir").textValue());
+ workingDir.getName() + "/job-runs/job-1",
index.get("workingDir").textValue());
}
@Test
public void testOutputIndexKeepsSpecialCharactersInWorkingDir() throws
IOException {
- // Job template names are not restricted, so the staging directory may
contain any character a
- // file name can. They must survive the JSON encoding and the '/'-joining
unchanged. Some of
- // these characters are only valid in POSIX file names, like the shell job
itself.
+ // The working directory may contain any character a file name can, e.g. a
staging directory of
+ // an earlier version, which was named after the job template. It must
survive the JSON encoding
+ // and the '/'-joining unchanged. Some of these characters are only valid
in POSIX file names,
+ // like the shell job itself.
String specialName =
"a b \"quoted\" back\\slash 䏿–‡ \t tab \n newline %20
#!$&'()*+,;=@[]{}~`^|<>?";
File jobDir =
diff --git a/docs/table-maintenance-service/optimizer-cli-reference.md
b/docs/table-maintenance-service/optimizer-cli-reference.md
index 81c35f0a1a..c93b2bde9b 100644
--- a/docs/table-maintenance-service/optimizer-cli-reference.md
+++ b/docs/table-maintenance-service/optimizer-cli-reference.md
@@ -347,7 +347,7 @@ CALL `rest_catalog`.system.expire_snapshots(
```bash
curl -sS "http://localhost:8090/api/metalakes/test/jobs/{job_id}" | jq
'.job.state'
-cat
/tmp/gravitino/jobs/staging/test/builtin-iceberg-expire-snapshots/{job_id}/stdout.log
+cat /tmp/gravitino/jobs/staging/job-runs/{job_id}/output.log
```
A successful run reports its state as `SUCCEEDED` and logs the counts it
removed:
diff --git a/docs/table-maintenance-service/optimizer-troubleshooting.md
b/docs/table-maintenance-service/optimizer-troubleshooting.md
index 099b9b4814..98e738021b 100644
--- a/docs/table-maintenance-service/optimizer-troubleshooting.md
+++ b/docs/table-maintenance-service/optimizer-troubleshooting.md
@@ -12,7 +12,7 @@ license: "This software is licensed under the Apache License
version 2."
Failures fall into three groups, matching where they occur in the workflow.
Command and argument errors surface immediately. Evaluation problems produce no
output rather than an error, which is what makes them confusing. Execution
failures happen inside Spark, so the real message is in the staging log rather
than the API response.
-Staging logs live under
`/tmp/gravitino/jobs/staging/{metalake}/{job_template_name}/{job_id}/`,
controlled by `gravitino.job.stagingDir`. Read `error.log` for failures and
`output.log` for results.
+Staging logs live under `/tmp/gravitino/jobs/staging/job-runs/{job_id}/`,
controlled by `gravitino.job.stagingDir`. Read `error.log` for failures and
`output.log` for results.
## Command and Argument Errors