This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-11982-1ff72ceb3d2a823a0fca64ec009984d75e846c28 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit c44f5b3388a8664a997b915f2d4b55c580ab3c90 Author: SEZ <[email protected]> AuthorDate: Sat Aug 29 02:09:04 2026 +0000 [Improve][Zeta] Expose job/pipeline restart and state-transition signals in job-info REST (#11982) --- docs/en/engines/zeta/rest-api-v1.md | 3 + docs/en/engines/zeta/rest-api-v2.md | 45 ++++- docs/zh/engines/zeta/rest-api-v1.md | 1 + docs/zh/engines/zeta/rest-api-v2.md | 42 +++- .../org/apache/seatunnel/engine/e2e/RestApiIT.java | 34 ++++ .../server/diagnostic/JobRuntimeDiagnostics.java | 141 ++++++++++++++ .../operation/GetJobDiagnosticsOperation.java | 80 ++++++++ .../seatunnel/engine/server/rest/RestConstant.java | 3 + .../engine/server/rest/service/BaseService.java | 44 +++++ .../engine/server/rest/service/JobInfoService.java | 5 +- .../ClientToServerOperationDataSerializerHook.java | 4 + .../diagnostic/JobRuntimeDiagnosticsTest.java | 212 +++++++++++++++++++++ 12 files changed, 611 insertions(+), 3 deletions(-) diff --git a/docs/en/engines/zeta/rest-api-v1.md b/docs/en/engines/zeta/rest-api-v1.md index 8b715327e4..5c353868cd 100644 --- a/docs/en/engines/zeta/rest-api-v1.md +++ b/docs/en/engines/zeta/rest-api-v1.md @@ -292,6 +292,9 @@ network: `jobId`, `jobName`, `jobStatus`, `createTime`, `jobDag`, `metrics` always be returned. `envOptions`, `pluginJarsUrls`, `isStartWithSavePoint` will return when job is running. `finishedTime`, `errorMsg` will return when job is finished. +A running job also returns a `diagnostics` block (state timestamps and per-pipeline restore counts), +see [REST API V2](rest-api-v2.md) for its fields. Only this endpoint returns it, `/running-jobs` does +not. #### Metrics field description diff --git a/docs/en/engines/zeta/rest-api-v2.md b/docs/en/engines/zeta/rest-api-v2.md index 2b0bab3245..5bd9b821c2 100644 --- a/docs/en/engines/zeta/rest-api-v2.md +++ b/docs/en/engines/zeta/rest-api-v2.md @@ -526,13 +526,43 @@ This endpoint helps troubleshoot why jobs stay in `PENDING` by showing the pendi }, "pluginJarsUrls": [ ], - "isStartWithSavePoint": false + "isStartWithSavePoint": false, + "diagnostics": { + "jobId": "", + "generatedAt": 1755000004000, + "stateTimestamps": { + "INITIALIZING": 1755000000000, + "CREATED": 1755000000200, + "SCHEDULED": 1755000001000, + "RUNNING": 1755000003000 + }, + "pipelines": [ + { + "pipelineId": 1, + "pipelineStatus": "RUNNING", + "restoreCount": 7, + "maxRestoreCount": 100, + "stateTimestamps": { + "INITIALIZING": 1755000000000, + "CREATED": 1755000000200, + "SCHEDULED": 1755000001100, + "DEPLOYING": 1755000002000, + "RUNNING": 1755000003500 + } + } + ], + "totalPipelineRestoreCount": 7 + } } ``` `jobId`, `jobName`, `jobStatus`, `createTime`, `jobDag`, `metrics` always be returned. `envOptions`, `pluginJarsUrls`, `isStartWithSavePoint` will return when job is running. `finishedTime`, `errorMsg` will return when job is finished. +`diagnostics` will return when the job is running and its diagnostics can be read from the master +node. It is auxiliary information: if it can not be obtained, the field is omitted instead of +failing the request. Only this endpoint returns it; `/running-jobs` does not, because collecting it +for every running job would cost one more round trip to the master per job. #### Metrics field description @@ -555,6 +585,19 @@ This endpoint helps troubleshoot why jobs stay in `PENDING` by showing the pendi | TableSinkWrite* | Per-table sink write attempts, key format `TableSinkWriteXXX#<table>` | | TableSinkCommitted* | Per-table sink committed metrics, key format `TableSinkCommittedXXX#<table>` | +#### Diagnostics field description + +| Field | Description | +| --- | --- | +| generatedAt | Epoch millis when this diagnostics block was collected | +| stateTimestamps | Epoch millis when the job entered each state. States never entered are omitted. A pipeline restart does not change the job state, so this alone does not show restarts | +| pipelines[].pipelineId | Pipeline id inside the job | +| pipelines[].pipelineStatus | Current pipeline state | +| pipelines[].restoreCount | How many times this pipeline has been restored since the job was submitted. A value that keeps growing while `jobStatus` stays `RUNNING` is a crash loop | +| pipelines[].maxRestoreCount | Restore limit of this pipeline, from the `job.retry.times` env option | +| pipelines[].stateTimestamps | Epoch millis when the pipeline entered each state. After a restore, the timestamps of the new attempt overwrite the previous ones | +| totalPipelineRestoreCount | Sum of `restoreCount` over all pipelines of the job | + When we can't get the job info, the response will be: ```json diff --git a/docs/zh/engines/zeta/rest-api-v1.md b/docs/zh/engines/zeta/rest-api-v1.md index 18875e4ed1..f1d43021c9 100644 --- a/docs/zh/engines/zeta/rest-api-v1.md +++ b/docs/zh/engines/zeta/rest-api-v1.md @@ -312,6 +312,7 @@ network: `jobId`, `jobName`, `jobStatus`, `createTime`, `jobDag`, `metrics` 字段总会返回. `envOptions`, `pluginJarsUrls`, `isStartWithSavePoint` 字段在Job在RUNNING状态时会返回 `finishedTime`, `errorMsg` 字段在Job结束时会返回,结束状态为不为RUNNING,可能为FINISHED,可能为CANCEL +运行中的Job还会返回 `diagnostics` 字段(状态时间戳与各 Pipeline 的恢复次数),字段说明见 [REST API V2](rest-api-v2.md)。该字段只在本接口返回,`/running-jobs` 不返回。 #### 指标字段说明 diff --git a/docs/zh/engines/zeta/rest-api-v2.md b/docs/zh/engines/zeta/rest-api-v2.md index 56438798a1..81ba2da7be 100644 --- a/docs/zh/engines/zeta/rest-api-v2.md +++ b/docs/zh/engines/zeta/rest-api-v2.md @@ -518,13 +518,40 @@ seatunnel: }, "pluginJarsUrls": [ ], - "isStartWithSavePoint": false + "isStartWithSavePoint": false, + "diagnostics": { + "jobId": "", + "generatedAt": 1755000004000, + "stateTimestamps": { + "INITIALIZING": 1755000000000, + "CREATED": 1755000000200, + "SCHEDULED": 1755000001000, + "RUNNING": 1755000003000 + }, + "pipelines": [ + { + "pipelineId": 1, + "pipelineStatus": "RUNNING", + "restoreCount": 7, + "maxRestoreCount": 100, + "stateTimestamps": { + "INITIALIZING": 1755000000000, + "CREATED": 1755000000200, + "SCHEDULED": 1755000001100, + "DEPLOYING": 1755000002000, + "RUNNING": 1755000003500 + } + } + ], + "totalPipelineRestoreCount": 7 + } } ``` `jobId`, `jobName`, `jobStatus`, `createTime`, `jobDag`, `metrics` 字段总会返回. `envOptions`, `pluginJarsUrls`, `isStartWithSavePoint` 字段在Job在RUNNING状态时会返回 `finishedTime`, `errorMsg` 字段在Job结束时会返回,结束状态为不为RUNNING,可能为FINISHED,可能为CANCEL +`diagnostics` 字段在Job运行中且能从Master节点读取到诊断信息时返回。它属于辅助信息,读取失败时该字段会被省略,不会导致请求失败。该字段只在本接口返回,`/running-jobs` 不返回:为列表中的每个Job收集诊断信息会额外增加一次到Master的请求。 #### 指标字段说明 @@ -547,6 +574,19 @@ seatunnel: | TableSinkWrite* | 按表汇总的 Sink 写入尝试,键格式 `TableSinkWriteXXX#<表>` | | TableSinkCommitted* | 按表汇总的 Sink 已提交指标,键格式 `TableSinkCommittedXXX#<表>` | +#### 诊断字段说明 + +| 字段 | 说明 | +| --- | --- | +| generatedAt | 采集该诊断信息的时间戳(毫秒) | +| stateTimestamps | Job 进入各状态的时间戳(毫秒),未进入过的状态不会出现。Pipeline 重启不会改变 Job 状态,因此仅凭该字段看不出重启 | +| pipelines[].pipelineId | Job 内的 Pipeline id | +| pipelines[].pipelineStatus | Pipeline 当前状态 | +| pipelines[].restoreCount | 自 Job 提交以来该 Pipeline 被恢复(重启)的次数。若 `jobStatus` 一直是 `RUNNING` 而该值不断增长,说明处于崩溃重启循环中 | +| pipelines[].maxRestoreCount | 该 Pipeline 的恢复次数上限,来自 `job.retry.times` env 配置 | +| pipelines[].stateTimestamps | Pipeline 进入各状态的时间戳(毫秒)。发生恢复后,新一轮的时间戳会覆盖上一轮 | +| totalPipelineRestoreCount | Job 下所有 Pipeline 的 `restoreCount` 之和 | + 当我们查询不到这个Job时,返回结果为: ```json diff --git a/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/RestApiIT.java b/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/RestApiIT.java index 6788d6dea5..ab0accdfa3 100644 --- a/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/RestApiIT.java +++ b/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/RestApiIT.java @@ -684,6 +684,40 @@ public class RestApiIT { }); } + @Test + public void testGetJobDiagnosticsOfRunningJob() { + // the diagnostics of a running job must be identical no matter which member serves the + // request: the master builds them locally, other members fetch them from the master + ports.forEach( + (key, value) -> + given().get( + HOST + + value + + node1Config + .getEngineConfig() + .getHttpConfig() + .getContextPath() + + RestConstant.REST_URL_JOB_INFO + + "/" + + clientJobProxy.getJobId()) + .then() + .statusCode(200) + .body("jobStatus", equalTo("RUNNING")) + .body( + "diagnostics.jobId", + equalTo(Long.toString(clientJobProxy.getJobId()))) + .body("diagnostics.generatedAt", notNullValue()) + .body("diagnostics.stateTimestamps.RUNNING", notNullValue()) + .body("diagnostics.pipelines", hasSize(1)) + .body("diagnostics.pipelines[0].pipelineId", equalTo(1)) + .body("diagnostics.pipelines[0].pipelineStatus", equalTo("RUNNING")) + .body("diagnostics.pipelines[0].restoreCount", equalTo(0)) + .body( + "diagnostics.pipelines[0].stateTimestamps.RUNNING", + notNullValue()) + .body("diagnostics.totalPipelineRestoreCount", equalTo(0))); + } + @Test public void testOverview() { Arrays.asList(node2, node1) diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/diagnostic/JobRuntimeDiagnostics.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/diagnostic/JobRuntimeDiagnostics.java new file mode 100644 index 0000000000..5acbc8307b --- /dev/null +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/diagnostic/JobRuntimeDiagnostics.java @@ -0,0 +1,141 @@ +/* + * 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.seatunnel.engine.server.diagnostic; + +import org.apache.seatunnel.engine.common.Constant; +import org.apache.seatunnel.engine.common.job.JobStatus; +import org.apache.seatunnel.engine.core.job.PipelineStatus; +import org.apache.seatunnel.engine.server.SeaTunnelServer; +import org.apache.seatunnel.engine.server.dag.physical.PhysicalPlan; +import org.apache.seatunnel.engine.server.dag.physical.SubPlan; +import org.apache.seatunnel.engine.server.master.JobMaster; + +import com.hazelcast.internal.json.JsonArray; +import com.hazelcast.internal.json.JsonObject; +import com.hazelcast.map.IMap; +import lombok.extern.slf4j.Slf4j; + +import java.util.Collections; +import java.util.List; + +/** + * Builds the {@code diagnostics} block of the job-info REST response. + * + * <p>Only signals the engine already tracks are exposed: the per-state entry timestamps kept in + * {@link Constant#IMAP_STATE_TIMESTAMPS} (job level and pipeline level) and the pipeline restore + * counter kept in {@link SubPlan}. Nothing here is recorded on the state transition path, so this + * is a pure read. + * + * <p>The pipeline part is only available where the {@link JobMaster} lives, therefore this must be + * called on the master node (REST callers on other members go through {@code + * GetJobDiagnosticsOperation}). + */ +@Slf4j +public final class JobRuntimeDiagnostics { + + public static final String JOB_ID = "jobId"; + public static final String GENERATED_AT = "generatedAt"; + public static final String STATE_TIMESTAMPS = "stateTimestamps"; + public static final String PIPELINES = "pipelines"; + public static final String PIPELINE_ID = "pipelineId"; + public static final String PIPELINE_STATUS = "pipelineStatus"; + public static final String RESTORE_COUNT = "restoreCount"; + public static final String MAX_RESTORE_COUNT = "maxRestoreCount"; + public static final String TOTAL_PIPELINE_RESTORE_COUNT = "totalPipelineRestoreCount"; + + private JobRuntimeDiagnostics() {} + + /** Collects the diagnostics of one job from the state timestamps map and the physical plan. */ + public static JsonObject build(SeaTunnelServer server, long jobId) { + JsonObject root = new JsonObject(); + root.add(JOB_ID, String.valueOf(jobId)); + root.add(GENERATED_AT, System.currentTimeMillis()); + + IMap<Object, Long[]> stateTimestampsMap = + server.getNodeEngine() + .getHazelcastInstance() + .getMap(Constant.IMAP_STATE_TIMESTAMPS); + root.add( + STATE_TIMESTAMPS, + toTimestampsJson(stateTimestampsMap.get(jobId), JobStatus.values())); + + JsonArray pipelines = new JsonArray(); + root.add(PIPELINES, pipelines); + + int totalRestoreCount = 0; + for (SubPlan subPlan : pipelineList(server, jobId)) { + PipelineStatus pipelineStatus = subPlan.getPipelineState(); + int restoreCount = subPlan.getPipelineRestoreNum(); + totalRestoreCount += restoreCount; + pipelines.add( + new JsonObject() + .add(PIPELINE_ID, subPlan.getPipelineId()) + .add( + PIPELINE_STATUS, + pipelineStatus == null ? null : pipelineStatus.toString()) + .add(RESTORE_COUNT, restoreCount) + .add(MAX_RESTORE_COUNT, subPlan.getPipelineMaxRestoreNum()) + .add( + STATE_TIMESTAMPS, + toTimestampsJson( + stateTimestampsMap.get(subPlan.getPipelineLocation()), + PipelineStatus.values()))); + } + root.add(TOTAL_PIPELINE_RESTORE_COUNT, totalRestoreCount); + return root; + } + + /** + * Returns the pipelines of a running job, or an empty list when the job is no longer + * coordinated by this member (finished job, master switch in progress). + */ + private static List<SubPlan> pipelineList(SeaTunnelServer server, long jobId) { + try { + JobMaster jobMaster = server.getCoordinatorService().getJobMaster(jobId); + if (jobMaster == null) { + return Collections.emptyList(); + } + PhysicalPlan physicalPlan = jobMaster.getPhysicalPlan(); + if (physicalPlan == null) { + return Collections.emptyList(); + } + return physicalPlan.getPipelineList(); + } catch (Throwable t) { + log.debug("Get pipeline diagnostics of job {} failed: {}", jobId, t.getMessage()); + return Collections.emptyList(); + } + } + + /** + * Renders a {@code stateTimestamps} array as a state name to epoch millis object, skipping + * states that were never entered. + */ + private static JsonObject toTimestampsJson(Long[] stateTimestamps, Enum<?>[] states) { + JsonObject json = new JsonObject(); + if (stateTimestamps == null) { + return json; + } + for (Enum<?> state : states) { + int ordinal = state.ordinal(); + if (ordinal < stateTimestamps.length && stateTimestamps[ordinal] != null) { + json.add(state.name(), stateTimestamps[ordinal].longValue()); + } + } + return json; + } +} diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/operation/GetJobDiagnosticsOperation.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/operation/GetJobDiagnosticsOperation.java new file mode 100644 index 0000000000..41c9b0bcb1 --- /dev/null +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/operation/GetJobDiagnosticsOperation.java @@ -0,0 +1,80 @@ +/* + * 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.seatunnel.engine.server.operation; + +import org.apache.seatunnel.engine.server.SeaTunnelServer; +import org.apache.seatunnel.engine.server.diagnostic.JobRuntimeDiagnostics; +import org.apache.seatunnel.engine.server.serializable.ClientToServerOperationDataSerializerHook; + +import com.hazelcast.nio.ObjectDataInput; +import com.hazelcast.nio.ObjectDataOutput; +import com.hazelcast.nio.serialization.IdentifiedDataSerializable; +import com.hazelcast.spi.impl.AllowedDuringPassiveState; +import com.hazelcast.spi.impl.operationservice.Operation; + +import java.io.IOException; + +/** + * Master operation that returns the runtime diagnostics (state timestamps, pipeline restore counts) + * of one job as a JSON string, so REST callers served by a non-master member get the same payload. + */ +public class GetJobDiagnosticsOperation extends Operation + implements IdentifiedDataSerializable, AllowedDuringPassiveState { + + private long jobId; + private String response; + + public GetJobDiagnosticsOperation() {} + + public GetJobDiagnosticsOperation(long jobId) { + this.jobId = jobId; + } + + @Override + public int getFactoryId() { + return ClientToServerOperationDataSerializerHook.FACTORY_ID; + } + + @Override + public int getClassId() { + return ClientToServerOperationDataSerializerHook.GET_JOB_DIAGNOSTICS_OPERATION; + } + + @Override + protected void writeInternal(ObjectDataOutput out) throws IOException { + super.writeInternal(out); + out.writeLong(jobId); + } + + @Override + protected void readInternal(ObjectDataInput in) throws IOException { + super.readInternal(in); + jobId = in.readLong(); + } + + @Override + public void run() { + SeaTunnelServer service = getService(); + response = JobRuntimeDiagnostics.build(service, jobId).toString(); + } + + @Override + public Object getResponse() { + return response; + } +} diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/RestConstant.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/RestConstant.java index f6da21b598..11964f1686 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/RestConstant.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/RestConstant.java @@ -57,6 +57,9 @@ public class RestConstant { public static final String ERROR_MSG = "errorMsg"; public static final String METRICS = "metrics"; + + public static final String DIAGNOSTICS = "diagnostics"; + public static final String LIMIT = "limit"; public static final String TABLE_SOURCE_RECEIVED_COUNT = "TableSourceReceivedCount"; diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/BaseService.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/BaseService.java index e2aab63fc6..ba7344de80 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/BaseService.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/BaseService.java @@ -46,9 +46,11 @@ import org.apache.seatunnel.engine.core.job.VertexInfo; import org.apache.seatunnel.engine.server.CoordinatorService; import org.apache.seatunnel.engine.server.SeaTunnelServer; import org.apache.seatunnel.engine.server.dag.DAGUtils; +import org.apache.seatunnel.engine.server.diagnostic.JobRuntimeDiagnostics; import org.apache.seatunnel.engine.server.master.JobHistoryService; import org.apache.seatunnel.engine.server.operation.CancelJobOperation; import org.apache.seatunnel.engine.server.operation.GetClusterHealthMetricsOperation; +import org.apache.seatunnel.engine.server.operation.GetJobDiagnosticsOperation; import org.apache.seatunnel.engine.server.operation.GetJobInfoOperation; import org.apache.seatunnel.engine.server.operation.GetJobMetricsOperation; import org.apache.seatunnel.engine.server.operation.GetJobStatusOperation; @@ -63,6 +65,7 @@ import com.hazelcast.cluster.Address; import com.hazelcast.cluster.Cluster; import com.hazelcast.cluster.Member; import com.hazelcast.instance.impl.Node; +import com.hazelcast.internal.json.Json; import com.hazelcast.internal.json.JsonArray; import com.hazelcast.internal.json.JsonObject; import com.hazelcast.internal.json.JsonValue; @@ -345,6 +348,16 @@ public abstract class BaseService { } protected JsonObject convertToJson(JobInfo jobInfo, long jobId) { + return convertToJson(jobInfo, jobId, true); + } + + /** + * @param withDiagnostics whether to add the {@code diagnostics} block. A listing of every + * running job builds this payload once per job, and every job already costs a master round + * trip when the request is not served by the master, so diagnostics are only collected for + * a request about one job. + */ + protected JsonObject convertToJson(JobInfo jobInfo, long jobId, boolean withDiagnostics) { JsonObject jobInfoJson = new JsonObject(); JobImmutableInformation jobImmutableInformation = @@ -451,6 +464,13 @@ public abstract class BaseService { RestConstant.METRICS, metricsToJsonObject(getJobMetrics(jobMetrics, jobDAGInfo))); + if (withDiagnostics) { + JsonObject diagnostics = getJobDiagnostics(jobId, seaTunnelServer); + if (diagnostics != null) { + jobInfoJson.add(RestConstant.DIAGNOSTICS, diagnostics); + } + } + if (jobStatus != null && jobStatus.isEndState()) { RUNNING_JOB_DAG_JSON_CACHE.remove(jobId); } @@ -458,6 +478,30 @@ public abstract class BaseService { return jobInfoJson; } + /** + * Reads the job runtime diagnostics from the master member, or {@code null} when they can not + * be obtained. Diagnostics are auxiliary information, so a failure here never fails the + * job-info response. + */ + private JsonObject getJobDiagnostics(long jobId, SeaTunnelServer masterSeaTunnelServer) { + try { + // the caller resolved it with getSeaTunnelServer(true), so it is either null or this + // node is the master and no further isMasterNode() check is needed + if (masterSeaTunnelServer != null) { + return JobRuntimeDiagnostics.build(masterSeaTunnelServer, jobId); + } + String response = + (String) + NodeEngineUtil.sendOperationToMasterNode( + nodeEngine, new GetJobDiagnosticsOperation(jobId)) + .join(); + return response == null ? null : Json.parse(response).asObject(); + } catch (Throwable t) { + log.debug("Get job {} diagnostics failed: {}", jobId, t.getMessage()); + return null; + } + } + private JobDAGInfo getRunningJobDAGInfo(long jobId, SeaTunnelServer masterSeaTunnelServer) { try { if (masterSeaTunnelServer != null && masterSeaTunnelServer.isMasterNode()) { diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/JobInfoService.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/JobInfoService.java index f5d6135ef6..bb73f58c47 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/JobInfoService.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/JobInfoService.java @@ -176,7 +176,10 @@ public class JobInfoService extends BaseService { Comparator.comparing( entry -> entry.getValue().getInitializationTimestamp(), Comparator.reverseOrder())) - .map(jobInfoEntry -> convertToJson(jobInfoEntry.getValue(), jobInfoEntry.getKey())) + .map( + jobInfoEntry -> + convertToJson( + jobInfoEntry.getValue(), jobInfoEntry.getKey(), false)) .collect(JsonArray::new, JsonArray::add, JsonArray::add); } diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/serializable/ClientToServerOperationDataSerializerHook.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/serializable/ClientToServerOperationDataSerializerHook.java index 4fb4e89b67..1b11469704 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/serializable/ClientToServerOperationDataSerializerHook.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/serializable/ClientToServerOperationDataSerializerHook.java @@ -24,6 +24,7 @@ import org.apache.seatunnel.engine.server.operation.GetCheckpointOverviewOperati import org.apache.seatunnel.engine.server.operation.GetClusterHealthMetricsOperation; import org.apache.seatunnel.engine.server.operation.GetJobCheckpointOperation; import org.apache.seatunnel.engine.server.operation.GetJobDetailStatusOperation; +import org.apache.seatunnel.engine.server.operation.GetJobDiagnosticsOperation; import org.apache.seatunnel.engine.server.operation.GetJobInfoOperation; import org.apache.seatunnel.engine.server.operation.GetJobMetricsOperation; import org.apache.seatunnel.engine.server.operation.GetJobStatusOperation; @@ -76,6 +77,7 @@ public final class ClientToServerOperationDataSerializerHook implements DataSeri public static final int GET_CHECKPOINT_HISTORY_OPERATION = 14; public static final int GET_NODE_HTTP_PORT_OPERATION = 15; public static final int GET_JOB_TASK_MAPPING_OPERATION = 16; + public static final int GET_JOB_DIAGNOSTICS_OPERATION = 17; public static final int FACTORY_ID = FactoryIdHelper.getFactoryId( @@ -135,6 +137,8 @@ public final class ClientToServerOperationDataSerializerHook implements DataSeri return new GetNodeHttpPortOperation(); case GET_JOB_TASK_MAPPING_OPERATION: return new GetJobTaskMappingOperation(); + case GET_JOB_DIAGNOSTICS_OPERATION: + return new GetJobDiagnosticsOperation(); default: throw new IllegalArgumentException("Unknown type id " + typeId); } diff --git a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/diagnostic/JobRuntimeDiagnosticsTest.java b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/diagnostic/JobRuntimeDiagnosticsTest.java new file mode 100644 index 0000000000..692e1fc2c5 --- /dev/null +++ b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/diagnostic/JobRuntimeDiagnosticsTest.java @@ -0,0 +1,212 @@ +/* + * 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.seatunnel.engine.server.diagnostic; + +import org.apache.seatunnel.engine.common.Constant; +import org.apache.seatunnel.engine.common.job.JobStatus; +import org.apache.seatunnel.engine.core.job.PipelineStatus; +import org.apache.seatunnel.engine.server.CoordinatorService; +import org.apache.seatunnel.engine.server.SeaTunnelServer; +import org.apache.seatunnel.engine.server.dag.physical.PhysicalPlan; +import org.apache.seatunnel.engine.server.dag.physical.PipelineLocation; +import org.apache.seatunnel.engine.server.dag.physical.SubPlan; +import org.apache.seatunnel.engine.server.master.JobMaster; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import com.hazelcast.core.HazelcastInstance; +import com.hazelcast.internal.json.JsonArray; +import com.hazelcast.internal.json.JsonObject; +import com.hazelcast.map.IMap; +import com.hazelcast.spi.impl.NodeEngineImpl; + +import java.util.Arrays; +import java.util.Collections; + +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +/** Covers the read-only job diagnostics payload exposed by the job-info REST endpoint. */ +public class JobRuntimeDiagnosticsTest { + + private static final long JOB_ID = 862858126931165185L; + + private SeaTunnelServer server; + private CoordinatorService coordinatorService; + private IMap<Object, Long[]> stateTimestampsMap; + + @BeforeEach + void setUp() { + server = mock(SeaTunnelServer.class); + coordinatorService = mock(CoordinatorService.class); + NodeEngineImpl nodeEngine = mock(NodeEngineImpl.class); + HazelcastInstance hazelcastInstance = mock(HazelcastInstance.class); + stateTimestampsMap = mock(IMap.class); + + when(server.getNodeEngine()).thenReturn(nodeEngine); + when(server.getCoordinatorService()).thenReturn(coordinatorService); + when(nodeEngine.getHazelcastInstance()).thenReturn(hazelcastInstance); + doReturn(stateTimestampsMap).when(hazelcastInstance).getMap(Constant.IMAP_STATE_TIMESTAMPS); + } + + private Long[] jobTimestamps() { + Long[] timestamps = new Long[JobStatus.values().length]; + timestamps[JobStatus.CREATED.ordinal()] = 1000L; + timestamps[JobStatus.SCHEDULED.ordinal()] = 2000L; + timestamps[JobStatus.RUNNING.ordinal()] = 3000L; + return timestamps; + } + + private SubPlan mockSubPlan(int pipelineId, PipelineStatus status, int restoreNum, int maxNum) { + SubPlan subPlan = mock(SubPlan.class); + when(subPlan.getPipelineId()).thenReturn(pipelineId); + when(subPlan.getPipelineState()).thenReturn(status); + when(subPlan.getPipelineRestoreNum()).thenReturn(restoreNum); + when(subPlan.getPipelineMaxRestoreNum()).thenReturn(maxNum); + when(subPlan.getPipelineLocation()).thenReturn(new PipelineLocation(JOB_ID, pipelineId)); + return subPlan; + } + + private void mockPipelines(SubPlan... subPlans) { + JobMaster jobMaster = mock(JobMaster.class); + PhysicalPlan physicalPlan = mock(PhysicalPlan.class); + when(coordinatorService.getJobMaster(JOB_ID)).thenReturn(jobMaster); + when(jobMaster.getPhysicalPlan()).thenReturn(physicalPlan); + when(physicalPlan.getPipelineList()).thenReturn(Arrays.asList(subPlans)); + } + + @Test + void testJobAndPipelineSignalsAreExposed() { + when(stateTimestampsMap.get(JOB_ID)).thenReturn(jobTimestamps()); + Long[] pipelineTimestamps = new Long[PipelineStatus.values().length]; + pipelineTimestamps[PipelineStatus.RUNNING.ordinal()] = 3500L; + when(stateTimestampsMap.get(new PipelineLocation(JOB_ID, 1))) + .thenReturn(pipelineTimestamps); + mockPipelines( + mockSubPlan(1, PipelineStatus.RUNNING, 7, 100), + mockSubPlan(2, PipelineStatus.RUNNING, 2, 100)); + + JsonObject diagnostics = JobRuntimeDiagnostics.build(server, JOB_ID); + + Assertions.assertEquals( + String.valueOf(JOB_ID), diagnostics.getString(JobRuntimeDiagnostics.JOB_ID, null)); + JsonObject jobStateTimestamps = + diagnostics.get(JobRuntimeDiagnostics.STATE_TIMESTAMPS).asObject(); + Assertions.assertEquals(1000L, jobStateTimestamps.getLong(JobStatus.CREATED.name(), -1L)); + Assertions.assertEquals(3000L, jobStateTimestamps.getLong(JobStatus.RUNNING.name(), -1L)); + // states never entered must not be rendered at all + Assertions.assertNull(jobStateTimestamps.get(JobStatus.FAILED.name())); + + JsonArray pipelines = diagnostics.get(JobRuntimeDiagnostics.PIPELINES).asArray(); + Assertions.assertEquals(2, pipelines.size()); + JsonObject firstPipeline = pipelines.get(0).asObject(); + Assertions.assertEquals(1, firstPipeline.getInt(JobRuntimeDiagnostics.PIPELINE_ID, -1)); + Assertions.assertEquals( + PipelineStatus.RUNNING.name(), + firstPipeline.getString(JobRuntimeDiagnostics.PIPELINE_STATUS, null)); + Assertions.assertEquals(7, firstPipeline.getInt(JobRuntimeDiagnostics.RESTORE_COUNT, -1)); + Assertions.assertEquals( + 100, firstPipeline.getInt(JobRuntimeDiagnostics.MAX_RESTORE_COUNT, -1)); + Assertions.assertEquals( + 3500L, + firstPipeline + .get(JobRuntimeDiagnostics.STATE_TIMESTAMPS) + .asObject() + .getLong(PipelineStatus.RUNNING.name(), -1L)); + // the second pipeline has no timestamps entry at all, which must not fail the payload + Assertions.assertTrue( + pipelines + .get(1) + .asObject() + .get(JobRuntimeDiagnostics.STATE_TIMESTAMPS) + .asObject() + .isEmpty()); + Assertions.assertEquals( + 9, diagnostics.getInt(JobRuntimeDiagnostics.TOTAL_PIPELINE_RESTORE_COUNT, -1)); + } + + @Test + void testJobWithoutJobMasterStillReportsStateTimestamps() { + when(stateTimestampsMap.get(JOB_ID)).thenReturn(jobTimestamps()); + when(coordinatorService.getJobMaster(JOB_ID)).thenReturn(null); + + JsonObject diagnostics = JobRuntimeDiagnostics.build(server, JOB_ID); + + Assertions.assertEquals( + 2000L, + diagnostics + .get(JobRuntimeDiagnostics.STATE_TIMESTAMPS) + .asObject() + .getLong(JobStatus.SCHEDULED.name(), -1L)); + Assertions.assertTrue(diagnostics.get(JobRuntimeDiagnostics.PIPELINES).asArray().isEmpty()); + Assertions.assertEquals( + 0, diagnostics.getInt(JobRuntimeDiagnostics.TOTAL_PIPELINE_RESTORE_COUNT, -1)); + } + + @Test + void testPipelineLookupFailureIsIsolated() { + when(stateTimestampsMap.get(JOB_ID)).thenReturn(jobTimestamps()); + doThrow(new RuntimeException("not the master node")) + .when(coordinatorService) + .getJobMaster(JOB_ID); + + JsonObject diagnostics = JobRuntimeDiagnostics.build(server, JOB_ID); + + Assertions.assertTrue(diagnostics.get(JobRuntimeDiagnostics.PIPELINES).asArray().isEmpty()); + Assertions.assertFalse( + diagnostics.get(JobRuntimeDiagnostics.STATE_TIMESTAMPS).asObject().isEmpty()); + } + + @Test + void testCleanedStateTimestampsRenderAsEmptyObject() { + when(stateTimestampsMap.get(JOB_ID)).thenReturn(null); + when(coordinatorService.getJobMaster(JOB_ID)).thenReturn(null); + + JsonObject diagnostics = JobRuntimeDiagnostics.build(server, JOB_ID); + + Assertions.assertTrue( + diagnostics.get(JobRuntimeDiagnostics.STATE_TIMESTAMPS).asObject().isEmpty()); + } + + @Test + void testShorterTimestampsArrayIsToleratedForRollingUpgrade() { + // an older member may have written an array shorter than the current enum + when(stateTimestampsMap.get(JOB_ID)).thenReturn(new Long[] {1000L}); + mockPipelines(mockSubPlan(1, null, 0, 3)); + when(stateTimestampsMap.get(new PipelineLocation(JOB_ID, 1))).thenReturn(new Long[] {null}); + + JsonObject diagnostics = JobRuntimeDiagnostics.build(server, JOB_ID); + + Assertions.assertEquals( + 1000L, + diagnostics + .get(JobRuntimeDiagnostics.STATE_TIMESTAMPS) + .asObject() + .getLong(JobStatus.values()[0].name(), -1L)); + JsonObject pipeline = + diagnostics.get(JobRuntimeDiagnostics.PIPELINES).asArray().get(0).asObject(); + Assertions.assertTrue(pipeline.get(JobRuntimeDiagnostics.PIPELINE_STATUS).isNull()); + Assertions.assertEquals( + Collections.emptyList(), + pipeline.get(JobRuntimeDiagnostics.STATE_TIMESTAMPS).asObject().names()); + } +}
