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 bf5fafe663 [#12716] feat(job): support output log retrieval for job
handler (#13250)
bf5fafe663 is described below
commit bf5fafe663d38dc7c9e220e6b3c943d1d0cbf8bf
Author: Jerry Shao <[email protected]>
AuthorDate: Mon Sep 21 19:32:05 2026 +0800
[#12716] feat(job): support output log retrieval for job handler (#13250)
### What changes were proposed in this pull request?
Adds the ability to retrieve a job's captured stdout/stderr output, end
to end:
- `JobExecutor#getJobStdout`/`getJobStderr(jobId, maxLines)` β new
default methods (empty list by
default), implemented by `LocalJobExecutor` via a bounded tail-read of
the job's
`output.log`/`error.log` (no full-file load).
- Core: `JobOperationDispatcher`/`JobManager`'s `getJob` gains an
`includeOutput` parameter;
`JobEntity` carries non-persisted `stdout`/`stderr` fields when
requested.
- REST: the existing `GET .../jobs/runs/{jobId}` endpoint gains an
`includeOutput` query
parameter (no new endpoint), so `listJobs`/plain `getJob` never carry
the extra payload.
- Java/Python clients: `SupportsJobs#getJob(jobId, includeOutput)`
(Java: default method
delegating to the new overload; Python: `get_job(job_id,
include_output=False)`), and
`JobHandle#stdout()`/`stderr()`.
- New global config `gravitino.job.outputMaxLines` (default 1000),
resolved by `JobManager` and
passed to the executor per call.
- Also fixes `ExceptionHandlers` not mapping
`UnsupportedOperationException` for job and
job-template operations, and makes `getJob(..., includeOutput=true)`
degrade gracefully to
empty output (instead of a false "job not found") when the job entity
still exists but the
executor's own output bookkeeping has expired.
- Docs: `gravitino.job.outputMaxLines` added to both config reference
tables, plus a new
"Get a Job's Output" section with REST/Java/Python examples.
### Why are the changes needed?
`JobExecutor` only exposed `submitJob`/`getJobStatus`/`cancelJob` β
there was no way to see why a
job failed without directly inspecting the local runner's staging
directory. This closes that
gap for the built-in local executor and lays the API surface for other
executors to do the same.
Fix: #12716 (subtask of epic #12667)
### Does this PR introduce _any_ user-facing change?
Yes:
- New optional `includeOutput` query parameter on `GET
.../jobs/runs/{jobId}`.
- New Java `SupportsJobs#getJob(String jobId, boolean includeOutput)`
and `JobHandle#stdout()`/
`stderr()`; new Python `get_job(job_id, include_output=False)` and
`JobHandle.stdout()`/
`stderr()`.
- New server config property `gravitino.job.outputMaxLines` (default
`1000`).
### How was this patch tested?
Added unit tests across `TestLocalJobExecutor`, `TestJobManager`,
`TestJobOperations`,
`TestJobDTO`, `TestSupportsJobs` (Java), and Python
`test_job_dto_serde`/`test_supports_jobs`.
Added an end-to-end integration test (`JobIT.testRunAndGetJobOutput`)
that runs a real local job,
polls to completion, and asserts the actual captured output. Full
`:core:test`/`:server:test`
suites re-verified green (2140+ tests) after a full automated code
review and fixes.
π€ Generated with [Claude Code](https://claude.com/claude-code)
Co-authored-by: Claude Sonnet 5 <[email protected]>
---
.../java/org/apache/gravitino/job/JobHandle.java | 28 ++
.../org/apache/gravitino/job/SupportsJobs.java | 29 ++
.../org/apache/gravitino/job/TestSupportsJobs.java | 133 +++++++++
.../apache/gravitino/client/GenericJobHandle.java | 12 +
.../apache/gravitino/client/GravitinoClient.java | 5 +
.../apache/gravitino/client/GravitinoMetalake.java | 8 +
.../apache/gravitino/client/TestSupportsJobs.java | 55 +++-
.../gravitino/client/integration/test/JobIT.java | 43 +++
.../client-python/gravitino/api/job/job_handle.py | 18 +-
.../gravitino/api/job/supports_jobs.py | 9 +-
.../gravitino/client/generic_job_handle.py | 10 +
.../gravitino/client/gravitino_client.py | 7 +-
.../gravitino/client/gravitino_metalake.py | 13 +-
clients/client-python/gravitino/dto/job/job_dto.py | 20 +-
.../tests/integration/test_supports_jobs.py | 40 +++
.../tests/unittests/dto/job/test_job_dto_serde.py | 26 ++
.../tests/unittests/test_supports_jobs.py | 36 +++
.../java/org/apache/gravitino/dto/job/JobDTO.java | 22 +-
.../org/apache/gravitino/dto/job/TestJobDTO.java | 45 ++-
.../main/java/org/apache/gravitino/Configs.java | 25 ++
.../gravitino/connector/job/JobExecutor.java | 46 +++
.../apache/gravitino/hook/JobHookDispatcher.java | 6 +-
.../java/org/apache/gravitino/job/JobManager.java | 83 +++++-
.../gravitino/job/JobOperationDispatcher.java | 41 ++-
.../job/JobTemplateValidationDispatcher.java | 6 +-
.../gravitino/job/local/LocalJobExecutor.java | 128 +++++++-
.../gravitino/job/local/LocalProcessBuilder.java | 21 +-
.../gravitino/job/local/ShellProcessBuilder.java | 4 +-
.../gravitino/job/local/SparkProcessBuilder.java | 4 +-
.../gravitino/listener/JobEventDispatcher.java | 7 +-
.../java/org/apache/gravitino/meta/JobEntity.java | 64 ++++
.../apache/gravitino/utils/MetadataObjectUtil.java | 2 +-
.../org/apache/gravitino/job/TestJobManager.java | 155 +++++++++-
.../gravitino/job/TestJobManagerMultiNode.java | 2 +-
.../job/TestJobTemplateValidationDispatcher.java | 4 +-
.../gravitino/job/local/TestLocalJobExecutor.java | 331 +++++++++++++++++++++
.../listener/api/event/TestJobEventDispatcher.java | 8 +-
.../gravitino/utils/TestMetadataObjectUtil.java | 2 +-
docs/gravitino-server-config.md | 2 +
docs/manage-jobs-in-gravitino.md | 54 ++++
docs/open-api/jobs.yaml | 57 ++++
.../server/web/rest/ExceptionHandlers.java | 6 +
.../gravitino/server/web/rest/JobOperations.java | 13 +-
.../server/web/rest/TestJobOperations.java | 140 +++++++++
44 files changed, 1699 insertions(+), 71 deletions(-)
diff --git a/api/src/main/java/org/apache/gravitino/job/JobHandle.java
b/api/src/main/java/org/apache/gravitino/job/JobHandle.java
index 7da6129df7..d9463fb053 100644
--- a/api/src/main/java/org/apache/gravitino/job/JobHandle.java
+++ b/api/src/main/java/org/apache/gravitino/job/JobHandle.java
@@ -19,6 +19,8 @@
package org.apache.gravitino.job;
import java.time.Instant;
+import java.util.Collections;
+import java.util.List;
import javax.annotation.Nullable;
/**
@@ -118,4 +120,30 @@ public interface JobHandle {
+ getClass().getName()
+ "; override this method");
}
+
+ /**
+ * Get the captured standard output of the job, as a list of lines. This is
only populated when
+ * the handle was obtained via {@link SupportsJobs#getJob(String, boolean)}
with {@code
+ * includeOutput=true}; handles obtained via the plain {@link
SupportsJobs#getJob(String)}, {@link
+ * SupportsJobs#listJobs()}, or {@link SupportsJobs#runJob(String,
java.util.Map)} always return
+ * an empty list here.
+ *
+ * @return the stdout lines of the job, or an empty list if not available
+ */
+ default List<String> stdout() {
+ return Collections.emptyList();
+ }
+
+ /**
+ * Get the captured standard error output of the job, as a list of lines.
This is only populated
+ * when the handle was obtained via {@link SupportsJobs#getJob(String,
boolean)} with {@code
+ * includeOutput=true}; handles obtained via the plain {@link
SupportsJobs#getJob(String)}, {@link
+ * SupportsJobs#listJobs()}, or {@link SupportsJobs#runJob(String,
java.util.Map)} always return
+ * an empty list here.
+ *
+ * @return the stderr lines of the job, or an empty list if not available
+ */
+ default List<String> stderr() {
+ return Collections.emptyList();
+ }
}
diff --git a/api/src/main/java/org/apache/gravitino/job/SupportsJobs.java
b/api/src/main/java/org/apache/gravitino/job/SupportsJobs.java
index b5ca1698e2..5bea40c5ab 100644
--- a/api/src/main/java/org/apache/gravitino/job/SupportsJobs.java
+++ b/api/src/main/java/org/apache/gravitino/job/SupportsJobs.java
@@ -126,6 +126,35 @@ public interface SupportsJobs {
*/
JobHandle getJob(String jobId) throws NoSuchJobException;
+ /**
+ * Retrieves a job by its ID, optionally including its captured
stdout/stderr output (see {@link
+ * JobHandle#stdout()}/{@link JobHandle#stderr()}).
+ *
+ * <p>Output is fetched live from the job executor on every call, not
persisted, so {@code
+ * includeOutput} should only be set to {@code true} when the output is
actually needed.
+ *
+ * <p>The default implementation delegates to {@link #getJob(String)} when
{@code includeOutput}
+ * is {@code false} - equivalent to a plain lookup, so existing implementors
of this interface
+ * keep compiling and behaving correctly without any changes - but throws
{@link
+ * UnsupportedOperationException} when {@code includeOutput} is {@code
true}, rather than silently
+ * ignoring the request and returning a job with no output. Implementors
that support output
+ * retrieval should override this method directly.
+ *
+ * @param jobId the ID of the job to retrieve
+ * @param includeOutput whether to also fetch and populate the job's
stdout/stderr output
+ * @return a handle to the job
+ * @throws NoSuchJobException if the job with the specified ID does not exist
+ * @throws UnsupportedOperationException if {@code includeOutput} is {@code
true} and this
+ * implementor does not support output retrieval
+ */
+ default JobHandle getJob(String jobId, boolean includeOutput) throws
NoSuchJobException {
+ if (!includeOutput) {
+ return getJob(jobId);
+ }
+ throw new UnsupportedOperationException(
+ "getJob(jobId, includeOutput=true) is not supported by " +
getClass().getName());
+ }
+
/**
* Cancel a job by its ID. This operation will attempt to cancel the job if
it is still running.
* This method will return immediately, user could use the job handle to
check the status of the
diff --git a/api/src/test/java/org/apache/gravitino/job/TestSupportsJobs.java
b/api/src/test/java/org/apache/gravitino/job/TestSupportsJobs.java
new file mode 100644
index 0000000000..de13a14926
--- /dev/null
+++ b/api/src/test/java/org/apache/gravitino/job/TestSupportsJobs.java
@@ -0,0 +1,133 @@
+/*
+ * 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 java.util.Collections;
+import java.util.List;
+import java.util.Map;
+import org.apache.gravitino.exceptions.InUseException;
+import org.apache.gravitino.exceptions.JobTemplateAlreadyExistsException;
+import org.apache.gravitino.exceptions.NoSuchJobException;
+import org.apache.gravitino.exceptions.NoSuchJobTemplateException;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+public class TestSupportsJobs {
+
+ /**
+ * Mimics a {@link SupportsJobs} implementor written before the {@code
includeOutput} overload was
+ * introduced - it only implements {@link #getJob(String)}, the sole
abstract method both before
+ * and after that change, so it must keep compiling and behaving correctly
without any
+ * modifications.
+ */
+ private static class LegacyJobsImpl implements SupportsJobs {
+
+ private final JobHandle handle;
+
+ LegacyJobsImpl(JobHandle handle) {
+ this.handle = handle;
+ }
+
+ @Override
+ public List<JobTemplate> listJobTemplates() {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public void registerJobTemplate(JobTemplate jobTemplate)
+ throws JobTemplateAlreadyExistsException {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public JobTemplate getJobTemplate(String jobTemplateName) throws
NoSuchJobTemplateException {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public boolean deleteJobTemplate(String jobTemplateName) throws
InUseException {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public List<JobHandle> listJobs(String jobTemplateName) throws
NoSuchJobTemplateException {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public List<JobHandle> listJobs() {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public JobHandle runJob(String jobTemplateName, Map<String, String>
jobConf)
+ throws NoSuchJobTemplateException {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public JobHandle getJob(String jobId) throws NoSuchJobException {
+ return handle;
+ }
+
+ @Override
+ public JobHandle cancelJob(String jobId) throws NoSuchJobException {
+ throw new UnsupportedOperationException();
+ }
+ }
+
+ private static class FakeJobHandle implements JobHandle {
+ @Override
+ public String jobTemplateName() {
+ return "template";
+ }
+
+ @Override
+ public String jobId() {
+ return "job-1";
+ }
+
+ @Override
+ public Status jobStatus() {
+ return Status.SUCCEEDED;
+ }
+ }
+
+ @Test
+ public void testLegacyImplementorStillCompilesAndBehavesCorrectly() throws
NoSuchJobException {
+ JobHandle handle = new FakeJobHandle();
+ SupportsJobs legacy = new LegacyJobsImpl(handle);
+
+ // The plain lookup, and the two-arg overload with includeOutput=false,
both delegate to the
+ // legacy implementor's only method and behave identically.
+ Assertions.assertSame(handle, legacy.getJob("job-1"));
+ Assertions.assertSame(handle, legacy.getJob("job-1", false));
+
+ // Requesting output from an implementor that never opted in throws a
clear signal rather
+ // than silently ignoring the request and returning a handle with no
output.
+ Assertions.assertThrows(
+ UnsupportedOperationException.class, () -> legacy.getJob("job-1",
true));
+ }
+
+ @Test
+ public void testDefaultJobHandleStdoutAndStderrAreEmpty() {
+ JobHandle handle = new FakeJobHandle();
+ Assertions.assertEquals(Collections.emptyList(), handle.stdout());
+ Assertions.assertEquals(Collections.emptyList(), handle.stderr());
+ }
+}
diff --git
a/clients/client-java/src/main/java/org/apache/gravitino/client/GenericJobHandle.java
b/clients/client-java/src/main/java/org/apache/gravitino/client/GenericJobHandle.java
index beb370b589..eb7f5f31ac 100644
---
a/clients/client-java/src/main/java/org/apache/gravitino/client/GenericJobHandle.java
+++
b/clients/client-java/src/main/java/org/apache/gravitino/client/GenericJobHandle.java
@@ -19,6 +19,8 @@
package org.apache.gravitino.client;
import java.time.Instant;
+import java.util.Collections;
+import java.util.List;
import org.apache.gravitino.dto.job.JobDTO;
import org.apache.gravitino.dto.util.DTOConverters;
import org.apache.gravitino.job.JobHandle;
@@ -69,4 +71,14 @@ public class GenericJobHandle implements JobHandle {
? null
: DTOConverters.fromDTO(jobDTO.runtimeJobTemplate());
}
+
+ @Override
+ public List<String> stdout() {
+ return jobDTO.stdout() == null ? Collections.emptyList() : jobDTO.stdout();
+ }
+
+ @Override
+ public List<String> stderr() {
+ return jobDTO.stderr() == null ? Collections.emptyList() : jobDTO.stderr();
+ }
}
diff --git
a/clients/client-java/src/main/java/org/apache/gravitino/client/GravitinoClient.java
b/clients/client-java/src/main/java/org/apache/gravitino/client/GravitinoClient.java
index 1ac10b77d0..4fbd927712 100644
---
a/clients/client-java/src/main/java/org/apache/gravitino/client/GravitinoClient.java
+++
b/clients/client-java/src/main/java/org/apache/gravitino/client/GravitinoClient.java
@@ -706,6 +706,11 @@ public class GravitinoClient extends GravitinoClientBase
return getMetalake().getJob(jobId);
}
+ @Override
+ public JobHandle getJob(String jobId, boolean includeOutput) throws
NoSuchJobException {
+ return getMetalake().getJob(jobId, includeOutput);
+ }
+
@Override
public JobHandle cancelJob(String jobId) throws NoSuchJobException {
return getMetalake().cancelJob(jobId);
diff --git
a/clients/client-java/src/main/java/org/apache/gravitino/client/GravitinoMetalake.java
b/clients/client-java/src/main/java/org/apache/gravitino/client/GravitinoMetalake.java
index d03f1922be..42149cbefb 100644
---
a/clients/client-java/src/main/java/org/apache/gravitino/client/GravitinoMetalake.java
+++
b/clients/client-java/src/main/java/org/apache/gravitino/client/GravitinoMetalake.java
@@ -1816,13 +1816,21 @@ public class GravitinoMetalake extends MetalakeDTO
@Override
public JobHandle getJob(String jobId) throws NoSuchJobException {
+ return getJob(jobId, false);
+ }
+
+ @Override
+ public JobHandle getJob(String jobId, boolean includeOutput) throws
NoSuchJobException {
Preconditions.checkArgument(StringUtils.isNotBlank(jobId), "job id must
not be null or empty");
+ Map<String, String> params =
+ includeOutput ? ImmutableMap.of("includeOutput", "true") :
Collections.emptyMap();
JobResponse resp =
restClient.get(
String.format(API_METALAKES_JOB_PATH,
RESTUtils.encodeString(this.name()))
+ "/"
+ RESTUtils.encodeString(jobId),
+ params,
JobResponse.class,
Collections.emptyMap(),
ErrorHandlers.jobErrorHandler());
diff --git
a/clients/client-java/src/test/java/org/apache/gravitino/client/TestSupportsJobs.java
b/clients/client-java/src/test/java/org/apache/gravitino/client/TestSupportsJobs.java
index ff6db13a09..d246cdbf3d 100644
---
a/clients/client-java/src/test/java/org/apache/gravitino/client/TestSupportsJobs.java
+++
b/clients/client-java/src/test/java/org/apache/gravitino/client/TestSupportsJobs.java
@@ -337,6 +337,41 @@ public class TestSupportsJobs extends TestBase {
Assertions.assertEquals(DTOConverters.fromDTO(runtimeJobTemplateDTO),
runtimeJobTemplate);
}
+ @Test
+ public void testGetJobWithOutput() throws JsonProcessingException {
+ String jobId = "job-1";
+ String jobTemplateName = "shell-job-template";
+ List<String> stdout = Lists.newArrayList("line1", "line2");
+ List<String> stderr = Lists.newArrayList("err1");
+ JobDTO expectedJob = newJobDTO(jobId, jobTemplateName, stdout, stderr);
+ JobResponse resp = new JobResponse(expectedJob);
+
+ buildMockResource(
+ Method.GET,
+ jobRunsPath() + "/" + jobId,
+ ImmutableMap.of("includeOutput", "true"),
+ null,
+ resp,
+ HttpStatus.SC_OK);
+
+ JobHandle actualHandle = metalake.getJob(jobId, true);
+ compare(expectedJob, actualHandle);
+ Assertions.assertEquals(stdout, actualHandle.stdout());
+ Assertions.assertEquals(stderr, actualHandle.stderr());
+
+ // Test throw NoSuchJobException
+ ErrorResponse errorResp =
+ ErrorResponse.notFound(NoSuchJobException.class.getSimpleName(), "mock
error");
+ buildMockResource(
+ Method.GET,
+ jobRunsPath() + "/" + jobId,
+ ImmutableMap.of("includeOutput", "true"),
+ null,
+ errorResp,
+ HttpStatus.SC_NOT_FOUND);
+ Assertions.assertThrows(NoSuchJobException.class, () ->
metalake.getJob(jobId, true));
+ }
+
@Test
public void testRunJob() throws JsonProcessingException {
String jobTemplateName = "shell-job-template";
@@ -459,6 +494,24 @@ public class TestSupportsJobs extends TestBase {
now,
startedAt,
finishedAt,
- runtimeJobTemplate);
+ runtimeJobTemplate,
+ null,
+ null);
+ }
+
+ private JobDTO newJobDTO(
+ String jobId, String templateName, List<String> stdout, List<String>
stderr) {
+ Instant now = Instant.now();
+ return new JobDTO(
+ jobId,
+ templateName,
+ JobHandle.Status.SUCCEEDED,
+ AuditDTO.builder().withCreator("test").withCreateTime(now).build(),
+ now,
+ now,
+ now,
+ null,
+ stdout,
+ stderr);
}
}
diff --git
a/clients/client-java/src/test/java/org/apache/gravitino/client/integration/test/JobIT.java
b/clients/client-java/src/test/java/org/apache/gravitino/client/integration/test/JobIT.java
index 380affdcde..08f38a9d24 100644
---
a/clients/client-java/src/test/java/org/apache/gravitino/client/integration/test/JobIT.java
+++
b/clients/client-java/src/test/java/org/apache/gravitino/client/integration/test/JobIT.java
@@ -511,6 +511,49 @@ public class JobIT extends BaseIT {
Assertions.assertEquals(runtimeJobTemplate,
retrievedJob.runtimeJobTemplate());
}
+ @Test
+ public void testRunAndGetJobOutput() {
+ JobTemplate template = builder.withName("test_run_get_output").build();
+ Assertions.assertDoesNotThrow(() ->
metalake.registerJobTemplate(template));
+
+ JobHandle jobHandle =
+ metalake.runJob(
+ template.name(),
+ ImmutableMap.of("arg1", "value1", "arg2", "success", "env_var",
"value2"));
+
+ // Plain getJob never carries output.
+ Assertions.assertTrue(jobHandle.stdout().isEmpty());
+ Assertions.assertTrue(jobHandle.stderr().isEmpty());
+
+ Awaitility.await()
+ .atMost(3, TimeUnit.MINUTES)
+ .until(
+ () -> {
+ JobHandle updatedJob = metalake.getJob(jobHandle.jobId());
+ return updatedJob.jobStatus() == JobHandle.Status.SUCCEEDED;
+ });
+
+ // getJob still never carries output, even after the job finishes.
+ JobHandle finishedJob = metalake.getJob(jobHandle.jobId());
+ Assertions.assertTrue(finishedJob.stdout().isEmpty());
+ Assertions.assertTrue(finishedJob.stderr().isEmpty());
+
+ // getJob(jobId, true) fetches the captured stdout/stderr.
+ JobHandle jobWithOutput = metalake.getJob(jobHandle.jobId(), true);
+ List<String> stdout = jobWithOutput.stdout();
+ Assertions.assertTrue(stdout.contains("starting test test job"));
+ Assertions.assertTrue(stdout.contains("in common script"));
+ Assertions.assertTrue(stdout.contains("value1"));
+ Assertions.assertTrue(stdout.contains("success"));
+ Assertions.assertTrue(stdout.contains("value2"));
+ // The test script never writes to stderr.
+ Assertions.assertTrue(jobWithOutput.stderr().isEmpty());
+
+ // Test get output for a non-existent job.
+ Assertions.assertThrows(
+ NoSuchJobException.class, () -> metalake.getJob("non_existent_job_id",
true));
+ }
+
@Test
public void testRunAndCancelJob() {
JobTemplate template = builder.withName("test_run_cancel").build();
diff --git a/clients/client-python/gravitino/api/job/job_handle.py
b/clients/client-python/gravitino/api/job/job_handle.py
index 4f9d749165..98dbda8d44 100644
--- a/clients/client-python/gravitino/api/job/job_handle.py
+++ b/clients/client-python/gravitino/api/job/job_handle.py
@@ -18,7 +18,7 @@
from abc import ABC, abstractmethod
from datetime import datetime
from enum import Enum
-from typing import Optional
+from typing import List, Optional
from gravitino.api.job.job_template import JobTemplate
@@ -79,3 +79,19 @@ class JobHandle(ABC):
this field was introduced.
"""
raise NotImplementedError("runtime_job_template is not implemented")
+
+ def stdout(self) -> List[str]:
+ """Returns the captured standard output of the job, as a list of
lines. This is only
+ populated when the handle was obtained via ``SupportsJobs.get_job``
with
+ ``include_output=True``; handles obtained via the plain ``get_job``,
``list_jobs``, or
+ ``run_job`` always return an empty list here.
+ """
+ return []
+
+ def stderr(self) -> List[str]:
+ """Returns the captured standard error output of the job, as a list of
lines. This is only
+ populated when the handle was obtained via ``SupportsJobs.get_job``
with
+ ``include_output=True``; handles obtained via the plain ``get_job``,
``list_jobs``, or
+ ``run_job`` always return an empty list here.
+ """
+ return []
diff --git a/clients/client-python/gravitino/api/job/supports_jobs.py
b/clients/client-python/gravitino/api/job/supports_jobs.py
index adb584fe40..b1a2b2b39d 100644
--- a/clients/client-python/gravitino/api/job/supports_jobs.py
+++ b/clients/client-python/gravitino/api/job/supports_jobs.py
@@ -144,12 +144,17 @@ class SupportsJobs(ABC):
pass
@abstractmethod
- def get_job(self, job_id: str) -> JobHandle:
+ def get_job(self, job_id: str, include_output: bool = False) -> JobHandle:
"""
- Retrieves a job by its ID.
+ Retrieves a job by its ID, optionally including its captured
stdout/stderr output (see
+ ``JobHandle.stdout``/``JobHandle.stderr``).
+
+ Output is fetched live from the job executor on every call, not
persisted, so
+ ``include_output`` should only be set to ``True`` when the output is
actually needed.
Args:
job_id: The ID of the job to retrieve.
+ include_output: Whether to also fetch and populate the job's
stdout/stderr output.
Returns:
JobHandle: The handle representing the job with the specified ID.
diff --git a/clients/client-python/gravitino/client/generic_job_handle.py
b/clients/client-python/gravitino/client/generic_job_handle.py
index ce0e2bbca8..4d5531b7e6 100644
--- a/clients/client-python/gravitino/client/generic_job_handle.py
+++ b/clients/client-python/gravitino/client/generic_job_handle.py
@@ -14,6 +14,8 @@
# KIND, either express or implied. See the License for the
# specific language governing permissions and limitations
# under the License.
+from typing import List
+
from gravitino.api.job.job_handle import JobHandle
from gravitino.client.dto_converters import DTOConverters
from gravitino.dto.job.job_dto import JobDTO
@@ -48,3 +50,11 @@ class GenericJobHandle(JobHandle):
if runtime_job_template_dto is None:
return None
return DTOConverters.from_job_template_dto(runtime_job_template_dto)
+
+ def stdout(self) -> List[str]:
+ stdout = self._job_dto.stdout()
+ return [] if stdout is None else stdout
+
+ def stderr(self) -> List[str]:
+ stderr = self._job_dto.stderr()
+ return [] if stderr is None else stderr
diff --git a/clients/client-python/gravitino/client/gravitino_client.py
b/clients/client-python/gravitino/client/gravitino_client.py
index b5644c2dc1..030c4ec15a 100644
--- a/clients/client-python/gravitino/client/gravitino_client.py
+++ b/clients/client-python/gravitino/client/gravitino_client.py
@@ -222,11 +222,12 @@ class GravitinoClient(GravitinoClientBase, SupportsJobs,
TagOperations):
"""
return self.get_metalake().list_jobs(job_template_name)
- def get_job(self, job_id: str) -> JobHandle:
- """Retrieves a job by its ID.
+ def get_job(self, job_id: str, include_output: bool = False) -> JobHandle:
+ """Retrieves a job by its ID, optionally including its captured
stdout/stderr output.
Args:
job_id: The ID of the job to retrieve.
+ include_output: Whether to also fetch and populate the job's
stdout/stderr output.
Returns:
The JobHandle object corresponding to the specified job ID.
@@ -234,7 +235,7 @@ class GravitinoClient(GravitinoClientBase, SupportsJobs,
TagOperations):
Raises:
NoSuchJobException: If no job with the specified ID exists.
"""
- return self.get_metalake().get_job(job_id)
+ return self.get_metalake().get_job(job_id, include_output)
def run_job(self, job_template_name: str, job_conf: Dict[str, str]) ->
JobHandle:
"""Runs a job using the specified job template and configuration.
diff --git a/clients/client-python/gravitino/client/gravitino_metalake.py
b/clients/client-python/gravitino/client/gravitino_metalake.py
index d4ec37e55d..42eac86e61 100644
--- a/clients/client-python/gravitino/client/gravitino_metalake.py
+++ b/clients/client-python/gravitino/client/gravitino_metalake.py
@@ -576,11 +576,15 @@ class GravitinoMetalake(
return [GenericJobHandle(dto) for dto in resp.jobs()]
- def get_job(self, job_id: str) -> JobHandle:
- """Retrieves a job by its ID.
+ def get_job(self, job_id: str, include_output: bool = False) -> JobHandle:
+ """Retrieves a job by its ID, optionally including its captured
stdout/stderr output.
+
+ Output is fetched live from the job executor on every call, not
persisted, so
+ ``include_output`` should only be set to ``True`` when the output is
actually needed.
Args:
job_id: The ID of the job to retrieve.
+ include_output: Whether to also fetch and populate the job's
stdout/stderr output.
Returns:
The JobHandle representing the job if found, otherwise raises an
exception.
@@ -595,7 +599,10 @@ class GravitinoMetalake(
f"{self.API_METALAKES_JOB_RUNS_PATH.format(encode_string(self.name()))}"
f"/{encode_string(job_id)}"
)
- response = self.rest_client.get(url, error_handler=JOB_ERROR_HANDLER)
+ params = {"includeOutput": "true"} if include_output else {}
+ response = self.rest_client.get(
+ url, params=params, error_handler=JOB_ERROR_HANDLER
+ )
resp = JobResponse.from_json(response.body, infer_missing=True)
resp.validate()
diff --git a/clients/client-python/gravitino/dto/job/job_dto.py
b/clients/client-python/gravitino/dto/job/job_dto.py
index dd855e923e..e58a155000 100644
--- a/clients/client-python/gravitino/dto/job/job_dto.py
+++ b/clients/client-python/gravitino/dto/job/job_dto.py
@@ -17,7 +17,7 @@
from dataclasses import dataclass, field
from datetime import datetime
-from typing import Dict, Optional
+from typing import Dict, List, Optional
from dataclasses_json import config, DataClassJsonMixin
@@ -88,6 +88,12 @@ class JobDTO(DataClassJsonMixin): # pylint:
disable=too-many-instance-attribute
decoder=_deserialize_runtime_job_template,
),
)
+ _stdout: Optional[List[str]] = field(
+ default=None, metadata=config(field_name="stdout")
+ )
+ _stderr: Optional[List[str]] = field(
+ default=None, metadata=config(field_name="stderr")
+ )
def __post_init__(self) -> None:
self._queued_at = _deserialize_datetime(self._queued_at)
@@ -132,6 +138,18 @@ class JobDTO(DataClassJsonMixin): # pylint:
disable=too-many-instance-attribute
"""
return self._runtime_job_template
+ def stdout(self) -> Optional[List[str]]:
+ """Returns the captured standard output of the job, as a list of
lines, or ``None`` if
+ output was not requested.
+ """
+ return self._stdout
+
+ def stderr(self) -> Optional[List[str]]:
+ """Returns the captured standard error output of the job, as a list of
lines, or ``None``
+ if output was not requested.
+ """
+ return self._stderr
+
def validate(self) -> None:
"""Validates the JobDTO, ensuring required fields are present and
non-empty."""
if self._job_id is None or not self._job_id.strip():
diff --git a/clients/client-python/tests/integration/test_supports_jobs.py
b/clients/client-python/tests/integration/test_supports_jobs.py
index f8e3321861..0cff0a682b 100644
--- a/clients/client-python/tests/integration/test_supports_jobs.py
+++ b/clients/client-python/tests/integration/test_supports_jobs.py
@@ -402,6 +402,46 @@ class TestSupportsJobs(IntegrationTestEnv):
retrieved_job = self._metalake.get_job(job_handle.job_id())
self.assertEqual(runtime_job_template,
retrieved_job.runtime_job_template())
+ def test_run_and_get_job_output(self):
+ template = self.builder.with_name("test_run_get_output").build()
+ self._metalake.register_job_template(template)
+
+ job_handle = self._metalake.run_job(
+ template.name, {"arg1": "value1", "arg2": "success", "env_var":
"value2"}
+ )
+
+ # Plain get_job never carries output.
+ self.assertEqual([], job_handle.stdout())
+ self.assertEqual([], job_handle.stderr())
+
+ self._wait_until(
+ lambda: self._metalake.get_job(job_handle.job_id()).job_status()
+ == JobHandle.Status.SUCCEEDED,
+ timeout=180,
+ )
+
+ # get_job still never carries output, even after the job finishes.
+ finished_job = self._metalake.get_job(job_handle.job_id())
+ self.assertEqual([], finished_job.stdout())
+ self.assertEqual([], finished_job.stderr())
+
+ # get_job(job_id, include_output=True) fetches the captured
stdout/stderr.
+ job_with_output = self._metalake.get_job(
+ job_handle.job_id(), include_output=True
+ )
+ stdout = job_with_output.stdout()
+ self.assertIn("starting test test job", stdout)
+ self.assertIn("in common script", stdout)
+ self.assertIn("value1", stdout)
+ self.assertIn("success", stdout)
+ self.assertIn("value2", stdout)
+ # The test script never writes to stderr.
+ self.assertEqual([], job_with_output.stderr())
+
+ # Test get output for a non-existent job.
+ with self.assertRaises(NoSuchJobException):
+ self._metalake.get_job("non_existent_job_id", include_output=True)
+
def test_run_and_cancel_job(self):
template = self.builder.with_name("test_run_cancel").build()
self._metalake.register_job_template(template)
diff --git
a/clients/client-python/tests/unittests/dto/job/test_job_dto_serde.py
b/clients/client-python/tests/unittests/dto/job/test_job_dto_serde.py
index 4e4eafb1db..f0d6cadaaa 100644
--- a/clients/client-python/tests/unittests/dto/job/test_job_dto_serde.py
+++ b/clients/client-python/tests/unittests/dto/job/test_job_dto_serde.py
@@ -67,6 +67,32 @@ class TestJobDTOSerDe(unittest.TestCase):
self.assertEqual(queued_at, deser_job_dto.queued_at())
self.assertIsNone(deser_job_dto.started_at())
self.assertIsNone(deser_job_dto.finished_at())
+ self.assertIsNone(deser_job_dto.stdout())
+ self.assertIsNone(deser_job_dto.stderr())
+
+ def test_ser_de_with_output(self):
+ stdout = ["line1", "line2"]
+ stderr = ["err1"]
+ job_dto = JobDTO(
+ _job_id="job-111",
+ _job_template_name="test_template",
+ _status=JobHandle.Status.SUCCEEDED,
+ _audit=AuditDTO(_creator="test",
_create_time=datetime.now(timezone.utc)),
+ _queued_at=datetime.now(timezone.utc),
+ _started_at=datetime.now(timezone.utc),
+ _finished_at=datetime.now(timezone.utc),
+ _stdout=stdout,
+ _stderr=stderr,
+ )
+
+ json_str = job_dto.to_json()
+ self.assertIn("stdout", json_str)
+ self.assertIn("stderr", json_str)
+
+ deser_job_dto = JobDTO.from_json(json_str)
+ self.assertEqual(job_dto, deser_job_dto)
+ self.assertEqual(stdout, deser_job_dto.stdout())
+ self.assertEqual(stderr, deser_job_dto.stderr())
def test_deserialize_from_string(self):
json_str = (
diff --git a/clients/client-python/tests/unittests/test_supports_jobs.py
b/clients/client-python/tests/unittests/test_supports_jobs.py
index 1167491a28..e9affbcb35 100644
--- a/clients/client-python/tests/unittests/test_supports_jobs.py
+++ b/clients/client-python/tests/unittests/test_supports_jobs.py
@@ -263,6 +263,38 @@ class TestSupportsJobs(unittest.TestCase):
job_handle.runtime_job_template(),
)
+ def test_get_job_with_output(self, *mock_methods):
+ gravitino_client = GravitinoClient(
+ uri="http://localhost:8090",
+ metalake_name=self._metalake_name,
+ )
+
+ job_template_name = "test_shell_job"
+ job_dto = self._new_job_dto(
+ job_template_name,
+ finished_at=datetime.now(timezone.utc),
+ started_at=datetime.now(timezone.utc),
+ stdout=["line1", "line2"],
+ stderr=["err1"],
+ )
+ resp = JobResponse(_job=job_dto, _code=0)
+ mock_resp = self._mock_http_response(resp.to_json())
+
+ with patch(
+ "gravitino.utils.http_client.HTTPClient.get",
return_value=mock_resp
+ ):
+ job_handle = gravitino_client.get_job(job_dto.job_id(),
include_output=True)
+ self._compare_job_handle(job_handle, job_dto)
+ self.assertEqual(["line1", "line2"], job_handle.stdout())
+ self.assertEqual(["err1"], job_handle.stderr())
+
+ # test with invalid input
+ with self.assertRaises(ValueError):
+ gravitino_client.get_job("", include_output=True)
+
+ with self.assertRaises(ValueError):
+ gravitino_client.get_job(None, include_output=True)
+
def test_cancel_job(self, *mock_methods):
gravitino_client = GravitinoClient(
uri="http://localhost:8090",
@@ -357,6 +389,8 @@ class TestSupportsJobs(unittest.TestCase):
finished_at: Optional[datetime] = None,
started_at: Optional[datetime] = None,
runtime_job_template: Optional[JobTemplateDTO] = None,
+ stdout: Optional[list] = None,
+ stderr: Optional[list] = None,
) -> JobDTO:
return JobDTO(
_job_id="job-123",
@@ -367,6 +401,8 @@ class TestSupportsJobs(unittest.TestCase):
_started_at=started_at,
_finished_at=finished_at,
_runtime_job_template=runtime_job_template,
+ _stdout=stdout,
+ _stderr=stderr,
)
def _compare_job_handle(self, job_handle: JobHandle, job_dto: JobDTO):
diff --git a/common/src/main/java/org/apache/gravitino/dto/job/JobDTO.java
b/common/src/main/java/org/apache/gravitino/dto/job/JobDTO.java
index 54b4abaf11..05745e39ad 100644
--- a/common/src/main/java/org/apache/gravitino/dto/job/JobDTO.java
+++ b/common/src/main/java/org/apache/gravitino/dto/job/JobDTO.java
@@ -30,6 +30,8 @@ import
com.fasterxml.jackson.databind.annotation.JsonSerialize;
import com.google.common.base.Preconditions;
import java.io.IOException;
import java.time.Instant;
+import java.util.List;
+import javax.annotation.Nullable;
import lombok.EqualsAndHashCode;
import lombok.Getter;
import lombok.ToString;
@@ -71,9 +73,17 @@ public class JobDTO {
@JsonProperty("runtimeJobTemplate")
private final JobTemplateDTO runtimeJobTemplate;
+ @JsonProperty("stdout")
+ @Nullable
+ private final List<String> stdout;
+
+ @JsonProperty("stderr")
+ @Nullable
+ private final List<String> stderr;
+
/** Default constructor for Jackson deserialization. */
private JobDTO() {
- this(null, null, null, null, null, null, null, null);
+ this(null, null, null, null, null, null, null, null, null, null);
}
/**
@@ -91,6 +101,10 @@ public class JobDTO {
* @param runtimeJobTemplate The resolved job template that was actually
submitted for execution,
* with placeholders replaced and referenced files downloaded, or null
for jobs run before
* this field was introduced.
+ * @param stdout The captured standard output of the job, as a list of
lines, or null if output
+ * was not requested.
+ * @param stderr The captured standard error output of the job, as a list of
lines, or null if
+ * output was not requested.
*/
public JobDTO(
String jobId,
@@ -100,7 +114,9 @@ public class JobDTO {
Instant queuedAt,
Instant startedAt,
Instant finishedAt,
- JobTemplateDTO runtimeJobTemplate) {
+ JobTemplateDTO runtimeJobTemplate,
+ @Nullable List<String> stdout,
+ @Nullable List<String> stderr) {
this.jobId = jobId;
this.jobTemplateName = jobTemplateName;
this.status = status;
@@ -109,6 +125,8 @@ public class JobDTO {
this.startedAt = startedAt;
this.finishedAt = finishedAt;
this.runtimeJobTemplate = runtimeJobTemplate;
+ this.stdout = stdout;
+ this.stderr = stderr;
}
/**
diff --git a/common/src/test/java/org/apache/gravitino/dto/job/TestJobDTO.java
b/common/src/test/java/org/apache/gravitino/dto/job/TestJobDTO.java
index 8065893bb0..3883913adc 100644
--- a/common/src/test/java/org/apache/gravitino/dto/job/TestJobDTO.java
+++ b/common/src/test/java/org/apache/gravitino/dto/job/TestJobDTO.java
@@ -19,8 +19,10 @@
package org.apache.gravitino.dto.job;
import com.fasterxml.jackson.core.JsonProcessingException;
+import com.google.common.collect.ImmutableList;
import com.google.common.collect.Lists;
import java.time.Instant;
+import java.util.List;
import org.apache.gravitino.dto.AuditDTO;
import org.apache.gravitino.job.JobHandle;
import org.apache.gravitino.job.JobTemplate;
@@ -44,6 +46,8 @@ public class TestJobDTO {
queuedAt,
startedAt,
finishedAt,
+ null,
+ null,
null);
Assertions.assertDoesNotThrow(jobDTO::validate);
@@ -58,6 +62,37 @@ public class TestJobDTO {
Assertions.assertEquals(queuedAt, deserJobDTO.queuedAt());
Assertions.assertEquals(startedAt, deserJobDTO.startedAt());
Assertions.assertEquals(finishedAt, deserJobDTO.finishedAt());
+ Assertions.assertNull(deserJobDTO.stdout());
+ Assertions.assertNull(deserJobDTO.stderr());
+ }
+
+ @Test
+ public void testSerDeWithOutput() throws JsonProcessingException {
+ List<String> stdout = ImmutableList.of("line1", "line2");
+ List<String> stderr = ImmutableList.of("err1");
+ JobDTO jobDTO =
+ new JobDTO(
+ "job-111",
+ "testTemplate",
+ JobHandle.Status.SUCCEEDED,
+
AuditDTO.builder().withCreator("test").withCreateTime(Instant.now()).build(),
+ Instant.now(),
+ Instant.now(),
+ Instant.now(),
+ null,
+ stdout,
+ stderr);
+
+ Assertions.assertDoesNotThrow(jobDTO::validate);
+
+ String serJson = JsonUtils.objectMapper().writeValueAsString(jobDTO);
+ Assertions.assertTrue(serJson.contains("\"stdout\""));
+ Assertions.assertTrue(serJson.contains("\"stderr\""));
+
+ JobDTO deserJobDTO = JsonUtils.objectMapper().readValue(serJson,
JobDTO.class);
+ Assertions.assertEquals(jobDTO, deserJobDTO);
+ Assertions.assertEquals(stdout, deserJobDTO.stdout());
+ Assertions.assertEquals(stderr, deserJobDTO.stderr());
}
@Test
@@ -72,6 +107,8 @@ public class TestJobDTO {
queuedAt,
null,
null,
+ null,
+ null,
null);
Assertions.assertDoesNotThrow(jobDTO::validate);
@@ -99,6 +136,8 @@ public class TestJobDTO {
queuedAt,
startedAt,
finishedAt,
+ null,
+ null,
null);
String serJson = JsonUtils.objectMapper().writeValueAsString(jobDTO);
@@ -174,7 +213,9 @@ public class TestJobDTO {
Instant.now(),
Instant.now(),
Instant.now(),
- runtimeJobTemplate);
+ runtimeJobTemplate,
+ null,
+ null);
Assertions.assertDoesNotThrow(jobDTO::validate);
@@ -198,6 +239,8 @@ public class TestJobDTO {
Instant.now(),
null,
null,
+ null,
+ null,
null);
Assertions.assertDoesNotThrow(jobDTO::validate);
diff --git a/core/src/main/java/org/apache/gravitino/Configs.java
b/core/src/main/java/org/apache/gravitino/Configs.java
index de6e2315c6..18e9805096 100644
--- a/core/src/main/java/org/apache/gravitino/Configs.java
+++ b/core/src/main/java/org/apache/gravitino/Configs.java
@@ -585,6 +585,31 @@ public class Configs {
.checkValue(value -> value > 0,
ConfigConstants.POSITIVE_NUMBER_ERROR_MSG)
.createWithDefault(5 * 60 * 1000L); // Default is 5 minutes
+ public static final ConfigEntry<Integer> JOB_OUTPUT_MAX_LINES =
+ new ConfigBuilder("gravitino.job.outputMaxLines")
+ .doc(
+ "The maximum number of lines returned by
JobExecutor#getJobStdout and "
+ + "JobExecutor#getJobStderr. This is resolved by JobManager
and passed as an "
+ + "argument to those two APIs, so all executors honor the
same cap.")
+ .version(ConfigConstants.VERSION_2_0_0)
+ .intConf()
+ .checkValue(value -> value > 0,
ConfigConstants.POSITIVE_NUMBER_ERROR_MSG)
+ .createWithDefault(1000);
+
+ public static final ConfigEntry<Integer> JOB_OUTPUT_MAX_BYTES =
+ new ConfigBuilder("gravitino.job.outputMaxBytes")
+ .doc(
+ "The maximum number of bytes read from the tail of a job's
captured stdout/stderr "
+ + "when retrieving its output. Bounds both the read cost and
the response size "
+ + "regardless of how the content is shaped (e.g. a single
very long line). "
+ + "This is resolved by JobManager and passed as an argument
to "
+ + "JobExecutor#getJobStdout and JobExecutor#getJobStderr, so
all executors "
+ + "honor the same cap.")
+ .version(ConfigConstants.VERSION_2_0_0)
+ .intConf()
+ .checkValue(value -> value > 0,
ConfigConstants.POSITIVE_NUMBER_ERROR_MSG)
+ .createWithDefault(256 * 1024); // 256KB
+
public static final ConfigEntry<Boolean> BLOCK_UNSAFE_REMOTE_URI =
new ConfigBuilder(FileFetcher.BLOCK_UNSAFE_REMOTE_URI_CONFIG)
.doc(
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 e14fb6dce1..6df99d9d9b 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
@@ -20,6 +20,8 @@
package org.apache.gravitino.connector.job;
import java.io.Closeable;
+import java.util.Collections;
+import java.util.List;
import java.util.Map;
import org.apache.gravitino.annotation.DeveloperApi;
import org.apache.gravitino.exceptions.NoSuchJobException;
@@ -112,4 +114,48 @@ public interface JobExecutor extends Closeable {
default boolean isJobStateNodeLocal() {
return false;
}
+
+ /**
+ * Get the captured standard output of the job, as a list of lines.
+ *
+ * <p>The default implementation returns an empty list, so implementors that
don't support output
+ * 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.
+ *
+ * @param jobId The unique identifier of the job.
+ * @param maxLines The maximum number of (most recent) lines to return,
resolved by the caller
+ * from the {@code gravitino.job.outputMaxLines} configuration.
+ * @param maxBytes The maximum number of (most recent) bytes to read from
the underlying output,
+ * resolved by the caller from the {@code gravitino.job.outputMaxBytes}
configuration. Bounds
+ * both the read cost and the response size regardless of how the
content is shaped.
+ * @return the stdout lines of the job, or an empty list if not available.
+ */
+ default List<String> getJobStdout(String jobId, int maxLines, int maxBytes) {
+ return Collections.emptyList();
+ }
+
+ /**
+ * Get the captured standard error output of the job, as a list of lines.
+ *
+ * <p>The default implementation returns an empty list, so implementors that
don't support output
+ * 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.
+ *
+ * @param jobId The unique identifier of the job.
+ * @param maxLines The maximum number of (most recent) lines to return,
resolved by the caller
+ * from the {@code gravitino.job.outputMaxLines} configuration.
+ * @param maxBytes The maximum number of (most recent) bytes to read from
the underlying output,
+ * resolved by the caller from the {@code gravitino.job.outputMaxBytes}
configuration. Bounds
+ * both the read cost and the response size regardless of how the
content is shaped.
+ * @return the stderr lines of the job, or an empty list if not available.
+ */
+ default List<String> getJobStderr(String jobId, int maxLines, int maxBytes) {
+ return Collections.emptyList();
+ }
}
diff --git
a/core/src/main/java/org/apache/gravitino/hook/JobHookDispatcher.java
b/core/src/main/java/org/apache/gravitino/hook/JobHookDispatcher.java
index 32badedc29..d6ca24eff2 100644
--- a/core/src/main/java/org/apache/gravitino/hook/JobHookDispatcher.java
+++ b/core/src/main/java/org/apache/gravitino/hook/JobHookDispatcher.java
@@ -91,8 +91,10 @@ public class JobHookDispatcher implements
JobOperationDispatcher {
}
@Override
- public JobEntity getJob(String metalake, String jobId) throws
NoSuchJobException {
- return jobOperationDispatcher.getJob(metalake, jobId);
+ public JobEntity getJob(
+ String metalake, String jobId, boolean includeOutput, Integer maxLines,
Integer maxBytes)
+ throws NoSuchJobException {
+ return jobOperationDispatcher.getJob(metalake, jobId, includeOutput,
maxLines, maxBytes);
}
@Override
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 a86e7d4d91..6fef520af7 100644
--- a/core/src/main/java/org/apache/gravitino/job/JobManager.java
+++ b/core/src/main/java/org/apache/gravitino/job/JobManager.java
@@ -109,6 +109,10 @@ public class JobManager implements JobOperationDispatcher {
private final long jobStagingDirKeepTimeInMs;
+ private final int jobOutputMaxLines;
+
+ private final int jobOutputMaxBytes;
+
@VisibleForTesting final ScheduledExecutorService cleanUpExecutor;
@VisibleForTesting final ScheduledExecutorService statusPullExecutor;
@@ -155,6 +159,9 @@ public class JobManager implements JobOperationDispatcher {
JOB_STAGING_DIR_CLEANUP_MIN_TIME_IN_MS);
}
+ this.jobOutputMaxLines = config.get(Configs.JOB_OUTPUT_MAX_LINES);
+ this.jobOutputMaxBytes = config.get(Configs.JOB_OUTPUT_MAX_BYTES);
+
this.cleanUpExecutor =
Executors.newSingleThreadScheduledExecutor(
runnable -> {
@@ -414,23 +421,69 @@ public class JobManager implements JobOperationDispatcher
{
}
@Override
- public JobEntity getJob(String metalake, String jobId) throws
NoSuchJobException {
+ public JobEntity getJob(
+ String metalake, String jobId, boolean includeOutput, Integer maxLines,
Integer maxBytes)
+ throws NoSuchJobException {
checkMetalake(NameIdentifierUtil.ofMetalake(metalake), entityStore);
NameIdentifier jobIdent = NameIdentifierUtil.ofJob(metalake, jobId);
- return TreeLockUtils.doWithTreeLock(
- jobIdent,
- LockType.READ,
- () -> {
- try {
- return entityStore.get(jobIdent, Entity.EntityType.JOB,
JobEntity.class);
- } catch (NoSuchEntityException e) {
- throw new NoSuchJobException(
- "Job with ID %s under metalake %s does not exist", jobId,
metalake);
- } catch (IOException ioe) {
- throw new RuntimeException(ioe);
- }
- });
+ JobEntity entity =
+ TreeLockUtils.doWithTreeLock(
+ jobIdent,
+ LockType.READ,
+ () -> {
+ try {
+ return entityStore.get(jobIdent, Entity.EntityType.JOB,
JobEntity.class);
+ } catch (NoSuchEntityException e) {
+ throw new NoSuchJobException(
+ "Job with ID %s under metalake %s does not exist", jobId,
metalake);
+ } catch (IOException ioe) {
+ throw new RuntimeException(ioe);
+ }
+ });
+
+ if (!includeOutput) {
+ return entity;
+ }
+
+ // maxLines/maxBytes are only meaningful (and only validated) when output
is actually
+ // requested - per the documented contract, they're ignored entirely when
includeOutput is
+ // false, so an invalid value must not fail a plain getJob call that never
uses them.
+ Preconditions.checkArgument(
+ maxLines == null || maxLines > 0, "maxLines must be positive if
specified");
+ Preconditions.checkArgument(
+ maxBytes == null || maxBytes > 0, "maxBytes must be positive if
specified");
+
+ // A caller-specified maxLines/maxBytes can only narrow the globally
configured cap, never
+ // widen it - the global configuration remains a hard upper bound on read
cost/response size.
+ int effectiveMaxLines = clampToGlobalMax(maxLines, jobOutputMaxLines);
+ int effectiveMaxBytes = clampToGlobalMax(maxBytes, jobOutputMaxBytes);
+
+ // 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
+ // 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);
+ List<String> stderr =
+ jobExecutor.getJobStderr(entity.jobExecutionId(), effectiveMaxLines,
effectiveMaxBytes);
+ if (stdout.isEmpty() && stderr.isEmpty() && LOG.isDebugEnabled()) {
+ LOG.debug(
+ "No output available for job {} under metalake {} - either it
produced none, or the "
+ + "job executor no longer has a record of it",
+ jobId,
+ metalake);
+ }
+ return entity.withOutput(stdout, stderr);
+ }
+
+ // A caller-specified value can only narrow the globally configured cap,
never widen it -
+ // null means "use the global default", and any non-null value is clamped to
at most that
+ // default so the global configuration always remains a hard upper bound.
+ private static int clampToGlobalMax(Integer requested, int globalMax) {
+ return requested == null ? globalMax : Math.min(requested, globalMax);
}
@Override
@@ -534,7 +587,7 @@ public class JobManager implements JobOperationDispatcher {
checkMetalake(NameIdentifierUtil.ofMetalake(metalake), entityStore);
// Retrieve the job entity, will throw NoSuchJobException if the job does
not exist.
- JobEntity jobEntity = getJob(metalake, jobId);
+ JobEntity jobEntity = getJob(metalake, jobId, false);
if (jobEntity.status() == JobHandle.Status.CANCELLING
|| jobEntity.status() == JobHandle.Status.CANCELLED
diff --git
a/core/src/main/java/org/apache/gravitino/job/JobOperationDispatcher.java
b/core/src/main/java/org/apache/gravitino/job/JobOperationDispatcher.java
index 795ae6aeae..f0e349a8cc 100644
--- a/core/src/main/java/org/apache/gravitino/job/JobOperationDispatcher.java
+++ b/core/src/main/java/org/apache/gravitino/job/JobOperationDispatcher.java
@@ -100,14 +100,51 @@ public interface JobOperationDispatcher extends Closeable
{
throws NoSuchJobTemplateException;
/**
- * Retrieves a job by its ID in the specified metalake.
+ * Retrieves a job by its ID in the specified metalake, optionally including
its captured
+ * stdout/stderr output (see {@link JobEntity#stdout()}/{@link
JobEntity#stderr()}), using the
+ * globally configured {@code gravitino.job.outputMaxLines}/{@code
outputMaxBytes} caps.
+ *
+ * <p>Output is fetched live from the {@code JobExecutor} on every call, not
persisted, so {@code
+ * includeOutput} should only be set to {@code true} when the caller
actually needs the output -
+ * it is never included in {@link #listJobs(String, Optional)}, and callers
that don't need output
+ * should pass {@code false} here.
+ *
+ * @param metalake the name of the metalake
+ * @param jobId the ID of the job to retrieve
+ * @param includeOutput whether to also fetch and attach the job's
stdout/stderr output
+ * @return the job entity associated with the specified ID
+ * @throws NoSuchJobException if no job with the specified ID exists
+ */
+ default JobEntity getJob(String metalake, String jobId, boolean
includeOutput)
+ throws NoSuchJobException {
+ return getJob(metalake, jobId, includeOutput, null, null);
+ }
+
+ /**
+ * Retrieves a job by its ID in the specified metalake, optionally including
its captured
+ * stdout/stderr output, with caller-specified caps on how much of it to
return.
+ *
+ * <p>{@code maxLines}/{@code maxBytes} let a caller ask for less output
than the globally
+ * configured {@code gravitino.job.outputMaxLines}/{@code outputMaxBytes}
caps, but never more - a
+ * value larger than the global cap is clamped down to it, so the global
configuration always
+ * remains a hard upper bound. Passing {@code null} for either uses the
global default. Both are
+ * ignored when {@code includeOutput} is {@code false}.
*
* @param metalake the name of the metalake
* @param jobId the ID of the job to retrieve
+ * @param includeOutput whether to also fetch and attach the job's
stdout/stderr output
+ * @param maxLines the maximum number of (most recent) output lines to
return, or {@code null} to
+ * use the global default
+ * @param maxBytes the maximum number of (most recent) output bytes to read,
or {@code null} to
+ * use the global default
* @return the job entity associated with the specified ID
* @throws NoSuchJobException if no job with the specified ID exists
+ * @throws IllegalArgumentException if {@code maxLines} or {@code maxBytes}
is specified and not
+ * positive
*/
- JobEntity getJob(String metalake, String jobId) throws NoSuchJobException;
+ JobEntity getJob(
+ String metalake, String jobId, boolean includeOutput, Integer maxLines,
Integer maxBytes)
+ throws NoSuchJobException;
/**
* Runs a job based on the specified job template and configuration in the
specified metalake.
diff --git
a/core/src/main/java/org/apache/gravitino/job/JobTemplateValidationDispatcher.java
b/core/src/main/java/org/apache/gravitino/job/JobTemplateValidationDispatcher.java
index 6a25735a5b..e54d478c60 100644
---
a/core/src/main/java/org/apache/gravitino/job/JobTemplateValidationDispatcher.java
+++
b/core/src/main/java/org/apache/gravitino/job/JobTemplateValidationDispatcher.java
@@ -114,8 +114,10 @@ public class JobTemplateValidationDispatcher implements
JobOperationDispatcher {
}
@Override
- public JobEntity getJob(String metalake, String jobId) throws
NoSuchJobException {
- return dispatcher.getJob(metalake, jobId);
+ public JobEntity getJob(
+ String metalake, String jobId, boolean includeOutput, Integer maxLines,
Integer maxBytes)
+ throws NoSuchJobException {
+ return dispatcher.getJob(metalake, jobId, includeOutput, maxLines,
maxBytes);
}
@Override
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 e441e7b207..d55d7b47f6 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
@@ -28,8 +28,15 @@ import static
org.apache.gravitino.job.local.LocalJobExecutorConfigs.WAITING_QUE
import com.google.common.annotations.VisibleForTesting;
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.charset.StandardCharsets;
+import java.util.List;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.BlockingQueue;
@@ -40,6 +47,7 @@ import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ThreadLocalRandom;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
+import javax.annotation.Nullable;
import org.apache.commons.lang3.tuple.Pair;
import org.apache.gravitino.connector.job.JobExecutor;
import org.apache.gravitino.exceptions.NoSuchJobException;
@@ -84,6 +92,10 @@ public class LocalJobExecutor implements JobExecutor {
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;
+
@Override
public void initialize(Map<String, String> configs) {
this.configs = configs;
@@ -165,6 +177,7 @@ public class LocalJobExecutor implements JobExecutor {
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.
@@ -195,6 +208,9 @@ 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));
}
return newJobId;
@@ -264,6 +280,16 @@ public class LocalJobExecutor implements JobExecutor {
return true;
}
+ @Override
+ public List<String> getJobStdout(String jobId, int maxLines, int maxBytes) {
+ return getJobOutput(jobId, LocalProcessBuilder.STDOUT_FILE_NAME, maxLines,
maxBytes);
+ }
+
+ @Override
+ public List<String> getJobStderr(String jobId, int maxLines, int maxBytes) {
+ return getJobOutput(jobId, LocalProcessBuilder.STDERR_FILE_NAME, maxLines,
maxBytes);
+ }
+
@Override
public void close() throws IOException {
// Mark the executor as finished to stop processing jobs
@@ -291,6 +317,7 @@ public class LocalJobExecutor implements JobExecutor {
// Stop the job status cleanup executor
jobStatusCleanupExecutor.shutdownNow();
jobStatus.clear();
+ jobWorkingDirs.clear();
}
public void runJob(Pair<String, JobTemplate> jobPair) {
@@ -377,9 +404,104 @@ public class LocalJobExecutor implements JobExecutor {
jobStatus
.entrySet()
.removeIf(
- entry ->
- entry.getValue().getRight() != UNEXPIRED_TIME_IN_MS
- && (currentTime - entry.getValue().getRight()) >=
jobStatusKeepTimeInMs);
+ entry -> {
+ boolean expired =
+ entry.getValue().getRight() != UNEXPIRED_TIME_IN_MS
+ && (currentTime - entry.getValue().getRight()) >=
jobStatusKeepTimeInMs;
+ if (expired) {
+ jobWorkingDirs.remove(entry.getKey());
+ }
+ return expired;
+ });
}
}
+
+ private List<String> getJobOutput(String jobId, String fileName, int
maxLines, int maxBytes) {
+ File workingDir = getWorkingDir(jobId);
+ if (workingDir == null) {
+ return ImmutableList.of();
+ }
+ return readLastLines(new File(workingDir, fileName), maxLines, maxBytes);
+ }
+
+ @Nullable
+ private File getWorkingDir(String jobId) {
+ return jobWorkingDirs.get(jobId);
+ }
+
+ private List<String> readLastLines(File file, int maxLines, int maxBytes) {
+ if (!file.exists()) {
+ // The job hasn't started (or hasn't produced this stream) yet.
+ 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) {
+ // 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
+ // handles, not a real I/O failure, so it must degrade to empty output
the same way.
+ return ImmutableList.of();
+ } catch (IOException e) {
+ // Any other I/O failure while reading an existing file is unexpected
and must not be
+ // silently reported as "no output" - that would be actively misleading
for the debugging
+ // use case this method exists for.
+ throw new RuntimeException("Failed to read job output file: " + file, e);
+ }
+
+ 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
+ // continuation byte - skip forward to the next character boundary so the
decoded content
+ // never begins with a corrupted replacement character. A '\n' byte never
appears inside a
+ // multi-byte UTF-8 character, so this can't skip past a real line
boundary.
+ while (contentStart < probeWindow.length && (probeWindow[contentStart] &
0xC0) == 0x80) {
+ contentStart++;
+ }
+
+ String content =
+ new String(
+ probeWindow, contentStart, probeWindow.length - contentStart,
StandardCharsets.UTF_8);
+ if (!windowStartsAtLineBoundary) {
+ // The window starts mid-file and the preceding byte isn't a line
terminator, so its first
+ // line may be a partial line whose true beginning fell outside the
window - drop up to and
+ // including the first newline, but only when something actually follows
it. If the only
+ // newline in the window is its very last character, that newline
terminates the single
+ // oversized line the window is entirely made of, rather than starting a
subsequent one -
+ // dropping through it would discard the whole line instead of
truncating it. Keep the
+ // content as-is in that case (and when no newline is found at all),
since a truncated line
+ // is more useful for debugging than silently returning nothing.
+ int firstNewline = content.indexOf('\n');
+ if (firstNewline >= 0 && firstNewline < content.length() - 1) {
+ content = content.substring(firstNewline + 1);
+ }
+ }
+
+ if (content.isEmpty()) {
+ return ImmutableList.of();
+ }
+
+ // Recognize both LF and CRLF line endings, so output captured with
Windows-style line
+ // endings doesn't leave a dangling '\r' on every returned line.
+ List<String> lines = Splitter.onPattern("\r\n|\n").splitToList(content);
+ if (content.endsWith("\n")) {
+ // A trailing separator produces a spurious empty trailing element -
drop it so a
+ // completed line isn't followed by a phantom blank one.
+ lines = lines.subList(0, lines.size() - 1);
+ }
+
+ int fromIndex = Math.max(0, lines.size() - maxLines);
+ return ImmutableList.copyOf(lines.subList(fromIndex, lines.size()));
+ }
}
diff --git
a/core/src/main/java/org/apache/gravitino/job/local/LocalProcessBuilder.java
b/core/src/main/java/org/apache/gravitino/job/local/LocalProcessBuilder.java
index 0338dc437f..f3c7f6228e 100644
--- a/core/src/main/java/org/apache/gravitino/job/local/LocalProcessBuilder.java
+++ b/core/src/main/java/org/apache/gravitino/job/local/LocalProcessBuilder.java
@@ -26,15 +26,30 @@ import org.apache.gravitino.job.SparkJobTemplate;
public abstract class LocalProcessBuilder {
+ /** The name of the file that captures the job process's standard output. */
+ public static final String STDOUT_FILE_NAME = "output.log";
+
+ /** The name of the file that captures the job process's standard error. */
+ public static final String STDERR_FILE_NAME = "error.log";
+
protected final JobTemplate jobTemplate;
protected final File workingDirectory;
protected LocalProcessBuilder(JobTemplate jobTemplate, Map<String, String>
configs) {
this.jobTemplate = jobTemplate;
- // Executable should be in the working directory, so we can figure out the
working directory
- // from the executable path.
- this.workingDirectory = new
File(jobTemplate.executable()).getAbsoluteFile().getParentFile();
+ this.workingDirectory = resolveWorkingDirectory(jobTemplate);
+ }
+
+ /**
+ * Resolves the working directory for a job template. The executable is
expected to be in the
+ * working directory, so the working directory can be derived from the
executable's path.
+ *
+ * @param jobTemplate the job template to resolve the working directory for
+ * @return the working directory for the job template
+ */
+ public static File resolveWorkingDirectory(JobTemplate jobTemplate) {
+ return new
File(jobTemplate.executable()).getAbsoluteFile().getParentFile();
}
public abstract Process start();
diff --git
a/core/src/main/java/org/apache/gravitino/job/local/ShellProcessBuilder.java
b/core/src/main/java/org/apache/gravitino/job/local/ShellProcessBuilder.java
index 940bba246c..b8ae065fbf 100644
--- a/core/src/main/java/org/apache/gravitino/job/local/ShellProcessBuilder.java
+++ b/core/src/main/java/org/apache/gravitino/job/local/ShellProcessBuilder.java
@@ -51,8 +51,8 @@ public class ShellProcessBuilder extends LocalProcessBuilder {
builder.directory(workingDirectory);
builder.environment().putAll(shellJobTemplate.environments());
- File outputFile = new File(workingDirectory, "output.log");
- File errorFile = new File(workingDirectory, "error.log");
+ File outputFile = new File(workingDirectory, STDOUT_FILE_NAME);
+ File errorFile = new File(workingDirectory, STDERR_FILE_NAME);
builder.redirectOutput(outputFile);
builder.redirectError(errorFile);
diff --git
a/core/src/main/java/org/apache/gravitino/job/local/SparkProcessBuilder.java
b/core/src/main/java/org/apache/gravitino/job/local/SparkProcessBuilder.java
index 06d9cef8e5..7d717b1dc8 100644
--- a/core/src/main/java/org/apache/gravitino/job/local/SparkProcessBuilder.java
+++ b/core/src/main/java/org/apache/gravitino/job/local/SparkProcessBuilder.java
@@ -140,8 +140,8 @@ public class SparkProcessBuilder extends
LocalProcessBuilder {
builder.directory(workingDirectory);
builder.environment().putAll(sparkJobTemplate.environments());
- File outputFile = new File(workingDirectory, "output.log");
- File errorFile = new File(workingDirectory, "error.log");
+ File outputFile = new File(workingDirectory, STDOUT_FILE_NAME);
+ File errorFile = new File(workingDirectory, STDERR_FILE_NAME);
builder.redirectOutput(outputFile);
builder.redirectError(errorFile);
diff --git
a/core/src/main/java/org/apache/gravitino/listener/JobEventDispatcher.java
b/core/src/main/java/org/apache/gravitino/listener/JobEventDispatcher.java
index 28e3f626fd..f8f4896fcd 100644
--- a/core/src/main/java/org/apache/gravitino/listener/JobEventDispatcher.java
+++ b/core/src/main/java/org/apache/gravitino/listener/JobEventDispatcher.java
@@ -201,12 +201,15 @@ public class JobEventDispatcher implements
JobOperationDispatcher {
}
@Override
- public JobEntity getJob(String metalake, String jobId) throws
NoSuchJobException {
+ public JobEntity getJob(
+ String metalake, String jobId, boolean includeOutput, Integer maxLines,
Integer maxBytes)
+ throws NoSuchJobException {
eventBus.dispatchEvent(
new GetJobPreEvent(PrincipalUtils.getCurrentUserName(), metalake,
jobId));
try {
- JobEntity job = jobOperationDispatcher.getJob(metalake, jobId);
+ JobEntity job =
+ jobOperationDispatcher.getJob(metalake, jobId, includeOutput,
maxLines, maxBytes);
eventBus.dispatchEvent(
new GetJobEvent(
PrincipalUtils.getCurrentUserName(), metalake,
JobInfo.fromJobEntity(job)));
diff --git a/core/src/main/java/org/apache/gravitino/meta/JobEntity.java
b/core/src/main/java/org/apache/gravitino/meta/JobEntity.java
index ff1e1f81a6..6c73983e01 100644
--- a/core/src/main/java/org/apache/gravitino/meta/JobEntity.java
+++ b/core/src/main/java/org/apache/gravitino/meta/JobEntity.java
@@ -22,6 +22,7 @@ package org.apache.gravitino.meta;
import com.google.common.collect.Maps;
import java.time.Instant;
import java.util.Collections;
+import java.util.List;
import java.util.Map;
import java.util.Objects;
import javax.annotation.Nullable;
@@ -80,6 +81,13 @@ public class JobEntity implements Entity, Auditable,
HasIdentifier {
private Long finishedAt;
private String runtimeJobTemplate;
+ // The stdout/stderr of the job, fetched live from the JobExecutor on demand
(e.g. via
+ // JobOperationDispatcher#getJob(String, String, boolean)). Deliberately not
included in
+ // fields()/equals()/hashCode() - this is never persisted, it only exists on
the in-memory copy
+ // returned when output was explicitly requested.
+ private List<String> stdout;
+ private List<String> stderr;
+
private JobEntity() {}
@Override
@@ -169,6 +177,52 @@ public class JobEntity implements Entity, Auditable,
HasIdentifier {
return runtimeJobTemplate;
}
+ /**
+ * Get the captured standard output of the job.
+ *
+ * @return the stdout lines of the job, or {@code null} if output was not
requested for this
+ * entity (see {@link #withOutput(List, List)}).
+ */
+ @Nullable
+ public List<String> stdout() {
+ return stdout;
+ }
+
+ /**
+ * Get the captured standard error output of the job.
+ *
+ * @return the stderr lines of the job, or {@code null} if output was not
requested for this
+ * entity (see {@link #withOutput(List, List)}).
+ */
+ @Nullable
+ public List<String> stderr() {
+ return stderr;
+ }
+
+ /**
+ * Returns a copy of this entity with the given stdout/stderr attached. This
entity is left
+ * unmodified.
+ *
+ * @param stdout the stdout lines to attach
+ * @param stderr the stderr lines to attach
+ * @return a new {@link JobEntity} with the given output attached
+ */
+ public JobEntity withOutput(List<String> stdout, List<String> stderr) {
+ return JobEntity.builder()
+ .withId(id)
+ .withJobExecutionId(jobExecutionId)
+ .withNamespace(namespace)
+ .withStatus(status)
+ .withJobTemplateName(jobTemplateName)
+ .withAuditInfo(auditInfo)
+ .withStartedAt(startedAt)
+ .withFinishedAt(finishedAt)
+ .withRuntimeJobTemplate(runtimeJobTemplate)
+ .withStdout(stdout)
+ .withStderr(stderr)
+ .build();
+ }
+
@Override
public AuditInfo auditInfo() {
return auditInfo;
@@ -270,6 +324,16 @@ public class JobEntity implements Entity, Auditable,
HasIdentifier {
return this;
}
+ public Builder withStdout(List<String> stdout) {
+ jobEntity.stdout = stdout;
+ return this;
+ }
+
+ public Builder withStderr(List<String> stderr) {
+ jobEntity.stderr = stderr;
+ return this;
+ }
+
public JobEntity build() {
jobEntity.validate();
return jobEntity;
diff --git
a/core/src/main/java/org/apache/gravitino/utils/MetadataObjectUtil.java
b/core/src/main/java/org/apache/gravitino/utils/MetadataObjectUtil.java
index 1a2d9f5b5f..445e3d640b 100644
--- a/core/src/main/java/org/apache/gravitino/utils/MetadataObjectUtil.java
+++ b/core/src/main/java/org/apache/gravitino/utils/MetadataObjectUtil.java
@@ -332,7 +332,7 @@ public class MetadataObjectUtil {
case JOB:
NameIdentifierUtil.checkJob(identifier);
try {
- env.internalJobOperationDispatcher().getJob(metalake,
object.fullName());
+ env.internalJobOperationDispatcher().getJob(metalake,
object.fullName(), false);
} catch (NoSuchJobException e) {
throw exceptionToThrowSupplier.get();
}
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 b3605dd576..79318a429a 100644
--- a/core/src/test/java/org/apache/gravitino/job/TestJobManager.java
+++ b/core/src/test/java/org/apache/gravitino/job/TestJobManager.java
@@ -20,6 +20,7 @@ package org.apache.gravitino.job;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyBoolean;
+import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.doNothing;
import static org.mockito.Mockito.doReturn;
@@ -595,7 +596,7 @@ public class TestJobManager {
.thenReturn(job);
// Get an existing job
- JobEntity retrievedJob = jobManager.getJob(metalake, job.name());
+ JobEntity retrievedJob = jobManager.getJob(metalake, job.name(), false);
Assertions.assertEquals(job, retrievedJob);
// Throw exception if job does not exist
@@ -607,7 +608,7 @@ public class TestJobManager {
Exception e =
Assertions.assertThrows(
- NoSuchJobException.class, () -> jobManager.getJob(metalake,
"non_existent"));
+ NoSuchJobException.class, () -> jobManager.getJob(metalake,
"non_existent", false));
Assertions.assertEquals(
"Job with ID non_existent under metalake test_metalake does not
exist", e.getMessage());
@@ -615,7 +616,135 @@ public class TestJobManager {
doThrow(new IOException("Entity store error"))
.when(entityStore)
.get(NameIdentifierUtil.ofJob(metalake, "job"), Entity.EntityType.JOB,
JobEntity.class);
- Assertions.assertThrows(RuntimeException.class, () ->
jobManager.getJob(metalake, "job"));
+ Assertions.assertThrows(
+ RuntimeException.class, () -> jobManager.getJob(metalake, "job",
false));
+ }
+
+ @Test
+ public void testGetJobWithOutput() throws IOException {
+ mockedMetalake
+ .when(() -> MetalakeManager.checkMetalake(metalakeIdent, entityStore))
+ .thenAnswer(a -> null);
+
+ JobEntity job = newJobEntity("shell_job", JobHandle.Status.SUCCEEDED);
+ when(entityStore.get(
+ NameIdentifierUtil.ofJob(metalake, job.name()),
Entity.EntityType.JOB, JobEntity.class))
+ .thenReturn(job);
+
+ List<String> stdout = ImmutableList.of("line1", "line2");
+ List<String> stderr = ImmutableList.of("err1");
+ // The default gravitino.job.outputMaxLines (1000) / outputMaxBytes
(256KB) values are what
+ // JobManager should resolve and pass through, since the test config
doesn't override them.
+ when(jobExecutor.getJobStdout(job.jobExecutionId(), 1000, 256 *
1024)).thenReturn(stdout);
+ when(jobExecutor.getJobStderr(job.jobExecutionId(), 1000, 256 *
1024)).thenReturn(stderr);
+
+ // includeOutput = true fetches and attaches the output.
+ JobEntity jobWithOutput = jobManager.getJob(metalake, job.name(), true);
+ Assertions.assertEquals(stdout, jobWithOutput.stdout());
+ Assertions.assertEquals(stderr, jobWithOutput.stderr());
+ // The rest of the entity is unaffected.
+ Assertions.assertEquals(job.jobExecutionId(),
jobWithOutput.jobExecutionId());
+ Assertions.assertEquals(job.status(), jobWithOutput.status());
+
+ // includeOutput = false never touches the executor for output.
+ JobEntity jobWithoutOutput = jobManager.getJob(metalake, job.name(),
false);
+ Assertions.assertNull(jobWithoutOutput.stdout());
+ Assertions.assertNull(jobWithoutOutput.stderr());
+
+ verify(jobExecutor, times(1)).getJobStdout(job.jobExecutionId(), 1000, 256
* 1024);
+ verify(jobExecutor, times(1)).getJobStderr(job.jobExecutionId(), 1000, 256
* 1024);
+ }
+
+ @Test
+ public void testGetJobWithOutputPerRequestCapsAreClampedToGlobal() throws
IOException {
+ mockedMetalake
+ .when(() -> MetalakeManager.checkMetalake(metalakeIdent, entityStore))
+ .thenAnswer(a -> null);
+
+ JobEntity job = newJobEntity("shell_job", JobHandle.Status.SUCCEEDED);
+ when(entityStore.get(
+ NameIdentifierUtil.ofJob(metalake, job.name()),
Entity.EntityType.JOB, JobEntity.class))
+ .thenReturn(job);
+
+ List<String> stdout = ImmutableList.of("line1");
+ List<String> stderr = ImmutableList.of();
+
+ // A per-request value smaller than the global default (1000 lines /
256KB) is honored as-is.
+ when(jobExecutor.getJobStdout(job.jobExecutionId(), 10,
1024)).thenReturn(stdout);
+ when(jobExecutor.getJobStderr(job.jobExecutionId(), 10,
1024)).thenReturn(stderr);
+ JobEntity narrower = jobManager.getJob(metalake, job.name(), true, 10,
1024);
+ Assertions.assertEquals(stdout, narrower.stdout());
+ verify(jobExecutor, times(1)).getJobStdout(job.jobExecutionId(), 10, 1024);
+ verify(jobExecutor, times(1)).getJobStderr(job.jobExecutionId(), 10, 1024);
+
+ // A per-request value larger than the global default is clamped down to
it - the global
+ // configuration remains a hard upper bound, not just a fallback default.
+ when(jobExecutor.getJobStdout(job.jobExecutionId(), 1000, 256 *
1024)).thenReturn(stdout);
+ when(jobExecutor.getJobStderr(job.jobExecutionId(), 1000, 256 *
1024)).thenReturn(stderr);
+ JobEntity clamped =
+ jobManager.getJob(metalake, job.name(), true, 1_000_000, 1024 * 1024 *
1024);
+ Assertions.assertEquals(stdout, clamped.stdout());
+ verify(jobExecutor, times(1)).getJobStdout(job.jobExecutionId(), 1000, 256
* 1024);
+ verify(jobExecutor, times(1)).getJobStderr(job.jobExecutionId(), 1000, 256
* 1024);
+
+ // Non-positive per-request values are rejected outright.
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () -> jobManager.getJob(metalake, job.name(), true, 0, null));
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () -> jobManager.getJob(metalake, job.name(), true, null, -1));
+ }
+
+ @Test
+ public void testGetJobIgnoresInvalidOutputCapsWhenOutputNotRequested()
throws IOException {
+ mockedMetalake
+ .when(() -> MetalakeManager.checkMetalake(metalakeIdent, entityStore))
+ .thenAnswer(a -> null);
+
+ JobEntity job = newJobEntity("shell_job", JobHandle.Status.SUCCEEDED);
+ when(entityStore.get(
+ NameIdentifierUtil.ofJob(metalake, job.name()),
Entity.EntityType.JOB, JobEntity.class))
+ .thenReturn(job);
+
+ // maxLines/maxBytes are documented as ignored when includeOutput is
false, so a non-positive
+ // value here must not fail the call - it's never even inspected.
+ JobEntity result = jobManager.getJob(metalake, job.name(), false, 0, -1);
+ Assertions.assertEquals(job.jobExecutionId(), result.jobExecutionId());
+ Assertions.assertNull(result.stdout());
+ Assertions.assertNull(result.stderr());
+ verify(jobExecutor, never()).getJobStdout(any(), anyInt(), anyInt());
+ verify(jobExecutor, never()).getJobStderr(any(), anyInt(), anyInt());
+ }
+
+ @Test
+ public void testGetJobWithOutputDegradesToEmptyWhenExecutorForgetsJob()
throws IOException {
+ // The job entity is confirmed to exist (via the entity store) before the
executor is ever
+ // consulted for output. The executor's own bookkeeping for a job's output
is separate and can
+ // legitimately expire or be lost (e.g. the local executor's in-memory
state ages out
+ // independently of the job entity/staging directory) -
JobExecutor#getJobStdout/getJobStderr
+ // report that as empty output, not an error, so getJob(..., true) must
still return the full
+ // entity with empty output rather than treat the job as missing.
+ mockedMetalake
+ .when(() -> MetalakeManager.checkMetalake(metalakeIdent, entityStore))
+ .thenAnswer(a -> null);
+
+ JobEntity job = newJobEntity("shell_job", JobHandle.Status.SUCCEEDED);
+ when(entityStore.get(
+ NameIdentifierUtil.ofJob(metalake, job.name()),
Entity.EntityType.JOB, JobEntity.class))
+ .thenReturn(job);
+
+ when(jobExecutor.getJobStdout(job.jobExecutionId(), 1000, 256 * 1024))
+ .thenReturn(Collections.emptyList());
+ when(jobExecutor.getJobStderr(job.jobExecutionId(), 1000, 256 * 1024))
+ .thenReturn(Collections.emptyList());
+
+ JobEntity jobWithOutput = jobManager.getJob(metalake, job.name(), true);
+ Assertions.assertEquals(Collections.emptyList(), jobWithOutput.stdout());
+ Assertions.assertEquals(Collections.emptyList(), jobWithOutput.stderr());
+ // The entity itself is still fully returned, not treated as missing.
+ Assertions.assertEquals(job.jobExecutionId(),
jobWithOutput.jobExecutionId());
+ Assertions.assertEquals(job.status(), jobWithOutput.status());
}
@Test
@@ -790,7 +919,7 @@ public class TestJobManager {
@Test
public void testCancelJobDoesNotReplayExecutorOnOccConflict() throws
IOException {
JobEntity job = newJobEntity("shell_job", JobHandle.Status.QUEUED);
- when(jobManager.getJob(metalake, job.name())).thenReturn(job);
+ when(jobManager.getJob(metalake, job.name(), false)).thenReturn(job);
doNothing().when(jobExecutor).cancelJob(job.jobExecutionId());
when(entityStore.update(any(), eq(JobEntity.class),
eq(Entity.EntityType.JOB), any()))
.thenThrow(new OptimisticLockException("job changed"));
@@ -806,7 +935,7 @@ public class TestJobManager {
.thenAnswer(a -> null);
JobEntity job = newJobEntity("shell_job", JobHandle.Status.QUEUED);
- when(jobManager.getJob(metalake, job.name())).thenReturn(job);
+ when(jobManager.getJob(metalake, job.name(), false)).thenReturn(job);
doNothing().when(jobExecutor).cancelJob(job.jobExecutionId());
stubEntityStoreUpdateToApply(job);
@@ -816,7 +945,7 @@ public class TestJobManager {
Assertions.assertEquals(JobHandle.Status.CANCELLING,
cancelledJob.status());
// Test cancel a nonexistent job
- when(jobManager.getJob(metalake, "non_existent"))
+ when(jobManager.getJob(metalake, "non_existent", false))
.thenThrow(new NoSuchJobException("Job does not exist"));
Exception e =
@@ -830,7 +959,7 @@ public class TestJobManager {
.forEach(
status -> {
JobEntity finishedJob = newJobEntity("shell_job", status);
- when(jobManager.getJob(metalake,
finishedJob.name())).thenReturn(finishedJob);
+ when(jobManager.getJob(metalake, finishedJob.name(),
false)).thenReturn(finishedJob);
JobEntity cancelledFinishedJob = jobManager.cancelJob(metalake,
finishedJob.name());
Assertions.assertEquals(
@@ -861,7 +990,7 @@ public class TestJobManager {
.thenAnswer(a -> null);
JobEntity job = newJobEntity("shell_job", JobHandle.Status.QUEUED);
- when(jobManager.getJob(metalake, job.name())).thenReturn(job);
+ when(jobManager.getJob(metalake, job.name(), false)).thenReturn(job);
doNothing().when(jobExecutor).cancelJob(job.jobExecutionId());
// Simulate the job having been deleted concurrently (e.g. by
legacy-timeline cleanup) in the
@@ -883,7 +1012,7 @@ public class TestJobManager {
// the time entityStore.update() re-fetches the entity, a concurrent
status poll has already
// persisted a terminal status - that must not be regressed back to
CANCELLING.
JobEntity queuedSnapshot = newJobEntity("shell_job",
JobHandle.Status.QUEUED);
- when(jobManager.getJob(metalake,
queuedSnapshot.name())).thenReturn(queuedSnapshot);
+ when(jobManager.getJob(metalake, queuedSnapshot.name(),
false)).thenReturn(queuedSnapshot);
doNothing().when(jobExecutor).cancelJob(queuedSnapshot.jobExecutionId());
JobEntity latestSucceeded =
@@ -915,7 +1044,7 @@ public class TestJobManager {
// snapshot and this update; re-applying CANCELLING here must not stamp a
fresh
// lastModifiedTime over the entity the other writer already wrote.
JobEntity queuedSnapshot = newJobEntity("shell_job",
JobHandle.Status.QUEUED);
- when(jobManager.getJob(metalake,
queuedSnapshot.name())).thenReturn(queuedSnapshot);
+ when(jobManager.getJob(metalake, queuedSnapshot.name(),
false)).thenReturn(queuedSnapshot);
doNothing().when(jobExecutor).cancelJob(queuedSnapshot.jobExecutionId());
JobEntity latestCancelling =
@@ -958,7 +1087,7 @@ public class TestJobManager {
.withAuditInfo(
AuditInfo.builder().withCreator("test").withCreateTime(Instant.now()).build())
.build();
- when(jobManager.getJob(metalake, job.name())).thenReturn(job);
+ when(jobManager.getJob(metalake, job.name(), false)).thenReturn(job);
doNothing().when(jobExecutor).cancelJob(job.jobExecutionId());
stubEntityStoreUpdateToApply(job);
@@ -1411,7 +1540,7 @@ public class TestJobManager {
JobEntity job =
newJobEntity("local-job-other-1", JobHandle.Status.STARTED,
Instant.now(), null);
- doReturn(job).when(jobManager).getJob(metalake, job.name());
+ doReturn(job).when(jobManager).getJob(metalake, job.name(), false);
when(jobExecutor.ownsJob(job.jobExecutionId())).thenReturn(false);
stubEntityStoreUpdateToApply(job);
@@ -1428,7 +1557,7 @@ public class TestJobManager {
.thenAnswer(a -> null);
JobEntity job = newJobEntity("local-job-mine-1", JobHandle.Status.STARTED,
Instant.now(), null);
- doReturn(job).when(jobManager).getJob(metalake, job.name());
+ doReturn(job).when(jobManager).getJob(metalake, job.name(), false);
doThrow(new
NoSuchJobException("lost")).when(jobExecutor).cancelJob(job.jobExecutionId());
stubEntityStoreUpdateToApply(job);
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 09bf25a3ca..9c5c705173 100644
--- a/core/src/test/java/org/apache/gravitino/job/TestJobManagerMultiNode.java
+++ b/core/src/test/java/org/apache/gravitino/job/TestJobManagerMultiNode.java
@@ -210,7 +210,7 @@ public class TestJobManagerMultiNode extends
TestJDBCBackend {
private JobEntity getJob(String jobName) {
// Read through node B, as it works the same from any node sharing the
metadata store.
- return nodeB.getJob(METALAKE, jobName);
+ return nodeB.getJob(METALAKE, jobName, false);
}
// Simulates that the keep time has elapsed since the job was last updated,
or finished.
diff --git
a/core/src/test/java/org/apache/gravitino/job/TestJobTemplateValidationDispatcher.java
b/core/src/test/java/org/apache/gravitino/job/TestJobTemplateValidationDispatcher.java
index 0f27df37cb..4c2d2e64e8 100644
---
a/core/src/test/java/org/apache/gravitino/job/TestJobTemplateValidationDispatcher.java
+++
b/core/src/test/java/org/apache/gravitino/job/TestJobTemplateValidationDispatcher.java
@@ -212,8 +212,8 @@ public class TestJobTemplateValidationDispatcher {
@Test
public void testGetJob() {
- validationDispatcher.getJob("metalake1", "job-123");
- verify(mockDispatcher).getJob("metalake1", "job-123");
+ validationDispatcher.getJob("metalake1", "job-123", false);
+ verify(mockDispatcher).getJob("metalake1", "job-123", false, null, null);
}
@Test
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 ec5f4d58d8..5a2dd71200 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,18 +18,23 @@
*/
package org.apache.gravitino.job.local;
+import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableMap;
import com.google.common.collect.Lists;
+import java.io.ByteArrayOutputStream;
import java.io.File;
import java.io.IOException;
import java.lang.reflect.Field;
import java.net.URL;
+import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.util.Collections;
+import java.util.List;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import org.apache.commons.io.FileUtils;
+import org.apache.commons.lang3.StringUtils;
import org.apache.gravitino.connector.job.JobExecutor;
import org.apache.gravitino.exceptions.NoSuchJobException;
import org.apache.gravitino.job.JobHandle;
@@ -50,6 +55,8 @@ import org.junit.jupiter.api.Test;
public class TestLocalJobExecutor {
+ private static final int DEFAULT_TEST_MAX_BYTES = 1_000_000;
+
private static JobExecutor jobExecutor;
private static JobTemplateEntity jobTemplateEntity;
@@ -259,6 +266,330 @@ public class TestLocalJobExecutor {
}
}
+ @Test
+ public void testGetJobOutputSuccessfully() throws IOException {
+ Map<String, String> jobConf =
+ ImmutableMap.of(
+ "arg1", "value1",
+ "arg2", "success",
+ "var", "value3");
+
+ JobTemplate template =
+ JobManager.createRuntimeJobTemplate(jobTemplateEntity, jobConf,
workingDir);
+
+ String jobId = jobExecutor.submitJob(template);
+ Awaitility.await()
+ .atMost(3, TimeUnit.MINUTES)
+ .until(() -> jobExecutor.getJobStatus(jobId) ==
JobHandle.Status.SUCCEEDED);
+
+ List<String> stdout = jobExecutor.getJobStdout(jobId, 1000,
DEFAULT_TEST_MAX_BYTES);
+ Assertions.assertEquals(6, stdout.size());
+ Assertions.assertEquals("starting test test job", stdout.get(0));
+ Assertions.assertEquals("in common script", stdout.get(1));
+ Assertions.assertTrue(stdout.get(2).startsWith("Submitting job with
name:"));
+ Assertions.assertEquals("value1", stdout.get(3));
+ Assertions.assertEquals("success", stdout.get(4));
+ Assertions.assertEquals("value3", stdout.get(5));
+
+ // The test script never writes to stderr.
+ Assertions.assertEquals(
+ Collections.emptyList(), jobExecutor.getJobStderr(jobId, 1000,
DEFAULT_TEST_MAX_BYTES));
+
+ // The full output has 6 lines; only the last 3 should be returned when
capped.
+ Assertions.assertEquals(
+ ImmutableList.of("value1", "success", "value3"),
+ jobExecutor.getJobStdout(jobId, 3, DEFAULT_TEST_MAX_BYTES));
+ }
+
+ @Test
+ public void testGetJobOutputForUnknownJobReturnsEmpty() {
+ // A job unknown to this executor - whether it never existed here, or its
bookkeeping has
+ // expired/been lost - reports empty output rather than throwing: the job
entity itself may
+ // still exist, and querying its output must not turn that into an error.
+ Assertions.assertEquals(
+ Collections.emptyList(),
+ jobExecutor.getJobStdout("no-such-job", 100, DEFAULT_TEST_MAX_BYTES));
+ Assertions.assertEquals(
+ Collections.emptyList(),
+ jobExecutor.getJobStderr("no-such-job", 100, DEFAULT_TEST_MAX_BYTES));
+ }
+
+ @Test
+ public void
testGetJobOutputReturnsEmptyWhenFileDisappearsBetweenExistsCheckAndRead()
+ throws IOException {
+ Map<String, String> jobConf =
+ ImmutableMap.of(
+ "arg1", "value1",
+ "arg2", "success",
+ "var", "value3");
+
+ JobTemplate template =
+ JobManager.createRuntimeJobTemplate(jobTemplateEntity, jobConf,
workingDir);
+
+ String jobId = jobExecutor.submitJob(template);
+ Awaitility.await()
+ .atMost(3, TimeUnit.MINUTES)
+ .until(() -> jobExecutor.getJobStatus(jobId) ==
JobHandle.Status.SUCCEEDED);
+
+ // Simulates JobManager#cleanUpStagingDirs racing with this read: the file
existed moments
+ // ago (the exists() check would pass), but is no longer a readable
regular file by the time
+ // the actual read is attempted - opening it as a RandomAccessFile throws
+ // FileNotFoundException, which must degrade to empty output rather than
propagate as a
+ // RuntimeException/500.
+ File outputLog = new File(workingDir, "output.log");
+ Assertions.assertTrue(outputLog.delete());
+ Assertions.assertTrue(outputLog.mkdir());
+
+ List<String> stdout = jobExecutor.getJobStdout(jobId, 100,
DEFAULT_TEST_MAX_BYTES);
+ Assertions.assertEquals(Collections.emptyList(), stdout);
+ }
+
+ @Test
+ public void testGetJobOutputForQueuedJobReturnsEmpty() throws IOException {
+ LocalJobExecutor exec = new LocalJobExecutor();
+ exec.initialize(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();
+ try {
+ Map<String, String> jobConf =
+ ImmutableMap.of(
+ "arg1", "value1",
+ "arg2", "success",
+ "var", "value3");
+
+ // Submit two jobs to a single-threaded executor - the second one stays
QUEUED until the
+ // first (which sleeps for a few seconds) finishes.
+ JobTemplate templateA =
+ JobManager.createRuntimeJobTemplate(jobTemplateEntity, jobConf,
workingDirA);
+ JobTemplate templateB =
+ JobManager.createRuntimeJobTemplate(jobTemplateEntity, jobConf,
workingDirB);
+ exec.submitJob(templateA);
+ String jobIdB = exec.submitJob(templateB);
+
+ Assertions.assertEquals(JobHandle.Status.QUEUED,
exec.getJobStatus(jobIdB));
+ Assertions.assertEquals(
+ Collections.emptyList(), exec.getJobStdout(jobIdB, 100,
DEFAULT_TEST_MAX_BYTES));
+ Assertions.assertEquals(
+ Collections.emptyList(), exec.getJobStderr(jobIdB, 100,
DEFAULT_TEST_MAX_BYTES));
+
+ Awaitility.await()
+ .atMost(3, TimeUnit.MINUTES)
+ .until(() -> exec.getJobStatus(jobIdB) ==
JobHandle.Status.SUCCEEDED);
+ } finally {
+ exec.close();
+ FileUtils.deleteDirectory(workingDirA);
+ FileUtils.deleteDirectory(workingDirB);
+ }
+ }
+
+ @Test
+ public void testGetJobOutputWithOversizedSingleLineIsBoundedByMaxBytes()
throws IOException {
+ Map<String, String> jobConf =
+ ImmutableMap.of(
+ "arg1", "value1",
+ "arg2", "success",
+ "var", "value3");
+
+ JobTemplate template =
+ JobManager.createRuntimeJobTemplate(jobTemplateEntity, jobConf,
workingDir);
+
+ String jobId = jobExecutor.submitJob(template);
+ Awaitility.await()
+ .atMost(3, TimeUnit.MINUTES)
+ .until(() -> jobExecutor.getJobStatus(jobId) ==
JobHandle.Status.SUCCEEDED);
+
+ // Overwrite the captured output with a single line far larger than the
byte window, with no
+ // trailing newline - the pathological case a byte-bounded tail read must
stay safe against
+ // (a naive line-oriented reverse reader can degrade badly on content
shaped like this).
+ int maxBytes = 1024;
+ String oversizedLine = StringUtils.repeat('x', maxBytes * 4);
+ FileUtils.writeStringToFile(new File(workingDir, "output.log"),
oversizedLine, "UTF-8");
+
+ List<String> stdout = jobExecutor.getJobStdout(jobId, 1000, maxBytes);
+ Assertions.assertEquals(1, stdout.size());
+ Assertions.assertTrue(stdout.get(0).length() <= maxBytes);
+ Assertions.assertTrue(oversizedLine.endsWith(stdout.get(0)));
+ }
+
+ @Test
+ public void
testGetJobOutputWithOversizedSingleLineEndingInNewlineIsTruncatedNotEmpty()
+ throws IOException {
+ Map<String, String> jobConf =
+ ImmutableMap.of(
+ "arg1", "value1",
+ "arg2", "success",
+ "var", "value3");
+
+ JobTemplate template =
+ JobManager.createRuntimeJobTemplate(jobTemplateEntity, jobConf,
workingDir);
+
+ String jobId = jobExecutor.submitJob(template);
+ Awaitility.await()
+ .atMost(3, TimeUnit.MINUTES)
+ .until(() -> jobExecutor.getJobStatus(jobId) ==
JobHandle.Status.SUCCEEDED);
+
+ // Same oversized single line as above, but terminated with a newline -
the shape almost
+ // every real shell command (echo, printf '...\n') actually produces. The
trailing newline is
+ // then the only '\n' in the window, and must not be mistaken for a
boundary to a subsequent
+ // line - the truncated line content must still come back, not an empty
list.
+ int maxBytes = 1024;
+ String oversizedLine = StringUtils.repeat('x', maxBytes * 4);
+ FileUtils.writeStringToFile(new File(workingDir, "output.log"),
oversizedLine + "\n", "UTF-8");
+
+ List<String> stdout = jobExecutor.getJobStdout(jobId, 1000, maxBytes);
+ Assertions.assertEquals(1, stdout.size());
+ Assertions.assertFalse(stdout.get(0).isEmpty());
+ Assertions.assertTrue(stdout.get(0).length() <= maxBytes);
+ Assertions.assertTrue(oversizedLine.endsWith(stdout.get(0)));
+ }
+
+ @Test
+ public void testGetJobOutputWithOrdinaryLinesFollowedByOversizedFinalLine()
throws IOException {
+ Map<String, String> jobConf =
+ ImmutableMap.of(
+ "arg1", "value1",
+ "arg2", "success",
+ "var", "value3");
+
+ JobTemplate template =
+ JobManager.createRuntimeJobTemplate(jobTemplateEntity, jobConf,
workingDir);
+
+ String jobId = jobExecutor.submitJob(template);
+ Awaitility.await()
+ .atMost(3, TimeUnit.MINUTES)
+ .until(() -> jobExecutor.getJobStatus(jobId) ==
JobHandle.Status.SUCCEEDED);
+
+ // A few ordinary lines followed by a newline-terminated final line that
alone exceeds the
+ // byte window - the window falls entirely within that last line, so its
trailing newline is
+ // again the only one visible, and the truncated tail of that line must
still be returned.
+ int maxBytes = 1024;
+ String oversizedLine = StringUtils.repeat('y', maxBytes * 4);
+ String content =
+ "error: something failed\nstack frame 1\nstack frame 2\n" +
oversizedLine + "\n";
+ FileUtils.writeStringToFile(new File(workingDir, "output.log"), content,
"UTF-8");
+
+ List<String> stdout = jobExecutor.getJobStdout(jobId, 1000, maxBytes);
+ Assertions.assertEquals(1, stdout.size());
+ Assertions.assertFalse(stdout.get(0).isEmpty());
+ Assertions.assertTrue(stdout.get(0).length() <= maxBytes);
+ Assertions.assertTrue(oversizedLine.endsWith(stdout.get(0)));
+ }
+
+ @Test
+ public void testGetJobOutputWindowSmallerThanRequestedLines() throws
IOException {
+ Map<String, String> jobConf =
+ ImmutableMap.of(
+ "arg1", "value1",
+ "arg2", "success",
+ "var", "value3");
+
+ JobTemplate template =
+ JobManager.createRuntimeJobTemplate(jobTemplateEntity, jobConf,
workingDir);
+
+ String jobId = jobExecutor.submitJob(template);
+ Awaitility.await()
+ .atMost(3, TimeUnit.MINUTES)
+ .until(() -> jobExecutor.getJobStatus(jobId) ==
JobHandle.Status.SUCCEEDED);
+
+ // Many short lines whose total size exceeds the byte window - only the
lines that fit in
+ // the window should come back, even though maxLines asks for far more
than that.
+ int lineCount = 2000;
+ StringBuilder content = new StringBuilder();
+ for (int i = 0; i < lineCount; i++) {
+ content.append("line").append(i).append('\n');
+ }
+ FileUtils.writeStringToFile(new File(workingDir, "output.log"),
content.toString(), "UTF-8");
+
+ int maxBytes = 100;
+ List<String> stdout = jobExecutor.getJobStdout(jobId, lineCount, maxBytes);
+ Assertions.assertFalse(stdout.isEmpty());
+ Assertions.assertTrue(stdout.size() < lineCount);
+ // It's a tail read, so the very last written line must always be present.
+ Assertions.assertEquals("line" + (lineCount - 1), stdout.get(stdout.size()
- 1));
+ }
+
+ @Test
+ public void testGetJobOutputStripsCrlfLineEndings() throws IOException {
+ Map<String, String> jobConf =
+ ImmutableMap.of(
+ "arg1", "value1",
+ "arg2", "success",
+ "var", "value3");
+
+ JobTemplate template =
+ JobManager.createRuntimeJobTemplate(jobTemplateEntity, jobConf,
workingDir);
+
+ String jobId = jobExecutor.submitJob(template);
+ Awaitility.await()
+ .atMost(3, TimeUnit.MINUTES)
+ .until(() -> jobExecutor.getJobStatus(jobId) ==
JobHandle.Status.SUCCEEDED);
+
+ // Windows-style line endings must not leave a trailing '\r' on the
returned lines.
+ FileUtils.writeStringToFile(
+ new File(workingDir, "output.log"), "line1\r\nline2\r\nline3\r\n",
"UTF-8");
+
+ List<String> stdout = jobExecutor.getJobStdout(jobId, 100,
DEFAULT_TEST_MAX_BYTES);
+ Assertions.assertEquals(ImmutableList.of("line1", "line2", "line3"),
stdout);
+ }
+
+ @Test
+ public void testGetJobOutputKeepsAllLinesWhenWindowAlignsOnLineBoundary()
throws IOException {
+ Map<String, String> jobConf =
+ ImmutableMap.of(
+ "arg1", "value1",
+ "arg2", "success",
+ "var", "value3");
+
+ JobTemplate template =
+ JobManager.createRuntimeJobTemplate(jobTemplateEntity, jobConf,
workingDir);
+
+ String jobId = jobExecutor.submitJob(template);
+ Awaitility.await()
+ .atMost(3, TimeUnit.MINUTES)
+ .until(() -> jobExecutor.getJobStatus(jobId) ==
JobHandle.Status.SUCCEEDED);
+
+ // "line1\n" is exactly 6 bytes, so a maxBytes of 12 makes the window
start exactly at the
+ // beginning of "line2" - a genuine line boundary, not a partial line.
Both "line2" and
+ // "line3" must be returned, not just the last one.
+ FileUtils.writeStringToFile(
+ new File(workingDir, "output.log"), "line1\nline2\nline3\n", "UTF-8");
+
+ List<String> stdout = jobExecutor.getJobStdout(jobId, 100, 12);
+ Assertions.assertEquals(ImmutableList.of("line2", "line3"), stdout);
+ }
+
+ @Test
+ public void
testGetJobOutputAvoidsCorruptingMultiByteUtf8CharacterAtWindowStart()
+ throws IOException {
+ Map<String, String> jobConf =
+ ImmutableMap.of(
+ "arg1", "value1",
+ "arg2", "success",
+ "var", "value3");
+
+ JobTemplate template =
+ JobManager.createRuntimeJobTemplate(jobTemplateEntity, jobConf,
workingDir);
+
+ String jobId = jobExecutor.submitJob(template);
+ Awaitility.await()
+ .atMost(3, TimeUnit.MINUTES)
+ .until(() -> jobExecutor.getJobStatus(jobId) ==
JobHandle.Status.SUCCEEDED);
+
+ // "AAAAA" + "δΈ" (3-byte UTF-8: 0xE4 0xB8 0xAD) + "BBBBB", with no
newlines at all. With
+ // maxBytes=7 the computed window start lands on the second byte of the
multi-byte character
+ // - the read must skip forward to the next character boundary rather than
emitting a
+ // replacement character for the split-up bytes.
+ ByteArrayOutputStream content = new ByteArrayOutputStream();
+ content.write("AAAAA".getBytes(StandardCharsets.UTF_8));
+ content.write(new byte[] {(byte) 0xE4, (byte) 0xB8, (byte) 0xAD});
+ content.write("BBBBB".getBytes(StandardCharsets.UTF_8));
+ FileUtils.writeByteArrayToFile(new File(workingDir, "output.log"),
content.toByteArray());
+
+ List<String> stdout = jobExecutor.getJobStdout(jobId, 100, 7);
+ Assertions.assertEquals(ImmutableList.of("BBBBB"), stdout);
+ }
+
@Test
public void testCancelJob() throws InterruptedException {
Map<String, String> jobConf =
diff --git
a/core/src/test/java/org/apache/gravitino/listener/api/event/TestJobEventDispatcher.java
b/core/src/test/java/org/apache/gravitino/listener/api/event/TestJobEventDispatcher.java
index cd9eaf3a7a..e43b85dba6 100644
---
a/core/src/test/java/org/apache/gravitino/listener/api/event/TestJobEventDispatcher.java
+++
b/core/src/test/java/org/apache/gravitino/listener/api/event/TestJobEventDispatcher.java
@@ -20,6 +20,7 @@
package org.apache.gravitino.listener.api.event;
import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyBoolean;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
@@ -341,7 +342,7 @@ public class TestJobEventDispatcher {
@Test
void testGetJobEvent() {
- dispatcher.getJob("metalake", jobInfo.jobId());
+ dispatcher.getJob("metalake", jobInfo.jobId(), false);
PreEvent preEvent = dummyEventListener.popPreEvent();
Assertions.assertEquals(
@@ -438,7 +439,7 @@ public class TestJobEventDispatcher {
void testGetJobFailureEvent() {
Assertions.assertThrowsExactly(
GravitinoRuntimeException.class,
- () -> failureDispatcher.getJob("metalake", jobInfo.jobId()));
+ () -> failureDispatcher.getJob("metalake", jobInfo.jobId(), false));
Event event = dummyEventListener.popPostEvent();
Assertions.assertInstanceOf(GetJobFailureEvent.class, event);
Assertions.assertEquals(
@@ -510,7 +511,8 @@ public class TestJobEventDispatcher {
// Job operations
when(dispatcher.listJobs(any(String.class), any(Optional.class)))
.thenReturn(Collections.singletonList(jobEntity));
- when(dispatcher.getJob(any(String.class),
any(String.class))).thenReturn(jobEntity);
+ when(dispatcher.getJob(any(String.class), any(String.class),
anyBoolean(), any(), any()))
+ .thenReturn(jobEntity);
when(dispatcher.runJob(any(String.class), any(String.class),
any(Map.class)))
.thenReturn(jobEntity);
when(dispatcher.cancelJob(any(String.class),
any(String.class))).thenReturn(jobEntity);
diff --git
a/core/src/test/java/org/apache/gravitino/utils/TestMetadataObjectUtil.java
b/core/src/test/java/org/apache/gravitino/utils/TestMetadataObjectUtil.java
index 1354dcc4e7..e5362ab74e 100644
--- a/core/src/test/java/org/apache/gravitino/utils/TestMetadataObjectUtil.java
+++ b/core/src/test/java/org/apache/gravitino/utils/TestMetadataObjectUtil.java
@@ -339,7 +339,7 @@ public class TestMetadataObjectUtil {
verify(accessControlDispatcher).getRole("metalake", "role");
verify(tagDispatcher).getTag("metalake", "tag");
verify(policyDispatcher).getPolicy("metalake", "policy");
- verify(jobDispatcher).getJob("metalake", "job");
+ verify(jobDispatcher).getJob("metalake", "job", false);
verify(jobDispatcher).getJobTemplate("metalake", "template");
}
diff --git a/docs/gravitino-server-config.md b/docs/gravitino-server-config.md
index 824d29ebf3..2bb8a90d08 100644
--- a/docs/gravitino-server-config.md
+++ b/docs/gravitino-server-config.md
@@ -547,6 +547,8 @@ server, are documented with those services. See
| `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) |
### Key Management
diff --git a/docs/manage-jobs-in-gravitino.md b/docs/manage-jobs-in-gravitino.md
index e210530b2c..7167f3e67c 100644
--- a/docs/manage-jobs-in-gravitino.md
+++ b/docs/manage-jobs-in-gravitino.md
@@ -224,6 +224,58 @@ cancelling = client.cancel_job(job_id)
Cancelling is a request rather than an instant. The job moves to `CANCELLING`
and then to
`CANCELLED`, and one that finishes first keeps the status it finished with.
+### Get a Job's Output
+
+A job's captured stdout/stderr can be fetched alongside its metadata by asking
for it explicitly.
+Output is fetched live from the job executor on every call rather than stored
in Gravitino, so it's
+only included when requested - a plain `getJob`/`get_job` call, or
`listJobs`/`list_jobs`, never
+returns it.
+
+<Tabs groupId='language' queryString>
+<TabItem value="shell" label="REST">
+
+```shell
+curl -X GET -H "Accept: application/vnd.gravitino.v1+json" \
+
"http://localhost:8090/api/metalakes/example/jobs/runs/{job_id}?includeOutput=true"
+```
+
+</TabItem>
+<TabItem value="java" label="Java">
+
+```java
+JobHandle job = client.getJob(jobId, true);
+List<String> stdout = job.stdout();
+List<String> stderr = job.stderr();
+```
+
+</TabItem>
+<TabItem value="python" label="Python">
+
+```python
+job = client.get_job(job_id, include_output=True)
+stdout = job.stdout()
+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
+by `gravitino.job.outputMaxLines` (line count) and
`gravitino.job.outputMaxBytes` (byte size),
+whichever limit is hit first.
+
+The REST API also accepts `outputMaxLines`/`outputMaxBytes` query parameters
to request less output
+than these global caps for a single call (e.g. a quick check that doesn't need
the full 1000
+lines) - a value larger than the global cap is clamped down to it, so the
global configuration
+always remains a hard upper bound:
+
+```shell
+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"
+```
+
### Job System Configuration
Configure the job system through the `gravitino.conf` file. The following are
the
@@ -235,6 +287,8 @@ default configurations:
| `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
diff --git a/docs/open-api/jobs.yaml b/docs/open-api/jobs.yaml
index 40a20b5060..342ac6f4e2 100644
--- a/docs/open-api/jobs.yaml
+++ b/docs/open-api/jobs.yaml
@@ -279,6 +279,10 @@ paths:
summary: Get job
operationId: getJob
description: Returns the specified job information in the specified
metalake
+ parameters:
+ - $ref: "#/components/parameters/includeOutput"
+ - $ref: "#/components/parameters/outputMaxLines"
+ - $ref: "#/components/parameters/outputMaxBytes"
responses:
"200":
description: Returns the job object
@@ -359,6 +363,43 @@ components:
schema:
type: string
default: ""
+ includeOutput:
+ name: includeOutput
+ in: query
+ description: >
+ 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.
+ required: false
+ schema:
+ type: boolean
+ default: false
+
+ outputMaxLines:
+ name: outputMaxLines
+ in: query
+ description: >
+ The maximum number of (most recent) output lines to return, when
includeOutput is true.
+ Ignored otherwise. Can only narrow the globally configured
gravitino.job.outputMaxLines
+ cap, never widen it - a value larger than the global cap is clamped
down to it. Omit to
+ use the global default.
+ required: false
+ schema:
+ type: integer
+ minimum: 1
+
+ outputMaxBytes:
+ name: outputMaxBytes
+ in: query
+ description: >
+ The maximum number of (most recent) output bytes to read, when
includeOutput is true.
+ Ignored otherwise. Can only narrow the globally configured
gravitino.job.outputMaxBytes
+ cap, never widen it - a value larger than the global cap is clamped
down to it. Omit to
+ use the global default.
+ required: false
+ schema:
+ type: integer
+ minimum: 1
jobId:
name: jobId
in: path
@@ -613,6 +654,22 @@ components:
nullable: true
allOf:
- $ref: "#/components/schemas/JobTemplate"
+ stdout:
+ type: array
+ nullable: true
+ items:
+ type: string
+ description: >
+ The captured standard output of the job, as a list of lines. Only
populated when the
+ job was fetched with includeOutput=true; null otherwise.
+ stderr:
+ type: array
+ nullable: true
+ items:
+ type: string
+ description: >
+ The captured standard error output of the job, as a list of lines.
Only populated when
+ the job was fetched with includeOutput=true; null otherwise.
TemplateUpdate:
diff --git
a/server/src/main/java/org/apache/gravitino/server/web/rest/ExceptionHandlers.java
b/server/src/main/java/org/apache/gravitino/server/web/rest/ExceptionHandlers.java
index 0a545267d5..0c87f998b1 100644
---
a/server/src/main/java/org/apache/gravitino/server/web/rest/ExceptionHandlers.java
+++
b/server/src/main/java/org/apache/gravitino/server/web/rest/ExceptionHandlers.java
@@ -1028,6 +1028,9 @@ public class ExceptionHandlers {
} else if (e instanceof ForbiddenException) {
return Utils.forbidden(errorMsg, e);
+ } else if (e instanceof UnsupportedOperationException) {
+ return Utils.unsupportedOperation(errorMsg, e);
+
} else {
return super.handle(op, jobTemplate, parent, e);
}
@@ -1063,6 +1066,9 @@ public class ExceptionHandlers {
} else if (e instanceof ForbiddenException) {
return Utils.forbidden(errorMsg, e);
+ } else if (e instanceof UnsupportedOperationException) {
+ return Utils.unsupportedOperation(errorMsg, e);
+
} else {
return super.handle(op, jobTemplate, parent, e);
}
diff --git
a/server/src/main/java/org/apache/gravitino/server/web/rest/JobOperations.java
b/server/src/main/java/org/apache/gravitino/server/web/rest/JobOperations.java
index ad9b7b5b70..863e619d4c 100644
---
a/server/src/main/java/org/apache/gravitino/server/web/rest/JobOperations.java
+++
b/server/src/main/java/org/apache/gravitino/server/web/rest/JobOperations.java
@@ -381,14 +381,19 @@ public class JobOperations {
public Response getJob(
@PathParam("metalake") @AuthorizationMetadata(type =
Entity.EntityType.METALAKE)
String metalake,
- @PathParam("jobId") @AuthorizationMetadata(type = Entity.EntityType.JOB)
String jobId) {
+ @PathParam("jobId") @AuthorizationMetadata(type = Entity.EntityType.JOB)
String jobId,
+ @QueryParam("includeOutput") @DefaultValue("false") boolean
includeOutput,
+ @QueryParam("outputMaxLines") Integer outputMaxLines,
+ @QueryParam("outputMaxBytes") Integer outputMaxBytes) {
LOG.info("Received request to get job {} in metalake {}", jobId, metalake);
try {
return Utils.doAs(
httpRequest,
() -> {
- JobEntity jobEntity = jobOperationDispatcher.getJob(metalake,
jobId);
+ JobEntity jobEntity =
+ jobOperationDispatcher.getJob(
+ metalake, jobId, includeOutput, outputMaxLines,
outputMaxBytes);
LOG.info("Retrieved job {} in metalake: {}", jobId, metalake);
return Utils.ok(new JobResponse(toDTO(jobEntity)));
});
@@ -536,7 +541,9 @@ public class JobOperations {
jobEntity.auditInfo().createTime(),
jobEntity.startedAtAsInstant(),
jobEntity.finishedAtAsInstant(),
- toRuntimeJobTemplateDTO(jobEntity));
+ toRuntimeJobTemplateDTO(jobEntity),
+ jobEntity.stdout(),
+ jobEntity.stderr());
}
/**
diff --git
a/server/src/test/java/org/apache/gravitino/server/web/rest/TestJobOperations.java
b/server/src/test/java/org/apache/gravitino/server/web/rest/TestJobOperations.java
index 399bd76eea..5714ec0d20 100644
---
a/server/src/test/java/org/apache/gravitino/server/web/rest/TestJobOperations.java
+++
b/server/src/test/java/org/apache/gravitino/server/web/rest/TestJobOperations.java
@@ -22,10 +22,13 @@ import static
javax.ws.rs.core.MediaType.APPLICATION_JSON_TYPE;
import static org.apache.gravitino.Configs.CACHE_ENABLED;
import static org.apache.gravitino.Configs.ENABLE_AUTHORIZATION;
import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.doNothing;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import com.codahale.metrics.annotation.ResponseMetered;
@@ -653,6 +656,23 @@ public class TestJobOperations extends JerseyTest {
Assertions.assertEquals(ErrorConstants.NOT_FOUND_CODE,
errorResp4.getCode());
Assertions.assertEquals(NoSuchJobTemplateException.class.getSimpleName(),
errorResp4.getType());
+ // Test throw UnsupportedOperationException
+ doThrow(new UnsupportedOperationException("mock error"))
+ .when(jobOperationDispatcher)
+ .alterJobTemplate(any(), any(), any());
+
+ Response resp5b =
+ target(jobTemplatePath())
+ .path(templateName)
+ .request(APPLICATION_JSON_TYPE)
+ .accept("application/vnd.gravitino.v1+json")
+ .put(Entity.entity(req, APPLICATION_JSON_TYPE));
+
+ Assertions.assertEquals(Response.Status.NOT_IMPLEMENTED.getStatusCode(),
resp5b.getStatus());
+ ErrorResponse errorResp4b = resp5b.readEntity(ErrorResponse.class);
+ Assertions.assertEquals(
+ UnsupportedOperationException.class.getSimpleName(),
errorResp4b.getType());
+
// Test throw RuntimeException
doThrow(new RuntimeException("mock error"))
.when(jobOperationDispatcher)
@@ -1094,6 +1114,126 @@ public class TestJobOperations extends JerseyTest {
return Instant.ofEpochMilli(epochMilli);
}
+ @Test
+ public void testGetJob() {
+ JobEntity job = newJobEntity("shell_template_1",
JobHandle.Status.SUCCEEDED);
+
+ when(jobOperationDispatcher.getJob(metalake, job.name(), false, null,
null)).thenReturn(job);
+
+ Response resp =
+ target(jobRunPath())
+ .path(job.name())
+ .request(APPLICATION_JSON_TYPE)
+ .accept("application/vnd.gravitino.v1+json")
+ .get();
+
+ Assertions.assertEquals(Response.Status.OK.getStatusCode(),
resp.getStatus());
+ JobResponse jobResp = resp.readEntity(JobResponse.class);
+ Assertions.assertEquals(0, jobResp.getCode());
+ Assertions.assertEquals(JobOperations.toDTO(job), jobResp.getJob());
+ // includeOutput defaults to false, so no output is fetched or returned.
+ Assertions.assertNull(jobResp.getJob().stdout());
+ Assertions.assertNull(jobResp.getJob().stderr());
+
+ verify(jobOperationDispatcher, never()).getJob(any(), any(), eq(true),
any(), any());
+ }
+
+ @Test
+ public void testGetJobWithOutput() {
+ JobEntity job = newJobEntity("shell_template_1",
JobHandle.Status.SUCCEEDED);
+ List<String> stdout = Lists.newArrayList("line1", "line2");
+ List<String> stderr = Lists.newArrayList("err1");
+ JobEntity jobWithOutput = job.withOutput(stdout, stderr);
+
+ when(jobOperationDispatcher.getJob(metalake, job.name(), true, null, null))
+ .thenReturn(jobWithOutput);
+
+ Response resp =
+ target(jobRunPath())
+ .path(job.name())
+ .queryParam("includeOutput", "true")
+ .request(APPLICATION_JSON_TYPE)
+ .accept("application/vnd.gravitino.v1+json")
+ .get();
+
+ Assertions.assertEquals(Response.Status.OK.getStatusCode(),
resp.getStatus());
+ JobResponse jobResp = resp.readEntity(JobResponse.class);
+ Assertions.assertEquals(0, jobResp.getCode());
+ Assertions.assertEquals(stdout, jobResp.getJob().stdout());
+ Assertions.assertEquals(stderr, jobResp.getJob().stderr());
+ }
+
+ @Test
+ public void testGetJobWithOutputCustomLimits() {
+ JobEntity job = newJobEntity("shell_template_1",
JobHandle.Status.SUCCEEDED);
+ List<String> stdout = Lists.newArrayList("line1");
+ List<String> stderr = Lists.newArrayList();
+ JobEntity jobWithOutput = job.withOutput(stdout, stderr);
+
+ // The outputMaxLines/outputMaxBytes query parameters are passed straight
through to the
+ // dispatcher as-is; clamping against the global configuration happens
downstream in
+ // JobManager, not in the REST layer.
+ when(jobOperationDispatcher.getJob(metalake, job.name(), true, 10, 1024))
+ .thenReturn(jobWithOutput);
+
+ Response resp =
+ target(jobRunPath())
+ .path(job.name())
+ .queryParam("includeOutput", "true")
+ .queryParam("outputMaxLines", "10")
+ .queryParam("outputMaxBytes", "1024")
+ .request(APPLICATION_JSON_TYPE)
+ .accept("application/vnd.gravitino.v1+json")
+ .get();
+
+ Assertions.assertEquals(Response.Status.OK.getStatusCode(),
resp.getStatus());
+ JobResponse jobResp = resp.readEntity(JobResponse.class);
+ Assertions.assertEquals(stdout, jobResp.getJob().stdout());
+ Assertions.assertEquals(stderr, jobResp.getJob().stderr());
+ }
+
+ @Test
+ public void testGetJobWithOutputUnsupportedOperation() {
+ JobEntity job = newJobEntity("shell_template_1",
JobHandle.Status.SUCCEEDED);
+
+ doThrow(new UnsupportedOperationException("output retrieval not
supported"))
+ .when(jobOperationDispatcher)
+ .getJob(metalake, job.name(), true, null, null);
+
+ Response resp =
+ target(jobRunPath())
+ .path(job.name())
+ .queryParam("includeOutput", "true")
+ .request(APPLICATION_JSON_TYPE)
+ .accept("application/vnd.gravitino.v1+json")
+ .get();
+
+ Assertions.assertEquals(Response.Status.NOT_IMPLEMENTED.getStatusCode(),
resp.getStatus());
+ ErrorResponse errorResp = resp.readEntity(ErrorResponse.class);
+ Assertions.assertEquals(
+ UnsupportedOperationException.class.getSimpleName(),
errorResp.getType());
+ }
+
+ @Test
+ public void testGetJobWithInvalidOutputMaxLines() {
+ JobEntity job = newJobEntity("shell_template_1",
JobHandle.Status.SUCCEEDED);
+
+ doThrow(new IllegalArgumentException("maxLines must be positive if
specified"))
+ .when(jobOperationDispatcher)
+ .getJob(metalake, job.name(), true, 0, null);
+
+ Response resp =
+ target(jobRunPath())
+ .path(job.name())
+ .queryParam("includeOutput", "true")
+ .queryParam("outputMaxLines", "0")
+ .request(APPLICATION_JSON_TYPE)
+ .accept("application/vnd.gravitino.v1+json")
+ .get();
+
+ Assertions.assertEquals(Response.Status.BAD_REQUEST.getStatusCode(),
resp.getStatus());
+ }
+
@Test
public void testRunJob() {
String templateName = "shell_template_1";