This is an automated email from the ASF dual-hosted git repository.

jerryshao 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 a4a21e3df3 [#12716] fix(job): retrieve local job output from any 
server sharing the staging directory (#13373)
a4a21e3df3 is described below

commit a4a21e3df365bc1f6ac36fe01a61a7129ae22eb4
Author: Jerry Shao <[email protected]>
AuthorDate: Tue Sep 22 19:26:10 2026 +0800

    [#12716] fix(job): retrieve local job output from any server sharing the 
staging directory (#13373)
    
    ### What changes were proposed in this pull request?
    
    Follow-up to #13250, for the multi-node gap raised in its review. With
    the local job executor, `getJob(includeOutput=true)` only returned a
    job's output on the server that ran it.
    
    - `LocalJobExecutor` no longer keeps job working directories in an
    in-memory map. When a job is submitted, it writes a write-once index
    file, `<gravitino.job.stagingDir>/.job-output-index/<jobId>.json`,
    containing `{"version":1,"workingDir":"<path relative to the staging
    directory>"}`. The output is found from the job id alone, so any server
    sharing the staging directory can return it. This also works after a
    restart or a template rename.
    - The output files stay where they were (`output.log` / `error.log` in
    the job's staging directory). The `JobExecutor` SPI and the `JobManager`
    logic are unchanged.
    - `JobExecutorFactory` passes `gravitino.job.stagingDir` to the built-in
    local job executor.
    - An hourly task removes index files that are more than an hour old and
    whose staging directory `JobManager` has already removed. An index
    therefore lives exactly as long as its job's output.
    - Output retrieval stays best-effort:
    - Failing to create or write the index only logs a warning. It never
    fails a job submission or the server startup.
    - An invalid index, or one pointing outside the staging directory,
    returns empty output.
    - Docs:
      - the shared staging directory requirement;
      - the index directory and the upgrade/downgrade behavior;
      - removed the old statement that output is lost on restart.
    
    ### Why are the changes needed?
    
    In a multi-server deployment behind a load balancer, a request for a
    job's output returned empty output whenever it reached a server other
    than the one that ran the job. That looks exactly like a job that
    printed nothing. The same happened on the server that ran the job, once
    the executor's in-memory status expired (after 1 hour) or after a
    restart, even though the staging directory is kept for
    `gravitino.job.stagingDirKeepTimeInMs` (7 days by default).
    
    Related: #12716 (follow-up to #13250), epic #12667. Pre-existing issues
    found along the way: #13371, #13372.
    
    ### Does this PR introduce _any_ user-facing change?
    
    - `getJob(includeOutput=true)` returns the output from any server
    sharing `gravitino.job.stagingDir`, also after server restarts, until
    the job's staging directory is removed.
    - A new directory, `<gravitino.job.stagingDir>/.job-output-index`. No
    configuration property is added or changed.
    - Jobs submitted before this change, or by a server not yet upgraded
    during a rolling upgrade, return empty output.
    - Servers that don't share the staging directory still only return
    output for the jobs they ran. This is documented, and they now log a
    warning.
    
    ### How was this patch tested?
    
    - `TestLocalJobExecutor`:
      - the index content;
    - reading the output from another executor instance, and after the job
    status expired;
      - special characters in the working directory;
      - invalid index contents and job ids;
      - the cleanup: age guard, interruption, missing directory;
      - index write failures, and concurrent initialization;
      - warning only once.
    - `TestJobManagerMultiNode`, reading a job's output:
      - from another node;
      - after the node that ran the job exits;
      - after a template rename;
      - for template names with special characters;
      - from a node that doesn't share the staging directory.
    - New `TestJobExecutorFactory`.
    - Ran:
      - `./gradlew :core:test :server:test -PskipITs`
      - `JobIT` (embedded mode)
      - `:core:spotlessCheck`, `:core:javadoc`, `:docs:build`
    
    🤖 Generated with [Claude Code](https://claude.com/claude-code)
    
    ---------
    
    Co-authored-by: Claude Opus 5 <[email protected]>
---
 .../main/java/org/apache/gravitino/Configs.java    |   6 +-
 .../gravitino/connector/job/JobExecutor.java       |  10 +-
 .../apache/gravitino/job/JobExecutorFactory.java   |  10 +-
 .../java/org/apache/gravitino/job/JobManager.java  |   6 +-
 .../gravitino/job/local/LocalJobExecutor.java      | 409 +++++++++++++--
 .../job/local/LocalJobExecutorConfigs.java         |   6 +
 .../gravitino/job/TestJobExecutorFactory.java      | 162 ++++++
 .../org/apache/gravitino/job/TestJobManager.java   |  15 +-
 .../gravitino/job/TestJobManagerMultiNode.java     | 146 +++++-
 .../gravitino/job/local/TestLocalJobExecutor.java  | 568 ++++++++++++++++++++-
 docs/gravitino-server-config.md                    |  21 +-
 docs/manage-jobs-in-gravitino.md                   |  47 +-
 docs/open-api/jobs.yaml                            |   3 +
 13 files changed, 1315 insertions(+), 94 deletions(-)

diff --git a/core/src/main/java/org/apache/gravitino/Configs.java 
b/core/src/main/java/org/apache/gravitino/Configs.java
index 18e9805096..5c016854cf 100644
--- a/core/src/main/java/org/apache/gravitino/Configs.java
+++ b/core/src/main/java/org/apache/gravitino/Configs.java
@@ -548,7 +548,11 @@ public class Configs {
 
   public static final ConfigEntry<String> JOB_STAGING_DIR =
       new ConfigBuilder("gravitino.job.stagingDir")
-          .doc("Directory for managing staging files when running jobs.")
+          .doc(
+              "Directory for managing staging files when running jobs. When 
multiple Gravitino "
+                  + "servers share the same metadata store, it must be on 
storage shared by all "
+                  + "servers, otherwise the output of a job run by the local 
job executor can "
+                  + "only be retrieved from the server that ran it.")
           .version(ConfigConstants.VERSION_1_0_0)
           .stringConf()
           .checkValue(StringUtils::isNotBlank, 
ConfigConstants.NOT_BLANK_ERROR_MSG)
diff --git 
a/core/src/main/java/org/apache/gravitino/connector/job/JobExecutor.java 
b/core/src/main/java/org/apache/gravitino/connector/job/JobExecutor.java
index 6df99d9d9b..f6e5243213 100644
--- a/core/src/main/java/org/apache/gravitino/connector/job/JobExecutor.java
+++ b/core/src/main/java/org/apache/gravitino/connector/job/JobExecutor.java
@@ -122,8 +122,9 @@ public interface JobExecutor extends Closeable {
    * retrieval don't need to override this method. Unlike {@link 
#getJobStatus(String)}/{@link
    * #cancelJob(String)}, this method never throws for a job the executor 
doesn't (or no longer)
    * know about - the job entity itself may still exist even after the 
executor's own bookkeeping
-   * for its output has expired or been lost (e.g. across a restart), so 
"unknown to this executor"
-   * is reported as empty output, not as an error.
+   * for its output has expired or been lost (e.g. when it's only kept in 
memory and the server
+   * restarts, or kept on storage this server can't reach), so "unknown to 
this executor" is
+   * reported as empty output, not as an error.
    *
    * @param jobId The unique identifier of the job.
    * @param maxLines The maximum number of (most recent) lines to return, 
resolved by the caller
@@ -144,8 +145,9 @@ public interface JobExecutor extends Closeable {
    * retrieval don't need to override this method. Unlike {@link 
#getJobStatus(String)}/{@link
    * #cancelJob(String)}, this method never throws for a job the executor 
doesn't (or no longer)
    * know about - the job entity itself may still exist even after the 
executor's own bookkeeping
-   * for its output has expired or been lost (e.g. across a restart), so 
"unknown to this executor"
-   * is reported as empty output, not as an error.
+   * for its output has expired or been lost (e.g. when it's only kept in 
memory and the server
+   * restarts, or kept on storage this server can't reach), so "unknown to 
this executor" is
+   * reported as empty output, not as an error.
    *
    * @param jobId The unique identifier of the job.
    * @param maxLines The maximum number of (most recent) lines to return, 
resolved by the caller
diff --git 
a/core/src/main/java/org/apache/gravitino/job/JobExecutorFactory.java 
b/core/src/main/java/org/apache/gravitino/job/JobExecutorFactory.java
index 808c8c9d26..dc0ec158a6 100644
--- a/core/src/main/java/org/apache/gravitino/job/JobExecutorFactory.java
+++ b/core/src/main/java/org/apache/gravitino/job/JobExecutorFactory.java
@@ -20,6 +20,7 @@ package org.apache.gravitino.job;
 
 import com.google.common.base.Preconditions;
 import com.google.common.collect.ImmutableMap;
+import com.google.common.collect.Maps;
 import java.util.Map;
 import org.apache.commons.lang3.StringUtils;
 import org.apache.gravitino.Config;
@@ -60,10 +61,17 @@ public class JobExecutorFactory {
         jobExecutorName);
 
     Map<String, String> configs =
-        config.getConfigsWithPrefix(JOB_EXECUTOR_CONF_PREFIX + jobExecutorName 
+ ".");
+        Maps.newHashMap(
+            config.getConfigsWithPrefix(JOB_EXECUTOR_CONF_PREFIX + 
jobExecutorName + "."));
     try {
       JobExecutor jobExecutor =
           (JobExecutor) 
Class.forName(clzName).getDeclaredConstructor().newInstance();
+      if (jobExecutor instanceof LocalJobExecutor) {
+        // The local job executor, and any subclass of it, keeps its output 
index under the job
+        // staging directory, so it must resolve paths against exactly the 
directory JobManager
+        // stages jobs in.
+        configs.put(LocalJobExecutorConfigs.STAGING_DIR, 
config.get(Configs.JOB_STAGING_DIR));
+      }
       jobExecutor.initialize(configs);
       return jobExecutor;
 
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 56614b15ba..7e67602a43 100644
--- a/core/src/main/java/org/apache/gravitino/job/JobManager.java
+++ b/core/src/main/java/org/apache/gravitino/job/JobManager.java
@@ -461,9 +461,9 @@ public class JobManager implements JobOperationDispatcher {
 
     // The job entity's existence was already confirmed above via the entity 
store, which is the
     // durable source of truth. The executor's own bookkeeping for a job's 
output is best-effort
-    // and can legitimately expire or be lost independently of the entity 
(e.g. LocalJobExecutor
-    // only retains in-memory output location for a limited time, and loses it 
entirely across a
-    // restart), so JobExecutor#getJobStdout/getJobStderr report a job unknown 
to the executor as
+    // and can legitimately be unavailable while the entity still exists (e.g. 
LocalJobExecutor
+    // can't reach the output of a job that ran on a server not sharing its 
staging directory),
+    // so JobExecutor#getJobStdout/getJobStderr report a job unknown to the 
executor as
     // empty output rather than an error - it never means "job does not exist" 
at this point.
     List<String> stdout =
         jobExecutor.getJobStdout(entity.jobExecutionId(), effectiveMaxLines, 
effectiveMaxBytes);
diff --git 
a/core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutor.java 
b/core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutor.java
index d55d7b47f6..296a6d89de 100644
--- a/core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutor.java
+++ b/core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutor.java
@@ -24,18 +24,32 @@ import static 
org.apache.gravitino.job.local.LocalJobExecutorConfigs.DEFAULT_MAX
 import static 
org.apache.gravitino.job.local.LocalJobExecutorConfigs.DEFAULT_WAITING_QUEUE_SIZE;
 import static 
org.apache.gravitino.job.local.LocalJobExecutorConfigs.JOB_STATUS_KEEP_TIME_MS;
 import static 
org.apache.gravitino.job.local.LocalJobExecutorConfigs.MAX_RUNNING_JOBS;
+import static 
org.apache.gravitino.job.local.LocalJobExecutorConfigs.STAGING_DIR;
 import static 
org.apache.gravitino.job.local.LocalJobExecutorConfigs.WAITING_QUEUE_SIZE;
 
+import com.fasterxml.jackson.databind.JsonNode;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.fasterxml.jackson.databind.node.ObjectNode;
 import com.google.common.annotations.VisibleForTesting;
+import com.google.common.base.Joiner;
 import com.google.common.base.Preconditions;
 import com.google.common.base.Splitter;
 import com.google.common.collect.ImmutableList;
 import com.google.common.collect.Maps;
-import java.io.File;
-import java.io.FileNotFoundException;
 import java.io.IOException;
-import java.io.RandomAccessFile;
+import java.nio.ByteBuffer;
+import java.nio.channels.ClosedByInterruptException;
+import java.nio.channels.SeekableByteChannel;
 import java.nio.charset.StandardCharsets;
+import java.nio.file.DirectoryStream;
+import java.nio.file.Files;
+import java.nio.file.InvalidPathException;
+import java.nio.file.LinkOption;
+import java.nio.file.NoSuchFileException;
+import java.nio.file.Path;
+import java.nio.file.Paths;
+import java.nio.file.StandardOpenOption;
+import java.util.Arrays;
 import java.util.List;
 import java.util.Map;
 import java.util.UUID;
@@ -47,16 +61,32 @@ import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.ThreadLocalRandom;
 import java.util.concurrent.ThreadPoolExecutor;
 import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.regex.Pattern;
 import javax.annotation.Nullable;
+import org.apache.commons.lang3.StringUtils;
 import org.apache.commons.lang3.tuple.Pair;
 import org.apache.gravitino.connector.job.JobExecutor;
 import org.apache.gravitino.exceptions.NoSuchJobException;
 import org.apache.gravitino.job.JobHandle;
 import org.apache.gravitino.job.JobTemplate;
 import org.apache.gravitino.job.SparkJobTemplate;
+import org.apache.gravitino.json.JsonUtils;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+/**
+ * A {@link JobExecutor} that runs jobs as processes on the Gravitino server 
itself.
+ *
+ * <p>Each job runs in its staging directory under {@code 
gravitino.job.stagingDir}, and its stdout
+ * and stderr are captured in files there. To find those files by job id 
alone, the executor writes
+ * a small, write-once index file per job under {@code 
<stagingDir>/.job-output-index}. In a
+ * multi-node deployment the staging directory must be shared by all Gravitino 
servers, otherwise a
+ * job's output can only be retrieved from the server that ran it.
+ *
+ * <p>{@link #initialize(Map)} requires {@link 
LocalJobExecutorConfigs#STAGING_DIR}, which {@code
+ * JobExecutorFactory} sets from {@code gravitino.job.stagingDir}.
+ */
 public class LocalJobExecutor implements JobExecutor {
 
   private static final Logger LOG = 
LoggerFactory.getLogger(LocalJobExecutor.class);
@@ -65,6 +95,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
+  // rename, so it doesn't collide with a metalake's staging directory. A 
metalake created before
+  // that check could still collide: the cleanup then skips its directories 
rather than failing.
+  private static final String OUTPUT_INDEX_DIR_NAME = ".job-output-index";
+
+  private static final String OUTPUT_INDEX_FILE_SUFFIX = ".json";
+
+  private static final int OUTPUT_INDEX_VERSION = 1;
+
+  private static final String OUTPUT_INDEX_VERSION_FIELD = "version";
+
+  private static final String OUTPUT_INDEX_WORKING_DIR_FIELD = "workingDir";
+
+  // How often indexes are scanned, not how long they are kept: an index is 
removed only once
+  // JobManager has removed the job's staging directory, so it lives exactly 
as long as the output.
+  private static final long OUTPUT_INDEX_CLEANUP_INTERVAL_IN_MS = 
TimeUnit.HOURS.toMillis(1);
+
+  // A newer index is never removed: on shared storage, a server may see a new 
index before it sees
+  // the job's staging directory created by another server.
+  private static final long OUTPUT_INDEX_MIN_AGE_IN_MS = 
TimeUnit.HOURS.toMillis(1);
+
+  private static final Pattern JOB_ID_PATTERN =
+      Pattern.compile(
+          LOCAL_JOB_PREFIX
+              + 
"[0-9a-f]{8}-[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}");
+
+  private static final ObjectMapper MAPPER = JsonUtils.anyFieldMapper();
+
   // A random id of this executor instance, it is embedded in every job id 
submitted by this
   // instance, so that each Gravitino server only tracks the jobs it runs 
itself.
   private String executorId;
@@ -88,13 +146,25 @@ public class LocalJobExecutor implements JobExecutor {
   private long jobStatusKeepTimeInMs;
   private ScheduledExecutorService jobStatusCleanupExecutor;
 
+  // Separate from jobStatusCleanupExecutor, so a long scan of the index 
directory on shared storage
+  // never delays the job status cleanup.
+  private ScheduledExecutorService outputIndexCleanupExecutor;
+
   private volatile boolean finished = false;
 
   private Map<String, Process> runningProcesses;
 
-  // The working directory of each job, used to locate its captured 
stdout/stderr files. Cleaned
-  // up together with jobStatus so the two maps stay in sync.
-  private Map<String, File> jobWorkingDirs;
+  // Not resolved with toRealPath(): every path is derived from the same 
configured string, so
+  // leaving symlinks unresolved keeps them comparable across servers.
+  private Path stagingRoot;
+
+  private Path outputIndexDir;
+
+  private final AtomicBoolean missingOutputIndexWarned = new 
AtomicBoolean(false);
+
+  private final AtomicBoolean invalidOutputIndexWarned = new 
AtomicBoolean(false);
+
+  private final AtomicBoolean unsafeOutputWarned = new AtomicBoolean(false);
 
   @Override
   public void initialize(Map<String, String> configs) {
@@ -103,6 +173,31 @@ public class LocalJobExecutor implements JobExecutor {
     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");
+    this.stagingRoot = Paths.get(stagingDir).toAbsolutePath().normalize();
+    this.outputIndexDir = stagingRoot.resolve(OUTPUT_INDEX_DIR_NAME);
+    // Safe when servers sharing the staging directory start concurrently: an 
already existing
+    // directory counts as created. Not fatal, output retrieval must never 
prevent the server from
+    // starting, and writeOutputIndex() tries again for every job.
+    try {
+      Files.createDirectories(outputIndexDir);
+    } catch (IOException e) {
+      LOG.warn(
+          "Failed to create the job output index directory {}, it's retried on 
job submission",
+          outputIndexDir,
+          e);
+    }
+    LOG.info(
+        "Job output of the local job executor is located through {}. In a 
multi-node deployment, "
+            + "gravitino.job.stagingDir ({}) must be shared by all Gravitino 
servers, otherwise "
+            + "a job's output can only be retrieved from the server that ran 
it.",
+        outputIndexDir,
+        stagingRoot);
+
     int waitingQueueSize =
         configs.containsKey(WAITING_QUEUE_SIZE)
             ? Integer.parseInt(configs.get(WAITING_QUEUE_SIZE))
@@ -175,9 +270,21 @@ public class LocalJobExecutor implements JobExecutor {
         jobStatusCleanupIntervalInMs,
         jobStatusCleanupIntervalInMs,
         TimeUnit.MILLISECONDS);
+    this.outputIndexCleanupExecutor =
+        Executors.newSingleThreadScheduledExecutor(
+            runnable -> {
+              Thread thread = new Thread(runnable);
+              thread.setName("LocalJobOutputIndexCleanup-" + thread.getId());
+              thread.setDaemon(true);
+              return thread;
+            });
+    outputIndexCleanupExecutor.scheduleWithFixedDelay(
+        this::cleanupOutputIndexes,
+        OUTPUT_INDEX_CLEANUP_INTERVAL_IN_MS,
+        OUTPUT_INDEX_CLEANUP_INTERVAL_IN_MS,
+        TimeUnit.MILLISECONDS);
 
     this.runningProcesses = Maps.newConcurrentMap();
-    this.jobWorkingDirs = Maps.newConcurrentMap();
 
     // Spark is optional for the local job executor, so a missing Spark 
installation must not fail
     // the server startup. Warn early instead; Spark jobs will be rejected at 
submission.
@@ -208,11 +315,10 @@ public class LocalJobExecutor implements JobExecutor {
       }
 
       jobStatus.put(newJobId, Pair.of(JobHandle.Status.QUEUED, 
UNEXPIRED_TIME_IN_MS));
-      // Retain the working directory so output.log/error.log can be located 
later without
-      // needing to keep the running Process object around after the job 
finishes.
-      jobWorkingDirs.put(newJobId, 
LocalProcessBuilder.resolveWorkingDirectory(jobTemplate));
     }
 
+    // Written only once the job is accepted, so a rejected submission leaves 
no index behind.
+    writeOutputIndex(newJobId, jobTemplate);
     return newJobId;
   }
 
@@ -314,10 +420,10 @@ public class LocalJobExecutor implements JobExecutor {
 
     jobExecutorService.shutdownNow();
 
-    // Stop the job status cleanup executor
+    // Stop the job status and output index cleanup executors
     jobStatusCleanupExecutor.shutdownNow();
+    outputIndexCleanupExecutor.shutdownNow();
     jobStatus.clear();
-    jobWorkingDirs.clear();
   }
 
   public void runJob(Pair<String, JobTemplate> jobPair) {
@@ -404,50 +510,268 @@ public class LocalJobExecutor implements JobExecutor {
       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);
-    if (workingDir == null) {
+    Path workingDir = locateWorkingDir(jobId);
+    if (workingDir == null || !isRealWorkingDir(jobId, workingDir)) {
       return ImmutableList.of();
     }
-    return readLastLines(new File(workingDir, fileName), maxLines, maxBytes);
+    return readLastLines(jobId, workingDir.resolve(fileName), maxLines, 
maxBytes);
+  }
+
+  // The job itself can replace its working directory, or a directory above 
it, with a symlink,
+  // e.g. to another job's directory, or to a directory that only exists on 
another server sharing
+  // the staging directory. Only read from the working directory if its real 
location is exactly
+  // where the index says it is.
+  private boolean isRealWorkingDir(String jobId, Path workingDir) {
+    Path expected;
+    Path actual;
+    try {
+      expected = 
stagingRoot.toRealPath().resolve(stagingRoot.relativize(workingDir));
+      actual = workingDir.toRealPath();
+    } catch (NoSuchFileException e) {
+      // Not started yet, or removed, e.g. by JobManager#cleanUpStagingDirs.
+      return false;
+    } catch (IOException e) {
+      throw new RuntimeException("Failed to resolve the working directory of 
job " + jobId, e);
+    }
+    if (!actual.equals(expected)) {
+      warnOnce(
+          unsafeOutputWarned,
+          "The working directory {} of job {} resolves to {} through a 
symlink, so its output "
+              + "isn't returned",
+          workingDir,
+          jobId,
+          actual);
+      return false;
+    }
+    return true;
+  }
+
+  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 File getWorkingDir(String jobId) {
-    return jobWorkingDirs.get(jobId);
+  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);
+    }
+
+    Path workingDir = parseOutputIndex(content);
+    if (workingDir == null) {
+      warnOnce(
+          invalidOutputIndexWarned,
+          "The output index {} of job {} is invalid or unsupported, so its 
output can't be "
+              + "retrieved. It may have been written by a newer Gravitino 
version.",
+          indexFile,
+          jobId);
+    }
+    return workingDir;
   }
 
-  private List<String> readLastLines(File file, int maxLines, int maxBytes) {
-    if (!file.exists()) {
+  // Returns null if the index can't be interpreted, including when it points 
outside the staging
+  // directory.
+  @Nullable
+  private Path parseOutputIndex(byte[] content) {
+    JsonNode index;
+    try {
+      index = MAPPER.readTree(content);
+    } catch (IOException e) {
+      return null;
+    }
+    if (index == null || !index.isObject()) {
+      return null;
+    }
+
+    JsonNode version = index.get(OUTPUT_INDEX_VERSION_FIELD);
+    if (version == null || !version.isInt() || version.intValue() != 
OUTPUT_INDEX_VERSION) {
+      return null;
+    }
+
+    JsonNode relativePath = index.get(OUTPUT_INDEX_WORKING_DIR_FIELD);
+    if (relativePath == null
+        || !relativePath.isTextual()
+        || StringUtils.isBlank(relativePath.textValue())
+        || relativePath.textValue().startsWith("/")) {
+      return null;
+    }
+
+    // Resolve the '/'-separated path against this server's staging directory, 
then reject anything
+    // that normalizes to outside of it (e.g. "../../etc") or to the staging 
directory itself: the
+    // index lives on shared storage and must never make this server read an 
arbitrary file.
+    Path workingDir = stagingRoot;
+    try {
+      for (String segment : 
Splitter.on('/').omitEmptyStrings().split(relativePath.textValue())) {
+        workingDir = workingDir.resolve(segment);
+      }
+    } catch (InvalidPathException e) {
+      // E.g. a NUL character, which no file name can contain.
+      return null;
+    }
+    workingDir = workingDir.normalize();
+    return workingDir.startsWith(stagingRoot) && 
!workingDir.equals(stagingRoot)
+        ? workingDir
+        : null;
+  }
+
+  private Path outputIndexFile(String jobId) {
+    return outputIndexDir.resolve(jobId + OUTPUT_INDEX_FILE_SUFFIX);
+  }
+
+  // Warns once per executor instance, a client polling the output would 
otherwise flood the log.
+  private static void warnOnce(AtomicBoolean warned, String message, Object... 
args) {
+    if (warned.compareAndSet(false, true)) {
+      LOG.warn(message + " Later occurrences are logged at DEBUG level.", 
args);
+    } else {
+      LOG.debug(message, args);
+    }
+  }
+
+  private List<String> readLastLines(String jobId, Path file, int maxLines, 
int maxBytes) {
+    if (!Files.exists(file, LinkOption.NOFOLLOW_LINKS)) {
       // The job hasn't started (or hasn't produced this stream) yet.
       return ImmutableList.of();
     }
+    // The job itself can replace its output file with a symlink, e.g. to a 
file that only exists
+    // on another server sharing the staging directory, so a symlink is never 
followed.
+    if (!Files.isRegularFile(file, LinkOption.NOFOLLOW_LINKS)) {
+      warnOnce(
+          unsafeOutputWarned,
+          "The output file {} of job {} is not a regular file, so it isn't 
returned",
+          file,
+          jobId);
+      return ImmutableList.of();
+    }
 
-    long fileLength = file.length();
-    int windowSize = (int) Math.min(fileLength, maxBytes);
-    long startOffset = fileLength - windowSize;
-
-    // Read one extra leading byte (when available) so we can tell whether the 
window's first
-    // line is already complete - i.e. the file byte immediately before the 
window is itself a
-    // line terminator - rather than always assuming it's a partial line and 
discarding it.
-    long readOffset = Math.max(0, startOffset - 1);
-    byte[] probeWindow = new byte[(int) (fileLength - readOffset)];
-    try (RandomAccessFile raf = new RandomAccessFile(file, "r")) {
-      raf.seek(readOffset);
-      raf.readFully(probeWindow);
-    } catch (FileNotFoundException e) {
+    long startOffset;
+    byte[] probeWindow;
+    // NOFOLLOW_LINKS also covers the file being replaced with a symlink after 
the check above.
+    try (SeekableByteChannel channel =
+        Files.newByteChannel(file, StandardOpenOption.READ, 
LinkOption.NOFOLLOW_LINKS)) {
+      long fileLength = channel.size();
+      int windowSize = (int) Math.min(fileLength, maxBytes);
+      startOffset = fileLength - windowSize;
+
+      // Read one extra leading byte (when available) so we can tell whether 
the window's first
+      // line is already complete - i.e. the file byte immediately before the 
window is itself a
+      // line terminator - rather than always assuming it's a partial line and 
discarding it.
+      long readOffset = Math.max(0, startOffset - 1);
+      ByteBuffer buffer = ByteBuffer.allocate((int) (fileLength - readOffset));
+      channel.position(readOffset);
+      while (buffer.hasRemaining() && channel.read(buffer) >= 0) {
+        // Keep reading until the window is full or the end of the file is 
reached.
+      }
+      probeWindow = Arrays.copyOf(buffer.array(), buffer.position());
+    } catch (NoSuchFileException e) {
       // The file existed at the check above but is gone now - most likely
       // JobManager#cleanUpStagingDirs deleted the staging directory 
concurrently with this call.
       // That's the same "output no longer available" situation the 
file-doesn't-exist check above
@@ -460,6 +784,9 @@ public class LocalJobExecutor implements JobExecutor {
       throw new RuntimeException("Failed to read job output file: " + file, e);
     }
 
+    if (probeWindow.length == 0) {
+      return ImmutableList.of();
+    }
     boolean windowStartsAtLineBoundary = startOffset == 0 || probeWindow[0] == 
'\n';
     int contentStart = startOffset == 0 ? 0 : 1;
     // The window may start mid-character if the file byte at contentStart 
happens to be a UTF-8
diff --git 
a/core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutorConfigs.java
 
b/core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutorConfigs.java
index 22cb5d0f72..c02eddb902 100644
--- 
a/core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutorConfigs.java
+++ 
b/core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutorConfigs.java
@@ -37,4 +37,10 @@ public class LocalJobExecutorConfigs {
   public static final long DEFAULT_JOB_STATUS_KEEP_TIME_MS = 60 * 60 * 1000; 
// 1 hour
 
   public static final String SPARK_HOME = "sparkHome";
+
+  /**
+   * The job staging directory, set by Gravitino from {@code 
gravitino.job.stagingDir} rather than
+   * by users. A value configured under the local job executor's prefix is 
overridden.
+   */
+  public static final String STAGING_DIR = "stagingDir";
 }
diff --git 
a/core/src/test/java/org/apache/gravitino/job/TestJobExecutorFactory.java 
b/core/src/test/java/org/apache/gravitino/job/TestJobExecutorFactory.java
new file mode 100644
index 0000000000..b8f1a8399c
--- /dev/null
+++ b/core/src/test/java/org/apache/gravitino/job/TestJobExecutorFactory.java
@@ -0,0 +1,162 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *  http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.gravitino.job;
+
+import com.google.common.collect.ImmutableMap;
+import java.io.File;
+import java.io.IOException;
+import java.nio.file.Files;
+import java.util.Map;
+import org.apache.commons.io.FileUtils;
+import org.apache.gravitino.Config;
+import org.apache.gravitino.Configs;
+import org.apache.gravitino.connector.job.JobExecutor;
+import org.apache.gravitino.exceptions.NoSuchJobException;
+import org.apache.gravitino.job.local.LocalJobExecutor;
+import org.apache.gravitino.job.local.LocalJobExecutorConfigs;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+public class TestJobExecutorFactory {
+
+  private static final String OUTPUT_INDEX_DIR_NAME = ".job-output-index";
+
+  private File testDir;
+
+  private File stagingDir;
+
+  private Config config;
+
+  @BeforeEach
+  public void setUp() throws IOException {
+    testDir = 
Files.createTempDirectory("gravitino-test-job-executor-factory").toFile();
+    stagingDir = new File(testDir, "staging");
+    config = new Config(false) {};
+    config.set(Configs.JOB_STAGING_DIR, stagingDir.getAbsolutePath());
+  }
+
+  @AfterEach
+  public void tearDown() throws IOException {
+    FileUtils.deleteDirectory(testDir);
+  }
+
+  @Test
+  public void testLocalJobExecutorUsesJobStagingDir() throws IOException {
+    try (JobExecutor executor = JobExecutorFactory.create(config)) {
+      Assertions.assertInstanceOf(LocalJobExecutor.class, executor);
+      Assertions.assertTrue(new File(stagingDir, 
OUTPUT_INDEX_DIR_NAME).isDirectory());
+    }
+  }
+
+  @Test
+  public void testLocalJobExecutorIgnoresConfiguredStagingDir() throws 
IOException {
+    File otherDir = new File(testDir, "other");
+    config.loadFromMap(
+        ImmutableMap.of(
+            "gravitino.jobExecutor.local." + 
LocalJobExecutorConfigs.STAGING_DIR,
+            otherDir.getAbsolutePath()),
+        key -> true);
+
+    try (JobExecutor executor = JobExecutorFactory.create(config)) {
+      Assertions.assertTrue(new File(stagingDir, 
OUTPUT_INDEX_DIR_NAME).isDirectory());
+      Assertions.assertFalse(otherDir.exists());
+    }
+  }
+
+  @Test
+  public void testLocalJobExecutorUnderCustomNameUsesJobStagingDir() throws 
IOException {
+    config.set(Configs.JOB_EXECUTOR, "mylocal");
+    config.loadFromMap(
+        ImmutableMap.of(
+            "gravitino.jobExecutor.mylocal.class", 
LocalJobExecutor.class.getCanonicalName()),
+        key -> true);
+
+    try (JobExecutor executor = JobExecutorFactory.create(config)) {
+      Assertions.assertInstanceOf(LocalJobExecutor.class, executor);
+      Assertions.assertTrue(new File(stagingDir, 
OUTPUT_INDEX_DIR_NAME).isDirectory());
+    }
+  }
+
+  @Test
+  public void testLocalJobExecutorSubclassUsesJobStagingDir() throws 
IOException {
+    config.set(Configs.JOB_EXECUTOR, "custom");
+    config.loadFromMap(
+        ImmutableMap.of(
+            "gravitino.jobExecutor.custom.class", 
CustomLocalJobExecutor.class.getName()),
+        key -> true);
+
+    try (JobExecutor executor = JobExecutorFactory.create(config)) {
+      Assertions.assertInstanceOf(CustomLocalJobExecutor.class, executor);
+      Assertions.assertTrue(new File(stagingDir, 
OUTPUT_INDEX_DIR_NAME).isDirectory());
+    }
+  }
+
+  @Test
+  public void testCustomJobExecutorConfigsAreUnchanged() throws IOException {
+    config.set(Configs.JOB_EXECUTOR, "recording");
+    config.loadFromMap(
+        ImmutableMap.of(
+            "gravitino.jobExecutor.recording.class",
+            RecordingJobExecutor.class.getName(),
+            "gravitino.jobExecutor.recording.foo",
+            "bar"),
+        key -> true);
+
+    try (JobExecutor executor = JobExecutorFactory.create(config)) {
+      Assertions.assertInstanceOf(RecordingJobExecutor.class, executor);
+      Map<String, String> configs = ((RecordingJobExecutor) executor).configs;
+      Assertions.assertEquals("bar", configs.get("foo"));
+      
Assertions.assertFalse(configs.containsKey(LocalJobExecutorConfigs.STAGING_DIR));
+    }
+  }
+
+  /** A user's subclass of the local job executor, inheriting its 
initialization. */
+  public static class CustomLocalJobExecutor extends LocalJobExecutor {}
+
+  /** A job executor that only records the configurations it's initialized 
with. */
+  public static class RecordingJobExecutor implements JobExecutor {
+
+    private Map<String, String> configs;
+
+    @Override
+    public void initialize(Map<String, String> configs) {
+      this.configs = configs;
+    }
+
+    @Override
+    public String submitJob(JobTemplate jobTemplate) {
+      throw new UnsupportedOperationException();
+    }
+
+    @Override
+    public JobHandle.Status getJobStatus(String jobId) throws 
NoSuchJobException {
+      throw new NoSuchJobException("No job found with ID: %s", jobId);
+    }
+
+    @Override
+    public void cancelJob(String jobId) throws NoSuchJobException {
+      throw new NoSuchJobException("No job found with ID: %s", jobId);
+    }
+
+    @Override
+    public void close() {}
+  }
+}
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 79318a429a..1091b38fec 100644
--- a/core/src/test/java/org/apache/gravitino/job/TestJobManager.java
+++ b/core/src/test/java/org/apache/gravitino/job/TestJobManager.java
@@ -32,6 +32,7 @@ import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
 import com.google.common.collect.ImmutableList;
+import com.google.common.collect.ImmutableMap;
 import com.google.common.collect.Lists;
 import com.sun.net.httpserver.HttpServer;
 import java.io.File;
@@ -45,6 +46,7 @@ import java.time.Instant;
 import java.time.temporal.ChronoUnit;
 import java.util.Collections;
 import java.util.List;
+import java.util.Map;
 import java.util.Optional;
 import java.util.Random;
 import java.util.UUID;
@@ -76,6 +78,7 @@ import 
org.apache.gravitino.exceptions.NoSuchMetalakeException;
 import org.apache.gravitino.exceptions.NonEmptyEntityException;
 import org.apache.gravitino.exceptions.OptimisticLockException;
 import org.apache.gravitino.job.local.LocalJobExecutor;
+import org.apache.gravitino.job.local.LocalJobExecutorConfigs;
 import org.apache.gravitino.json.JsonUtils;
 import org.apache.gravitino.lock.LockManager;
 import org.apache.gravitino.meta.AuditInfo;
@@ -1570,15 +1573,19 @@ public class TestJobManager {
   public void testPullJobStatusAcrossServersWithLocalJobExecutors() throws 
Exception {
     // Reproduces the multi-node deployment: two servers share the same 
metadata store, and each
     // has its own local job executor. The server that didn't run the job must 
not mark it FAILED.
+    File stagingRoot = 
Files.createTempDirectory("gravitino-test-multi-node-job").toFile();
+    File jobStagingDir = new File(stagingRoot, "job");
+    Assertions.assertTrue(jobStagingDir.mkdirs());
+    Map<String, String> executorConfigs =
+        ImmutableMap.of(LocalJobExecutorConfigs.STAGING_DIR, 
stagingRoot.getAbsolutePath());
     LocalJobExecutor ownerExecutor = new LocalJobExecutor();
     LocalJobExecutor otherExecutor = new LocalJobExecutor();
-    ownerExecutor.initialize(Collections.emptyMap());
-    otherExecutor.initialize(Collections.emptyMap());
+    ownerExecutor.initialize(executorConfigs);
+    otherExecutor.initialize(executorConfigs);
     JobManager ownerManager =
         Mockito.spy(new JobManager(config, entityStore, idGenerator, 
ownerExecutor));
     JobManager otherManager =
         Mockito.spy(new JobManager(config, entityStore, idGenerator, 
otherExecutor));
-    File jobStagingDir = 
Files.createTempDirectory("gravitino-test-multi-node-job").toFile();
 
     try {
       JobTemplate jobTemplate =
@@ -1610,7 +1617,7 @@ public class TestJobManager {
     } finally {
       ownerManager.close();
       otherManager.close();
-      FileUtils.deleteDirectory(jobStagingDir);
+      FileUtils.deleteDirectory(stagingRoot);
     }
   }
 
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 9c5c705173..a74d982db4 100644
--- a/core/src/test/java/org/apache/gravitino/job/TestJobManagerMultiNode.java
+++ b/core/src/test/java/org/apache/gravitino/job/TestJobManagerMultiNode.java
@@ -18,6 +18,7 @@
  */
 package org.apache.gravitino.job;
 
+import com.google.common.collect.ImmutableList;
 import com.google.common.collect.ImmutableMap;
 import com.google.common.collect.Lists;
 import java.io.File;
@@ -37,6 +38,7 @@ import org.apache.gravitino.cache.NoOpsCache;
 import org.apache.gravitino.connector.job.JobExecutor;
 import org.apache.gravitino.exceptions.NoSuchJobException;
 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.JobEntity;
@@ -63,6 +65,8 @@ public class TestJobManagerMultiNode extends TestJDBCBackend {
 
   private static final String TEMPLATE = "sleep_job";
 
+  private static final String ECHO_TEMPLATE = "echo_job";
+
   // Active jobs not updated for this long are expired, and finished jobs are 
cleaned up after it.
   // The tests move the job timestamps back instead of waiting for it to 
elapse.
   private static final long JOB_KEEP_TIME_IN_MS = TimeUnit.HOURS.toMillis(1);
@@ -83,9 +87,8 @@ public class TestJobManagerMultiNode extends TestJDBCBackend {
   public void setUpNodes() throws Exception {
     testDir = 
Files.createTempDirectory("gravitino-test-job-multi-node").toFile();
 
-    Config config = new Config(false) {};
-    config.set(Configs.JOB_STAGING_DIR, new File(testDir, 
"staging").getAbsolutePath());
-    config.set(Configs.JOB_STAGING_DIR_KEEP_TIME_IN_MS, JOB_KEEP_TIME_IN_MS);
+    // Both nodes share the same staging directory, like a deployment on 
shared storage.
+    Config config = newConfig(new File(testDir, "staging"));
     FieldUtils.writeField(GravitinoEnv.getInstance(), "lockManager", new 
LockManager(config), true);
 
     // Both nodes share the same metadata store, backed by the relational 
backend under test.
@@ -96,11 +99,10 @@ public class TestJobManagerMultiNode extends 
TestJDBCBackend {
 
     createAndInsertMakeLake(METALAKE);
     backend.insert(newSleepJobTemplateEntity(), false);
+    backend.insert(newEchoJobTemplateEntity(), false);
 
-    executorA = new LocalJobExecutor();
-    executorA.initialize(Collections.emptyMap());
-    executorB = new LocalJobExecutor();
-    executorB.initialize(Collections.emptyMap());
+    executorA = newLocalJobExecutor(config);
+    executorB = newLocalJobExecutor(config);
     nodeA = newJobManager(config, entityStore, executorA);
     nodeB = newJobManager(config, entityStore, executorB);
   }
@@ -198,6 +200,102 @@ public class TestJobManagerMultiNode extends 
TestJDBCBackend {
     Assertions.assertFalse(jobExists(cancellingJob.name()));
   }
 
+  @TestTemplate
+  public void testGetJobOutputFromAnotherNode() throws IOException {
+    JobEntity job = runEchoJobOnNodeA("a");
+
+    // Node B didn't run the job, it reads the output from the shared staging 
directory.
+    JobEntity jobWithOutput = nodeB.getJob(METALAKE, job.name(), true);
+    Assertions.assertEquals(ImmutableList.of("hello a"), 
jobWithOutput.stdout());
+    Assertions.assertEquals(ImmutableList.of("oops a"), 
jobWithOutput.stderr());
+
+    // The output doesn't depend on the node that ran the job being alive.
+    nodeA.close();
+    nodeA = null;
+    Assertions.assertEquals(
+        ImmutableList.of("hello a"), nodeB.getJob(METALAKE, job.name(), 
true).stdout());
+  }
+
+  @TestTemplate
+  public void testGetJobOutputAfterTemplateRenamed() throws IOException {
+    JobEntity job = runEchoJobOnNodeA("b");
+
+    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.
+    JobEntity jobWithOutput = nodeB.getJob(METALAKE, job.name(), true);
+    Assertions.assertEquals(newName, jobWithOutput.jobTemplateName());
+    Assertions.assertEquals(ImmutableList.of("hello b"), 
jobWithOutput.stdout());
+  }
+
+  @TestTemplate
+  public void testGetJobOutputOfTemplateNamedWithSpecialCharacters() throws 
IOException {
+    // Template names are not restricted, and name a directory level of the 
job staging directory.
+    for (String templateName : ImmutableList.of("etl job \"v2\" 中文 #1", 
"team/etl")) {
+      backend.insert(
+          newScriptJobTemplateEntity(
+              templateName, "echo \"hello $1\"\necho \"oops $1\" >&2\n", 
"{{name}}"),
+          false);
+
+      JobEntity job = nodeA.runJob(METALAKE, templateName, 
ImmutableMap.of("name", "d"));
+      Awaitility.await()
+          .atMost(1, TimeUnit.MINUTES)
+          .until(() -> executorA.getJobStatus(job.jobExecutionId()) == 
JobHandle.Status.SUCCEEDED);
+
+      JobEntity jobWithOutput = nodeB.getJob(METALAKE, job.name(), true);
+      Assertions.assertEquals(ImmutableList.of("hello d"), 
jobWithOutput.stdout(), templateName);
+      Assertions.assertEquals(ImmutableList.of("oops d"), 
jobWithOutput.stderr(), templateName);
+    }
+  }
+
+  @TestTemplate
+  public void testGetJobOutputDoesNotFollowSymlinkFromAnotherNode() throws 
IOException {
+    // A file on the node reading the output, outside the staging directory.
+    File secret = new File(testDir, "secret.txt");
+    Files.writeString(secret.toPath(), "top secret\n");
+    // The job replaces its own output file with an absolute symlink while 
running.
+    String templateName = "symlink_job";
+    backend.insert(
+        newScriptJobTemplateEntity(
+            templateName, "echo \"hello\"\nln -sf \"$1\" output.log\n", 
"{{target}}"),
+        false);
+
+    JobEntity job =
+        nodeA.runJob(METALAKE, templateName, ImmutableMap.of("target", 
secret.getAbsolutePath()));
+    Awaitility.await()
+        .atMost(1, TimeUnit.MINUTES)
+        .until(() -> executorA.getJobStatus(job.jobExecutionId()) == 
JobHandle.Status.SUCCEEDED);
+
+    Assertions.assertEquals(
+        Collections.emptyList(), nodeB.getJob(METALAKE, job.name(), 
true).stdout());
+  }
+
+  @TestTemplate
+  public void testGetJobOutputFromNodeNotSharingStagingDir() throws 
IOException {
+    JobEntity job = runEchoJobOnNodeA("c");
+
+    // A node with a staging directory of its own can't reach the output, and 
reports none
+    // instead of failing.
+    Config otherConfig = newConfig(new File(testDir, "other-staging"));
+    JobManager nodeC = newJobManager(otherConfig, entityStore, 
newLocalJobExecutor(otherConfig));
+    try {
+      JobEntity jobWithOutput = nodeC.getJob(METALAKE, job.name(), true);
+      Assertions.assertEquals(Collections.emptyList(), jobWithOutput.stdout());
+      Assertions.assertEquals(Collections.emptyList(), jobWithOutput.stderr());
+    } finally {
+      nodeC.close();
+    }
+  }
+
+  private JobEntity runEchoJobOnNodeA(String name) throws IOException {
+    JobEntity job = nodeA.runJob(METALAKE, ECHO_TEMPLATE, 
ImmutableMap.of("name", name));
+    Awaitility.await()
+        .atMost(1, TimeUnit.MINUTES)
+        .until(() -> executorA.getJobStatus(job.jobExecutionId()) == 
JobHandle.Status.SUCCEEDED);
+    return job;
+  }
+
   private JobEntity runLongJobOnNodeA() throws IOException {
     JobEntity job = nodeA.runJob(METALAKE, TEMPLATE, 
ImmutableMap.of("seconds", "600"));
     Awaitility.await()
@@ -253,20 +351,31 @@ public class TestJobManagerMultiNode extends 
TestJDBCBackend {
   }
 
   private JobTemplateEntity newSleepJobTemplateEntity() throws IOException {
-    File script = new File(testDir, "sleep-job.sh");
     // Exec the sleep, so that killing the job process also stops the sleep.
-    Files.writeString(script.toPath(), "#!/bin/bash\nexec sleep \"$1\"\n");
+    return newScriptJobTemplateEntity(TEMPLATE, "exec sleep \"$1\"\n", 
"{{seconds}}");
+  }
+
+  private JobTemplateEntity newEchoJobTemplateEntity() throws IOException {
+    return newScriptJobTemplateEntity(
+        ECHO_TEMPLATE, "echo \"hello $1\"\necho \"oops $1\" >&2\n", 
"{{name}}");
+  }
+
+  private JobTemplateEntity newScriptJobTemplateEntity(
+      String name, String scriptBody, String argument) throws IOException {
+    // Not named after the template, whose name may contain any character.
+    File script = Files.createTempFile(testDir.toPath(), "job", 
".sh").toFile();
+    Files.writeString(script.toPath(), "#!/bin/bash\n" + scriptBody);
     Assertions.assertTrue(script.setExecutable(true));
 
     return JobTemplateEntity.builder()
         .withId(RandomIdGenerator.INSTANCE.nextId())
-        .withName(TEMPLATE)
+        .withName(name)
         .withNamespace(NamespaceUtil.ofJobTemplate(METALAKE))
         .withTemplateContent(
             JobTemplateEntity.TemplateContent.builder()
                 .withJobType(JobTemplate.JobType.SHELL)
                 .withExecutable(script.getAbsolutePath())
-                .withArguments(Lists.newArrayList("{{seconds}}"))
+                .withArguments(Lists.newArrayList(argument))
                 .withEnvironments(Collections.emptyMap())
                 .withCustomFields(Collections.emptyMap())
                 .withScripts(Collections.emptyList())
@@ -276,6 +385,21 @@ public class TestJobManagerMultiNode extends 
TestJDBCBackend {
         .build();
   }
 
+  private static Config newConfig(File stagingDir) {
+    Config config = new Config(false) {};
+    config.set(Configs.JOB_STAGING_DIR, stagingDir.getAbsolutePath());
+    config.set(Configs.JOB_STAGING_DIR_KEEP_TIME_IN_MS, JOB_KEEP_TIME_IN_MS);
+    return config;
+  }
+
+  // Configured the way JobExecutorFactory configures it for a server with 
this configuration.
+  private static LocalJobExecutor newLocalJobExecutor(Config config) {
+    LocalJobExecutor executor = new LocalJobExecutor();
+    executor.initialize(
+        ImmutableMap.of(LocalJobExecutorConfigs.STAGING_DIR, 
config.get(Configs.JOB_STAGING_DIR)));
+    return executor;
+  }
+
   private static JobManager newJobManager(
       Config config, EntityStore entityStore, JobExecutor jobExecutor) {
     JobManager jobManager =
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 5a2dd71200..6c5441dd79 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
@@ -18,6 +18,7 @@
  */
 package org.apache.gravitino.job.local;
 
+import com.fasterxml.jackson.databind.JsonNode;
 import com.google.common.collect.ImmutableList;
 import com.google.common.collect.ImmutableMap;
 import com.google.common.collect.Lists;
@@ -28,10 +29,15 @@ import java.lang.reflect.Field;
 import java.net.URL;
 import java.nio.charset.StandardCharsets;
 import java.nio.file.Files;
+import java.nio.file.attribute.FileTime;
 import java.util.Collections;
 import java.util.List;
 import java.util.Map;
 import java.util.UUID;
+import java.util.concurrent.CyclicBarrier;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
 import java.util.concurrent.TimeUnit;
 import org.apache.commons.io.FileUtils;
 import org.apache.commons.lang3.StringUtils;
@@ -42,9 +48,18 @@ import org.apache.gravitino.job.JobManager;
 import org.apache.gravitino.job.JobTemplate;
 import org.apache.gravitino.job.ShellJobTemplate;
 import org.apache.gravitino.job.SparkJobTemplate;
+import org.apache.gravitino.json.JsonUtils;
 import org.apache.gravitino.meta.AuditInfo;
 import org.apache.gravitino.meta.JobTemplateEntity;
 import org.apache.gravitino.utils.NamespaceUtil;
+import org.apache.logging.log4j.Level;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.core.LogEvent;
+import org.apache.logging.log4j.core.LoggerContext;
+import org.apache.logging.log4j.core.appender.AbstractAppender;
+import org.apache.logging.log4j.core.config.AbstractConfiguration;
+import org.apache.logging.log4j.core.config.LoggerConfig;
+import org.apache.logging.log4j.core.layout.PatternLayout;
 import org.awaitility.Awaitility;
 import org.junit.jupiter.api.AfterAll;
 import org.junit.jupiter.api.AfterEach;
@@ -52,21 +67,27 @@ import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.BeforeAll;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.function.Executable;
 
 public class TestLocalJobExecutor {
 
   private static final int DEFAULT_TEST_MAX_BYTES = 1_000_000;
 
+  private static final String OUTPUT_INDEX_DIR_NAME = ".job-output-index";
+
   private static JobExecutor jobExecutor;
 
   private static JobTemplateEntity jobTemplateEntity;
 
+  private static File stagingRoot;
+
   private static File workingDir;
 
   @BeforeAll
   public static void setUpClass() throws IOException {
+    stagingRoot = 
Files.createTempDirectory("gravitino-test-local-job-staging").toFile();
     jobExecutor = new LocalJobExecutor();
-    jobExecutor.initialize(Collections.emptyMap());
+    jobExecutor.initialize(withStagingDir(Collections.emptyMap()));
 
     URL testJobScriptUrl = 
TestLocalJobExecutor.class.getResource("/test-job.sh");
     Assertions.assertNotNull(testJobScriptUrl);
@@ -102,11 +123,15 @@ public class TestLocalJobExecutor {
       jobExecutor.close();
       jobExecutor = null;
     }
+    if (stagingRoot != null) {
+      FileUtils.deleteDirectory(stagingRoot);
+    }
   }
 
   @BeforeEach
   public void setUp() throws IOException {
-    workingDir = 
Files.createTempDirectory("gravitino-test-local-job-executor").toFile();
+    // Jobs are staged under the staging directory, like JobManager does.
+    workingDir = Files.createTempDirectory(stagingRoot.toPath(), 
"job").toFile();
   }
 
   @AfterEach
@@ -165,7 +190,7 @@ public class TestLocalJobExecutor {
 
     LocalJobExecutor anotherExecutor = new LocalJobExecutor();
     try {
-      anotherExecutor.initialize(Collections.emptyMap());
+      anotherExecutor.initialize(withStagingDir(Collections.emptyMap()));
       Assertions.assertNotEquals(executor.executorId(), 
anotherExecutor.executorId());
       Assertions.assertFalse(anotherExecutor.ownsJob(jobId));
       Assertions.assertThrows(NoSuchJobException.class, () -> 
anotherExecutor.getJobStatus(jobId));
@@ -182,7 +207,8 @@ public class TestLocalJobExecutor {
   public void testRunJobsConcurrentlyUpToMaxRunningJobs() throws IOException {
     LocalJobExecutor executor = new LocalJobExecutor();
     try {
-      
executor.initialize(ImmutableMap.of(LocalJobExecutorConfigs.MAX_RUNNING_JOBS, 
"2"));
+      executor.initialize(
+          
withStagingDir(ImmutableMap.of(LocalJobExecutorConfigs.MAX_RUNNING_JOBS, "2")));
       String jobId1 = executor.submitJob(newSleepJobTemplate("sleep-1"));
       String jobId2 = executor.submitJob(newSleepJobTemplate("sleep-2"));
       String jobId3 = executor.submitJob(newSleepJobTemplate("sleep-3"));
@@ -236,7 +262,8 @@ public class TestLocalJobExecutor {
     File sparkHome = new File(workingDir, "spark");
     LocalJobExecutor exec = new LocalJobExecutor();
     exec.initialize(
-        ImmutableMap.of(LocalJobExecutorConfigs.SPARK_HOME, 
sparkHome.getAbsolutePath()));
+        withStagingDir(
+            ImmutableMap.of(LocalJobExecutorConfigs.SPARK_HOME, 
sparkHome.getAbsolutePath())));
 
     try {
       SparkJobTemplate template =
@@ -247,9 +274,12 @@ public class TestLocalJobExecutor {
               .build();
 
       // spark-submit does not exist, the job is rejected at submission 
instead of being queued.
+      int indexCountBeforeRejection = outputIndexFileCount();
       IllegalArgumentException e =
           Assertions.assertThrows(IllegalArgumentException.class, () -> 
exec.submitJob(template));
       Assertions.assertTrue(e.getMessage().contains("spark-submit is not found 
or not executable"));
+      // A rejected submission leaves no output index behind.
+      Assertions.assertEquals(indexCountBeforeRejection, 
outputIndexFileCount());
 
       // Once spark-submit is available, the same job is accepted.
       File sparkSubmit = new File(sparkHome, "bin/spark-submit");
@@ -347,10 +377,10 @@ public class TestLocalJobExecutor {
   @Test
   public void testGetJobOutputForQueuedJobReturnsEmpty() throws IOException {
     LocalJobExecutor exec = new LocalJobExecutor();
-    exec.initialize(ImmutableMap.of(LocalJobExecutorConfigs.MAX_RUNNING_JOBS, 
"1"));
+    
exec.initialize(withStagingDir(ImmutableMap.of(LocalJobExecutorConfigs.MAX_RUNNING_JOBS,
 "1")));
 
-    File workingDirA = 
Files.createTempDirectory("gravitino-test-local-job-executor-a").toFile();
-    File workingDirB = 
Files.createTempDirectory("gravitino-test-local-job-executor-b").toFile();
+    File workingDirA = Files.createTempDirectory(stagingRoot.toPath(), 
"job-a").toFile();
+    File workingDirB = Files.createTempDirectory(stagingRoot.toPath(), 
"job-b").toFile();
     try {
       Map<String, String> jobConf =
           ImmutableMap.of(
@@ -664,7 +694,8 @@ public class TestLocalJobExecutor {
       throws NoSuchFieldException, IllegalAccessException {
     LocalJobExecutor exec = new LocalJobExecutor();
 
-    
exec.initialize(ImmutableMap.of(LocalJobExecutorConfigs.JOB_STATUS_KEEP_TIME_MS,
 "1"));
+    exec.initialize(
+        
withStagingDir(ImmutableMap.of(LocalJobExecutorConfigs.JOB_STATUS_KEEP_TIME_MS, 
"1")));
 
     Field field = 
LocalJobExecutor.class.getDeclaredField("jobStatusKeepTimeInMs");
     field.setAccessible(true);
@@ -677,7 +708,8 @@ public class TestLocalJobExecutor {
       throws NoSuchFieldException, IllegalAccessException {
     LocalJobExecutor exec = new LocalJobExecutor();
 
-    
exec.initialize(ImmutableMap.of(LocalJobExecutorConfigs.JOB_STATUS_KEEP_TIME_MS,
 "11"));
+    exec.initialize(
+        
withStagingDir(ImmutableMap.of(LocalJobExecutorConfigs.JOB_STATUS_KEEP_TIME_MS, 
"11")));
 
     Field field = 
LocalJobExecutor.class.getDeclaredField("jobStatusKeepTimeInMs");
     field.setAccessible(true);
@@ -685,6 +717,522 @@ public class TestLocalJobExecutor {
     Assertions.assertEquals(11L, actualValue);
   }
 
+  @Test
+  public void testInitializeWithoutStagingDirFails() {
+    LocalJobExecutor exec = new LocalJobExecutor();
+    IllegalArgumentException e =
+        Assertions.assertThrows(
+            IllegalArgumentException.class, () -> 
exec.initialize(Collections.emptyMap()));
+    Assertions.assertTrue(e.getMessage().contains("staging directory"));
+  }
+
+  @Test
+  public void testConcurrentInitializationSharingNewStagingDir() throws 
Exception {
+    // Servers sharing a staging directory may start at the same time, all 
creating the index
+    // directory.
+    File newStagingRoot = new File(stagingRoot, "concurrent-staging");
+    int executorCount = 8;
+    CyclicBarrier barrier = new CyclicBarrier(executorCount);
+    ExecutorService pool = Executors.newFixedThreadPool(executorCount);
+    List<LocalJobExecutor> executors = 
Collections.synchronizedList(Lists.newArrayList());
+    try {
+      List<Future<?>> futures = Lists.newArrayList();
+      for (int i = 0; i < executorCount; i++) {
+        futures.add(
+            pool.submit(
+                () -> {
+                  LocalJobExecutor executor = new LocalJobExecutor();
+                  barrier.await();
+                  executor.initialize(
+                      ImmutableMap.of(
+                          LocalJobExecutorConfigs.STAGING_DIR, 
newStagingRoot.getAbsolutePath()));
+                  executors.add(executor);
+                  return null;
+                }));
+      }
+      for (Future<?> future : futures) {
+        future.get(1, TimeUnit.MINUTES);
+      }
+
+      Assertions.assertEquals(executorCount, executors.size());
+      Assertions.assertTrue(new File(newStagingRoot, 
OUTPUT_INDEX_DIR_NAME).isDirectory());
+    } finally {
+      pool.shutdownNow();
+      for (LocalJobExecutor executor : executors) {
+        executor.close();
+      }
+    }
+  }
+
+  @Test
+  public void testInitializeSucceedsWhenOutputIndexDirCannotBeCreated() throws 
IOException {
+    // A regular file where a directory is expected makes creating the index 
directory fail.
+    File blocker = new File(stagingRoot, "blocker");
+    FileUtils.writeStringToFile(blocker, "not a directory", 
StandardCharsets.UTF_8);
+    LocalJobExecutor exec = new LocalJobExecutor();
+    try {
+      Assertions.assertDoesNotThrow(
+          () ->
+              exec.initialize(
+                  ImmutableMap.of(
+                      LocalJobExecutorConfigs.STAGING_DIR,
+                      new File(blocker, "staging").getAbsolutePath())));
+    } finally {
+      exec.close();
+      FileUtils.deleteQuietly(blocker);
+    }
+  }
+
+  @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");
+    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());
+  }
+
+  @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.
+    String specialName =
+        "a b \"quoted\" back\\slash 中文 \t tab \n newline %20 
#!$&'()*+,;=@[]{}~`^|<>?";
+    File jobDir =
+        new File(workingDir, "metalake" + File.separator + specialName + 
File.separator + "job-1");
+    Assertions.assertTrue(jobDir.mkdirs());
+    String jobId = runSucceededJob(jobDir);
+
+    JsonNode index = 
JsonUtils.anyFieldMapper().readTree(outputIndexFile(jobId));
+    Assertions.assertEquals(
+        workingDir.getName() + "/metalake/" + specialName + "/job-1",
+        index.get("workingDir").textValue());
+
+    LocalJobExecutor anotherExecutor = new LocalJobExecutor();
+    try {
+      anotherExecutor.initialize(withStagingDir(Collections.emptyMap()));
+      Assertions.assertEquals(
+          6, anotherExecutor.getJobStdout(jobId, 1000, 
DEFAULT_TEST_MAX_BYTES).size());
+    } finally {
+      anotherExecutor.close();
+    }
+  }
+
+  @Test
+  public void testGetJobOutputFromAnotherExecutorInstance() throws IOException 
{
+    String jobId = runSucceededJob(workingDir);
+    List<String> stdout = jobExecutor.getJobStdout(jobId, 1000, 
DEFAULT_TEST_MAX_BYTES);
+    Assertions.assertEquals(6, stdout.size());
+
+    // Another server sharing the staging directory, or this server after a 
restart.
+    LocalJobExecutor anotherExecutor = new LocalJobExecutor();
+    try {
+      anotherExecutor.initialize(withStagingDir(Collections.emptyMap()));
+      Assertions.assertFalse(anotherExecutor.ownsJob(jobId));
+      Assertions.assertEquals(
+          stdout, anotherExecutor.getJobStdout(jobId, 1000, 
DEFAULT_TEST_MAX_BYTES));
+      Assertions.assertEquals(
+          Collections.emptyList(),
+          anotherExecutor.getJobStderr(jobId, 1000, DEFAULT_TEST_MAX_BYTES));
+    } finally {
+      anotherExecutor.close();
+    }
+  }
+
+  @Test
+  public void testGetJobOutputAfterJobStatusExpired() throws IOException {
+    LocalJobExecutor exec = new LocalJobExecutor();
+    try {
+      // The status of a finished job is dropped from memory almost right away.
+      exec.initialize(
+          
withStagingDir(ImmutableMap.of(LocalJobExecutorConfigs.JOB_STATUS_KEEP_TIME_MS, 
"10")));
+      String jobId = exec.submitJob(newRuntimeJobTemplate(workingDir));
+      // A queued or running job's status never expires, so a missing status 
means it finished.
+      Awaitility.await().atMost(3, TimeUnit.MINUTES).until(() -> 
!hasJobStatus(exec, jobId));
+
+      Assertions.assertEquals(6, exec.getJobStdout(jobId, 1000, 
DEFAULT_TEST_MAX_BYTES).size());
+    } finally {
+      exec.close();
+    }
+  }
+
+  @Test
+  public void testGetJobOutputWhenWorkingDirIsOutsideStagingDir() throws 
IOException {
+    File outsideDir = 
Files.createTempDirectory("gravitino-test-local-job-outside").toFile();
+    try {
+      String jobId = runSucceededJob(outsideDir);
+      Assertions.assertTrue(new File(outsideDir, "output.log").length() > 0);
+
+      Assertions.assertFalse(outputIndexFile(jobId).exists());
+      Assertions.assertEquals(
+          Collections.emptyList(), jobExecutor.getJobStdout(jobId, 100, 
DEFAULT_TEST_MAX_BYTES));
+    } finally {
+      FileUtils.deleteDirectory(outsideDir);
+    }
+  }
+
+  @Test
+  public void testGetJobOutputWithInvalidOutputIndexReturnsEmpty() throws 
IOException {
+    String jobId = runSucceededJob(workingDir);
+    String validIndex = FileUtils.readFileToString(outputIndexFile(jobId), 
StandardCharsets.UTF_8);
+
+    // A readable output outside the staging directory, next to it.
+    File outsideDir = 
Files.createTempDirectory("gravitino-test-local-job-outside").toFile();
+    try {
+      FileUtils.writeStringToFile(
+          new File(outsideDir, "output.log"), "outside\n", 
StandardCharsets.UTF_8);
+
+      List<String> invalidIndexes =
+          ImmutableList.of(
+              "{\"version\":1,\"workingDir\":\"../" + outsideDir.getName() + 
"\"}",
+              "{\"version\":1,\"workingDir\":\"" + 
outsideDir.getAbsolutePath() + "\"}",
+              "{\"version\":1,\"workingDir\":\".\"}",
+              "{\"version\":1,\"workingDir\":\"\"}",
+              "{\"version\":2,\"workingDir\":\"" + workingDir.getName() + 
"\"}",
+              "{\"workingDir\":\"" + workingDir.getName() + "\"}",
+              // No file name can contain a NUL character.
+              "{\"version\":1,\"workingDir\":\"a\\u0000b\"}",
+              // Truncated by a crash while being written.
+              "{\"version\":1,\"workingDir\":\"",
+              "");
+      for (String invalidIndex : invalidIndexes) {
+        FileUtils.writeStringToFile(outputIndexFile(jobId), invalidIndex, 
StandardCharsets.UTF_8);
+        Assertions.assertEquals(
+            Collections.emptyList(),
+            jobExecutor.getJobStdout(jobId, 100, DEFAULT_TEST_MAX_BYTES),
+            invalidIndex);
+      }
+
+      FileUtils.writeStringToFile(outputIndexFile(jobId), validIndex, 
StandardCharsets.UTF_8);
+      Assertions.assertEquals(
+          6, jobExecutor.getJobStdout(jobId, 100, 
DEFAULT_TEST_MAX_BYTES).size());
+    } finally {
+      FileUtils.deleteDirectory(outsideDir);
+    }
+  }
+
+  @Test
+  public void testGetJobOutputDoesNotFollowSymlinks() throws IOException {
+    // A readable file outside the staging directory, next to it.
+    File outsideDir = 
Files.createTempDirectory("gravitino-test-local-job-outside").toFile();
+    File anotherJobDir = Files.createTempDirectory(stagingRoot.toPath(), 
"another").toFile();
+    LocalJobExecutor anotherExecutor = new LocalJobExecutor();
+    try {
+      FileUtils.writeStringToFile(
+          new File(outsideDir, "output.log"), "outside\n", 
StandardCharsets.UTF_8);
+      runSucceededJob(anotherJobDir);
+      // Another server sharing the staging directory.
+      anotherExecutor.initialize(withStagingDir(Collections.emptyMap()));
+
+      // The job replaces its output file with a symlink to a file outside the 
staging directory.
+      File jobDir = new File(workingDir, "output-file");
+      Assertions.assertTrue(jobDir.mkdirs());
+      String outputFileJobId = runSucceededJob(jobDir);
+      replaceWithSymlink(new File(jobDir, "output.log"), new File(outsideDir, 
"output.log"));
+
+      // The job replaces its working directory with a symlink to a directory 
outside.
+      jobDir = new File(workingDir, "working-dir");
+      Assertions.assertTrue(jobDir.mkdirs());
+      String workingDirJobId = runSucceededJob(jobDir);
+      replaceWithSymlink(jobDir, outsideDir);
+
+      // The job replaces a directory above its working directory with a 
symlink to the one of
+      // another job, e.g. of another metalake: tpl/job -> another-tpl/job.
+      File anotherTemplateDir = new File(workingDir, "another-tpl");
+      FileUtils.copyDirectory(anotherJobDir, new File(anotherTemplateDir, 
"job"));
+      jobDir = new File(workingDir, "tpl" + File.separator + "job");
+      Assertions.assertTrue(jobDir.mkdirs());
+      String parentDirJobId = runSucceededJob(jobDir);
+      replaceWithSymlink(jobDir.getParentFile(), anotherTemplateDir);
+
+      for (String jobId : ImmutableList.of(outputFileJobId, workingDirJobId, 
parentDirJobId)) {
+        for (JobExecutor executor : ImmutableList.of(jobExecutor, 
anotherExecutor)) {
+          Assertions.assertEquals(
+              Collections.emptyList(),
+              executor.getJobStdout(jobId, 100, DEFAULT_TEST_MAX_BYTES),
+              jobId);
+        }
+      }
+    } finally {
+      anotherExecutor.close();
+      FileUtils.deleteDirectory(outsideDir);
+      FileUtils.deleteDirectory(anotherJobDir);
+    }
+  }
+
+  @Test
+  public void testGetJobOutputWithInvalidJobIdReturnsEmpty() {
+    // The job id is used as a file name, so it must never be able to carry 
path elements.
+    List<String> jobIds =
+        ImmutableList.of(
+            "local-job-../../etc/passwd",
+            "../local-job-00000000-" + UUID.randomUUID(),
+            // A job id from before the executor id was introduced.
+            "local-job-" + UUID.randomUUID(),
+            // A well-formed job id without an output index.
+            "local-job-00000000-" + UUID.randomUUID(),
+            "");
+    for (String jobId : jobIds) {
+      Assertions.assertEquals(
+          Collections.emptyList(),
+          jobExecutor.getJobStdout(jobId, 100, DEFAULT_TEST_MAX_BYTES),
+          jobId);
+    }
+  }
+
+  @Test
+  public void testCleanupOutputIndexes() throws IOException {
+    File keptDir = Files.createTempDirectory(stagingRoot.toPath(), 
"kept").toFile();
+    File removedDir = Files.createTempDirectory(stagingRoot.toPath(), 
"removed").toFile();
+    File recentDir = Files.createTempDirectory(stagingRoot.toPath(), 
"recent").toFile();
+    String keptJobId = runSucceededJob(keptDir);
+    String removedJobId = runSucceededJob(removedDir);
+    String recentJobId = runSucceededJob(recentDir);
+    // Like JobManager removing the staging directory of an expired job.
+    FileUtils.deleteDirectory(removedDir);
+    // A new index whose staging directory this server doesn't see yet, e.g. 
through a stale cache
+    // of the shared storage.
+    FileUtils.deleteDirectory(recentDir);
+
+    // An index this server can't interpret, e.g. written by a newer server, 
and an unrelated file.
+    File unsupportedIndex = outputIndexFile("local-job-00000000-" + 
UUID.randomUUID());
+    FileUtils.writeStringToFile(unsupportedIndex, "{\"version\":2}", 
StandardCharsets.UTF_8);
+    File unrelatedFile = new File(stagingRoot, OUTPUT_INDEX_DIR_NAME + 
File.separator + "README");
+    FileUtils.writeStringToFile(unrelatedFile, "keep me", 
StandardCharsets.UTF_8);
+    for (String jobId : ImmutableList.of(keptJobId, removedJobId)) {
+      ageOutputIndex(outputIndexFile(jobId));
+    }
+    ageOutputIndex(unsupportedIndex);
+    try {
+      ((LocalJobExecutor) jobExecutor).cleanupOutputIndexes();
+
+      Assertions.assertFalse(outputIndexFile(removedJobId).exists());
+      Assertions.assertTrue(outputIndexFile(keptJobId).exists());
+      Assertions.assertTrue(outputIndexFile(recentJobId).exists());
+      Assertions.assertTrue(unsupportedIndex.exists());
+      Assertions.assertTrue(unrelatedFile.exists());
+      Assertions.assertEquals(
+          6, jobExecutor.getJobStdout(keptJobId, 100, 
DEFAULT_TEST_MAX_BYTES).size());
+    } finally {
+      FileUtils.deleteQuietly(unsupportedIndex);
+      FileUtils.deleteQuietly(unrelatedFile);
+      FileUtils.deleteDirectory(keptDir);
+    }
+  }
+
+  @Test
+  public void testCleanupOutputIndexesStopsWhenInterrupted() throws Throwable {
+    List<File> removedIndexes = Lists.newArrayList();
+    for (int i = 0; i < 3; i++) {
+      File removedDir = Files.createTempDirectory(stagingRoot.toPath(), 
"removed").toFile();
+      File removedIndex = outputIndexFile(runSucceededJob(removedDir));
+      FileUtils.deleteDirectory(removedDir);
+      ageOutputIndex(removedIndex);
+      removedIndexes.add(removedIndex);
+    }
+
+    // close() interrupts the cleanup thread. Each read of an index would then 
fail, and log a
+    // warning, if the cleanup didn't stop.
+    List<String> warnings =
+        captureWarnings(
+            () -> {
+              Thread.currentThread().interrupt();
+              try {
+                ((LocalJobExecutor) jobExecutor).cleanupOutputIndexes();
+              } finally {
+                Assertions.assertTrue(Thread.interrupted());
+              }
+            });
+    Assertions.assertEquals(Collections.emptyList(), warnings);
+    removedIndexes.forEach(index -> Assertions.assertTrue(index.exists()));
+
+    ((LocalJobExecutor) jobExecutor).cleanupOutputIndexes();
+    removedIndexes.forEach(index -> Assertions.assertFalse(index.exists()));
+  }
+
+  @Test
+  public void testInvalidOutputIndexIsWarnedOnce() throws Throwable {
+    String jobId = runSucceededJob(workingDir);
+    FileUtils.writeStringToFile(
+        outputIndexFile(jobId), "{\"version\":2,\"workingDir\":\"x\"}", 
StandardCharsets.UTF_8);
+
+    // A fresh executor, which hasn't warned about an invalid index yet. A 
client polling the
+    // output must not flood the log.
+    LocalJobExecutor exec = new LocalJobExecutor();
+    try {
+      exec.initialize(withStagingDir(Collections.emptyMap()));
+      List<String> warnings =
+          captureWarnings(
+              () -> {
+                for (int i = 0; i < 3; i++) {
+                  Assertions.assertEquals(
+                      Collections.emptyList(),
+                      exec.getJobStdout(jobId, 100, DEFAULT_TEST_MAX_BYTES));
+                }
+              });
+      Assertions.assertEquals(1, warnings.size(), warnings.toString());
+      Assertions.assertTrue(warnings.get(0).contains("is invalid or 
unsupported"), warnings.get(0));
+    } finally {
+      exec.close();
+    }
+  }
+
+  @Test
+  public void testCleanupOutputIndexesSkipsDirectories() throws Throwable {
+    // E.g. the staging directory of a metalake created before metalake names 
were checked.
+    File directory = outputIndexFile("local-job-00000000-" + 
UUID.randomUUID());
+    Assertions.assertTrue(directory.mkdirs());
+    ageOutputIndex(directory);
+    try {
+      List<String> warnings =
+          captureWarnings(() -> ((LocalJobExecutor) 
jobExecutor).cleanupOutputIndexes());
+      Assertions.assertEquals(Collections.emptyList(), warnings);
+      Assertions.assertTrue(directory.isDirectory());
+    } finally {
+      FileUtils.deleteDirectory(directory);
+    }
+  }
+
+  @Test
+  public void testCleanupOutputIndexesWithoutIndexDir() throws IOException {
+    FileUtils.deleteDirectory(new File(stagingRoot, OUTPUT_INDEX_DIR_NAME));
+    Assertions.assertDoesNotThrow(() -> ((LocalJobExecutor) 
jobExecutor).cleanupOutputIndexes());
+  }
+
+  @Test
+  public void testSubmitJobSucceedsWhenOutputIndexCannotBeWritten() throws 
IOException {
+    // A regular file where the index directory is expected makes writing the 
index fail.
+    File indexDir = new File(stagingRoot, OUTPUT_INDEX_DIR_NAME);
+    FileUtils.deleteDirectory(indexDir);
+    FileUtils.writeStringToFile(indexDir, "not a directory", 
StandardCharsets.UTF_8);
+    try {
+      String jobId = runSucceededJob(workingDir);
+      Assertions.assertTrue(new File(workingDir, "output.log").length() > 0);
+      Assertions.assertFalse(outputIndexFile(jobId).exists());
+    } finally {
+      FileUtils.deleteQuietly(indexDir);
+    }
+  }
+
+  @Test
+  public void testGetJobOutputFailsWhenOutputIndexCannotBeRead() throws 
IOException {
+    String jobId = runSucceededJob(workingDir);
+    // A directory where the index file is expected can't be read, unlike a 
missing index.
+    File indexFile = outputIndexFile(jobId);
+    Assertions.assertTrue(indexFile.delete());
+    Assertions.assertTrue(indexFile.mkdir());
+    try {
+      Assertions.assertThrows(
+          RuntimeException.class,
+          () -> jobExecutor.getJobStdout(jobId, 100, DEFAULT_TEST_MAX_BYTES));
+    } finally {
+      FileUtils.deleteDirectory(indexFile);
+    }
+  }
+
+  @Test
+  public void testOutputIndexDirIsRecreatedWhenRemoved() throws IOException {
+    FileUtils.deleteDirectory(new File(stagingRoot, OUTPUT_INDEX_DIR_NAME));
+
+    String jobId = runSucceededJob(workingDir);
+    Assertions.assertTrue(outputIndexFile(jobId).exists());
+    Assertions.assertEquals(6, jobExecutor.getJobStdout(jobId, 100, 
DEFAULT_TEST_MAX_BYTES).size());
+  }
+
+  private static Map<String, String> withStagingDir(Map<String, String> 
configs) {
+    return ImmutableMap.<String, String>builder()
+        .putAll(configs)
+        .put(LocalJobExecutorConfigs.STAGING_DIR, 
stagingRoot.getAbsolutePath())
+        .build();
+  }
+
+  private static JobTemplate newRuntimeJobTemplate(File jobDir) {
+    return JobManager.createRuntimeJobTemplate(
+        jobTemplateEntity,
+        ImmutableMap.of("arg1", "value1", "arg2", "success", "var", "value3"),
+        jobDir);
+  }
+
+  private static String runSucceededJob(File jobDir) {
+    String jobId = jobExecutor.submitJob(newRuntimeJobTemplate(jobDir));
+    Awaitility.await()
+        .atMost(3, TimeUnit.MINUTES)
+        .until(() -> jobExecutor.getJobStatus(jobId) == 
JobHandle.Status.SUCCEEDED);
+    return jobId;
+  }
+
+  private static boolean hasJobStatus(JobExecutor executor, String jobId) {
+    try {
+      executor.getJobStatus(jobId);
+      return true;
+    } catch (NoSuchJobException e) {
+      return false;
+    }
+  }
+
+  private static File outputIndexFile(String jobId) {
+    return new File(stagingRoot, OUTPUT_INDEX_DIR_NAME + File.separator + 
jobId + ".json");
+  }
+
+  // Returns the warnings LocalJobExecutor logs while running the action.
+  private static List<String> captureWarnings(Executable action) throws 
Throwable {
+    LoggerContext context =
+        (LoggerContext) 
LogManager.getContext(LocalJobExecutor.class.getClassLoader(), false);
+    AbstractConfiguration configuration = (AbstractConfiguration) 
context.getConfiguration();
+    WarningCollector collector = new WarningCollector();
+    collector.start();
+    configuration.addAppender(collector);
+    LoggerConfig loggerConfig =
+        new LoggerConfig(LocalJobExecutor.class.getName(), Level.WARN, false);
+    loggerConfig.addAppender(collector, Level.WARN, null);
+    configuration.addLogger(LocalJobExecutor.class.getName(), loggerConfig);
+    context.updateLoggers();
+    try {
+      action.execute();
+    } finally {
+      configuration.removeLogger(LocalJobExecutor.class.getName());
+      collector.stop();
+      configuration.removeAppender(collector.getName());
+      context.updateLoggers();
+    }
+    return ImmutableList.copyOf(collector.messages);
+  }
+
+  private static void replaceWithSymlink(File file, File target) throws 
IOException {
+    FileUtils.forceDelete(file);
+    Files.createSymbolicLink(file.toPath(), target.toPath().toAbsolutePath());
+  }
+
+  // Makes the index old enough for the cleanup to consider it.
+  private static void ageOutputIndex(File indexFile) throws IOException {
+    Files.setLastModifiedTime(
+        indexFile.toPath(),
+        FileTime.fromMillis(System.currentTimeMillis() - 
TimeUnit.HOURS.toMillis(2)));
+  }
+
+  private static int outputIndexFileCount() {
+    File[] indexFiles = new File(stagingRoot, 
OUTPUT_INDEX_DIR_NAME).listFiles();
+    return indexFiles == null ? 0 : indexFiles.length;
+  }
+
+  private static class WarningCollector extends AbstractAppender {
+    private final List<String> messages = 
Collections.synchronizedList(Lists.newArrayList());
+
+    WarningCollector() {
+      super("localJobExecutorWarnings", null, 
PatternLayout.createDefaultLayout(), true, null);
+    }
+
+    @Override
+    public void append(LogEvent event) {
+      messages.add(event.getMessage().getFormattedMessage());
+    }
+  }
+
   private JobTemplate newSleepJobTemplate(String name) throws IOException {
     // The job runs in the directory of its executable, so give each job its 
own directory.
     File jobDir = new File(workingDir, name);
diff --git a/docs/gravitino-server-config.md b/docs/gravitino-server-config.md
index 2bb8a90d08..1fb42dea4a 100644
--- a/docs/gravitino-server-config.md
+++ b/docs/gravitino-server-config.md
@@ -146,6 +146,11 @@ a three second poll, and a server that cannot keep its 
caches current exits rath
 metadata it knows to be stale. Point the load balancer's health check at `GET 
/health/ready` so a
 server that has lost its database stops receiving traffic.
 
+Jobs run by the default `local` job executor keep their output in 
`gravitino.job.stagingDir`. Put
+that directory on storage shared by all servers, for example an NFS mount, so 
that a request for a
+job's output can be served by any server. Otherwise only the server that ran 
the job can return
+it, and the others return empty output. See [Manage 
Jobs](manage-jobs-in-gravitino.md).
+
 ## Server Configuration
 
 Every property in this section belongs in 
`${GRAVITINO_HOME}/conf/gravitino.conf`, one
@@ -541,14 +546,14 @@ server, are documented with those services. See
 
 #### Jobs
 
-| Configuration Item                     | Description                         
                                                                       | 
Default Value                 |
-|----------------------------------------|------------------------------------------------------------------------------------------------------------|-------------------------------|
-| `gravitino.job.executor`               | Executor that runs jobs. Implement 
your own and name it here to replace the built-in one.                  | 
`local`                       |
-| `gravitino.job.stagingDir`             | Directory holding staging files for 
running jobs.                                                          | 
`/tmp/gravitino/jobs/staging` |
-| `gravitino.job.stagingDirKeepTimeInMs` | How long in milliseconds a finished 
job's staging files are kept. Use at least 10 minutes outside testing. | 
`604800000` (7 days)          |
-| `gravitino.job.statusPullIntervalInMs` | Interval in milliseconds between 
job status polls. Use at least 1 minute outside testing.                  | 
`300000` (5 minutes)          |
-| `gravitino.job.outputMaxLines`         | Maximum number of lines returned 
when fetching a job's stdout/stderr output.                               | 
`1000`                        |
-| `gravitino.job.outputMaxBytes`         | Maximum number of bytes read from 
the tail of a job's stdout/stderr when fetching its output.              | 
`262144` (256KB)              |
+| Configuration Item                     | Description                         
                                                                                
                                           | Default Value                 |
+|----------------------------------------|----------------------------------------------------------------------------------------------------------------------------------------------------------------|-------------------------------|
+| `gravitino.job.executor`               | Executor that runs jobs. Implement 
your own and name it here to replace the built-in one.                          
                                            | `local`                       |
+| `gravitino.job.stagingDir`             | Directory holding staging files for 
running jobs. With multiple servers, put it on storage shared by all servers so 
that any server can return a job's output. | `/tmp/gravitino/jobs/staging` |
+| `gravitino.job.stagingDirKeepTimeInMs` | How long in milliseconds a finished 
job's staging files are kept. Use at least 10 minutes outside testing.          
                                           | `604800000` (7 days)          |
+| `gravitino.job.statusPullIntervalInMs` | Interval in milliseconds between 
job status polls. Use at least 1 minute outside testing.                        
                                              | `300000` (5 minutes)          |
+| `gravitino.job.outputMaxLines`         | Maximum number of lines returned 
when fetching a job's stdout/stderr output.                                     
                                              | `1000`                        |
+| `gravitino.job.outputMaxBytes`         | Maximum number of bytes read from 
the tail of a job's stdout/stderr when fetching its output.                     
                                             | `262144` (256KB)              |
 
 ### Key Management
 
diff --git a/docs/manage-jobs-in-gravitino.md b/docs/manage-jobs-in-gravitino.md
index 7167f3e67c..019e273ea3 100644
--- a/docs/manage-jobs-in-gravitino.md
+++ b/docs/manage-jobs-in-gravitino.md
@@ -260,9 +260,10 @@ stderr = job.stderr()
 </TabItem>
 </Tabs>
 
-Output is only kept for as long as the job executor retains it - for the local 
job executor, that's
-tied to `gravitino.jobExecutor.local.jobStatusKeepTimeInMs` below, and it's 
lost entirely across a
-server restart. What's returned is always the tail of the output (the most 
recent content), capped
+Output is only kept for as long as the job executor retains it. The local job 
executor keeps it in
+the job's staging directory, so it's available until the staging directory is 
cleaned up
+(`gravitino.job.stagingDirKeepTimeInMs` after the job finishes), also across 
server restarts.
+What's returned is always the tail of the output (the most recent content), 
capped
 by `gravitino.job.outputMaxLines` (line count) and 
`gravitino.job.outputMaxBytes` (byte size),
 whichever limit is hit first.
 
@@ -276,19 +277,38 @@ 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
+  yet, have no index file and return empty output.
+- Deleting this directory makes the output of existing jobs unavailable. After 
downgrading to an
+  earlier version, it isn't used anymore and can be removed.
+- A missing index returns empty output, but an index that exists and can't be 
read, for example
+  while the shared storage is unavailable, fails the request with an error 
rather than returning
+  empty output that looks like the job printed nothing.
+
 ### Job System Configuration
 
 Configure the job system through the `gravitino.conf` file. The following are 
the
 default configurations:
 
-| Property name                          | Description                         
                                              | Default value                 | 
Required |
-|----------------------------------------|-----------------------------------------------------------------------------------|-------------------------------|----------|
-| `gravitino.job.stagingDir`             | Directory for managing the staging 
files when running jobs                        | `/tmp/gravitino/jobs/staging` 
| No       |
-| `gravitino.job.executor`               | The job executor to use for running 
jobs                                          | `local`                       | 
No       |
-| `gravitino.job.stagingDirKeepTimeInMs` | The time in milliseconds to keep 
the staging directory after the job is completed | `604800000` (7 days)         
 | No       |
-| `gravitino.job.statusPullIntervalInMs` | The interval in milliseconds to 
pull the job status from the job executor         | `300000` (5 minutes)        
  | No       |
-| `gravitino.job.outputMaxLines`         | The maximum number of lines 
returned when fetching a job's stdout/stderr output   | `1000`                  
      | No       |
-| `gravitino.job.outputMaxBytes`         | The maximum number of bytes read 
from the tail of a job's stdout/stderr output    | `262144` (256KB)             
 | No       |
+| Property name                          | Description                         
                                                                                
                                                 | Default value                
 | Required |
+|----------------------------------------|----------------------------------------------------------------------------------------------------------------------------------------------------------------------|-------------------------------|----------|
+| `gravitino.job.stagingDir`             | Directory for managing the staging 
files when running jobs. Must be shared by all servers in a multi-server 
deployment, see [Get a Job's Output](#get-a-jobs-output) | 
`/tmp/gravitino/jobs/staging` | No       |
+| `gravitino.job.executor`               | The job executor to use for running 
jobs                                                                            
                                                 | `local`                      
 | No       |
+| `gravitino.job.stagingDirKeepTimeInMs` | The time in milliseconds to keep 
the staging directory after the job is completed                                
                                                    | `604800000` (7 days)      
    | No       |
+| `gravitino.job.statusPullIntervalInMs` | The interval in milliseconds to 
pull the job status from the job executor                                       
                                                     | `300000` (5 minutes)     
     | No       |
+| `gravitino.job.outputMaxLines`         | The maximum number of lines 
returned when fetching a job's stdout/stderr output                             
                                                         | `1000`               
         | No       |
+| `gravitino.job.outputMaxBytes`         | The maximum number of bytes read 
from the tail of a job's stdout/stderr output                                   
                                                    | `262144` (256KB)          
    | No       |
 
 #### Configurations for Local Job Executor
 
@@ -302,6 +322,9 @@ The following are the default configurations for the local 
job executor:
 | `gravitino.jobExecutor.local.jobStatusKeepTimeInMs` | The time in 
milliseconds to keep the job status in the local job executor                   
                                                      | `3600000` (1 hour)      
               | No       |
 | `gravitino.jobExecutor.local.sparkHome`             | The home directory of 
Spark, Gravitino checks this configuration firstly and then `SPARK_HOME` env. 
Either of them should be set to run Spark job | `None`                          
       | No       |
 
+The local job executor always uses `gravitino.job.stagingDir` as its staging 
directory, the same one
+the job system stages jobs in. A `gravitino.jobExecutor.local.stagingDir` 
setting is ignored.
+
 The local job executor runs up to `gravitino.jobExecutor.local.maxRunningJobs` 
jobs at the same
 time, each in its own process on the Gravitino server host, and queues the 
others. Make sure the
 host has enough resources for that many jobs, or lower this value, especially 
when running Spark
@@ -315,6 +338,8 @@ only tracks the jobs it runs itself:
 - Cancelling a job on a server that doesn't run it marks the job as 
`CANCELLING`, and the server
   running the job cancels it the next time it pulls job statuses. This can 
take up to
   `gravitino.job.statusPullIntervalInMs`.
+- A job's output can be read from any server only if 
`gravitino.job.stagingDir` is shared by all
+  servers, see [Get a Job's Output](#get-a-jobs-output).
 - If a server exits while running jobs, nobody can track these jobs anymore. 
When such a job has
   not been updated for `gravitino.job.stagingDirKeepTimeInMs`, it is marked as 
`FAILED`, or as
   `CANCELLED` if it was being cancelled. Like other finished jobs, it is then 
kept for another
diff --git a/docs/open-api/jobs.yaml b/docs/open-api/jobs.yaml
index 342ac6f4e2..856a999fb1 100644
--- a/docs/open-api/jobs.yaml
+++ b/docs/open-api/jobs.yaml
@@ -370,6 +370,9 @@ components:
         Whether to also fetch and include the job's captured stdout/stderr 
output. Output is
         fetched live from the job executor on every request, so this should 
only be set to true
         when the caller actually needs it - it is never included when this 
parameter is omitted.
+        With the local job executor and multiple Gravitino servers, 
gravitino.job.stagingDir must
+        be shared by all servers for any server to return the output; 
otherwise a server that
+        didn't run the job returns empty output.
       required: false
       schema:
         type: boolean

Reply via email to