Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -249,16 +249,24 @@ protected JobTopology getJobTopology(JobDetailsInfo jobDetailsInfo) {

var metrics = new HashMap<JobVertexID, IOMetrics>();
var finished = new HashSet<JobVertexID>();
jobDetailsInfo
.getJobVertexInfos()
.forEach(
d -> {
if (d.getExecutionState() == ExecutionState.FINISHED) {
finished.add(d.getJobVertexID());
}
metrics.put(
d.getJobVertexID(), IOMetrics.from(d.getJobVertexMetrics()));
});
for (var d : jobDetailsInfo.getJobVertexInfos()) {
if (d.getExecutionState() == ExecutionState.FINISHED) {
finished.add(d.getJobVertexID());
}
var ioMetricsInfo = d.getJobVertexMetrics();
if (!ioMetricsInfo.isRecordsReadComplete()
|| !ioMetricsInfo.isRecordsWrittenComplete()) {
throw new NotReadyException(
String.format(
"Vertex %s has incomplete IO metrics (read-records-complete=%s,"
+ " write-records-complete=%s). Skipping metric"
+ " collection until metrics are complete.",
d.getJobVertexID(),
ioMetricsInfo.isRecordsReadComplete(),
ioMetricsInfo.isRecordsWrittenComplete()));
}
metrics.put(d.getJobVertexID(), IOMetrics.from(ioMetricsInfo));
}

return JobTopology.fromJsonPlan(
json, slotSharingGroupIdMap, maxParallelismMap, metrics, finished);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -354,6 +354,136 @@ public void testJobTopologyParsingThrowsNotReadyException() throws Exception {
NotReadyException.class, () -> metricsCollector.getJobTopology(jobDetailsInfo));
}

@Test
public void testIncompleteIoMetricsThrowsNotReadyException() throws Exception {
String s =
"{\n"
+ " \"jid\": \"bb8f15efbb37f2ce519f55cdc0e049bf\",\n"
+ " \"name\": \"State machine job\",\n"
+ " \"isStoppable\": false,\n"
+ " \"state\": \"RUNNING\",\n"
+ " \"start-time\": 1707893512027,\n"
+ " \"end-time\": -1,\n"
+ " \"duration\": 214716,\n"
+ " \"maxParallelism\": -1,\n"
+ " \"now\": 1707893726743,\n"
+ " \"timestamps\": {\n"
+ " \"SUSPENDED\": 0,\n"
+ " \"CREATED\": 1707893512139,\n"
+ " \"FAILING\": 0,\n"
+ " \"FAILED\": 0,\n"
+ " \"INITIALIZING\": 1707893512027,\n"
+ " \"RECONCILING\": 0,\n"
+ " \"RUNNING\": 1707893512217,\n"
+ " \"RESTARTING\": 0,\n"
+ " \"CANCELLING\": 0,\n"
+ " \"FINISHED\": 0,\n"
+ " \"CANCELED\": 0\n"
+ " },\n"
+ " \"vertices\": [\n"
+ " {\n"
+ " \"id\": \"bc764cd8ddf7a0cff126f51c16239658\",\n"
+ " \"name\": \"Source: Custom Source\",\n"
+ " \"maxParallelism\": 128,\n"
+ " \"parallelism\": 2,\n"
+ " \"status\": \"RUNNING\",\n"
+ " \"start-time\": 1707893517277,\n"
+ " \"end-time\": -1,\n"
+ " \"duration\": 209466,\n"
+ " \"tasks\": {\n"
+ " \"DEPLOYING\": 0,\n"
+ " \"INITIALIZING\": 0,\n"
+ " \"SCHEDULED\": 0,\n"
+ " \"CANCELING\": 0,\n"
+ " \"CANCELED\": 0,\n"
+ " \"RECONCILING\": 0,\n"
+ " \"RUNNING\": 2,\n"
+ " \"FAILED\": 0,\n"
+ " \"CREATED\": 0,\n"
+ " \"FINISHED\": 0\n"
+ " },\n"
+ " \"metrics\": {\n"
+ " \"read-bytes\": 0,\n"
+ " \"read-bytes-complete\": true,\n"
+ " \"write-bytes\": 0,\n"
+ " \"write-bytes-complete\": false,\n"
+ " \"read-records\": 0,\n"
+ " \"read-records-complete\": true,\n"
+ " \"write-records\": 0,\n"
+ " \"write-records-complete\": false,\n"
+ " \"accumulated-backpressured-time\": 0,\n"
+ " \"accumulated-idle-time\": 0,\n"
+ " \"accumulated-busy-time\": 0\n"
+ " }\n"
+ " },\n"
+ " {\n"
+ " \"id\": \"20ba6b65f97481d5570070de90e4e791\",\n"
+ " \"name\": \"Flat Map -> Sink: Print to Std. Out\",\n"
+ " \"maxParallelism\": 128,\n"
+ " \"parallelism\": 2,\n"
+ " \"status\": \"RUNNING\",\n"
+ " \"start-time\": 1707893517280,\n"
+ " \"end-time\": -1,\n"
+ " \"duration\": 209463,\n"
+ " \"tasks\": {\n"
+ " \"DEPLOYING\": 0,\n"
+ " \"INITIALIZING\": 0,\n"
+ " \"SCHEDULED\": 0,\n"
+ " \"CANCELING\": 0,\n"
+ " \"CANCELED\": 0,\n"
+ " \"RECONCILING\": 0,\n"
+ " \"RUNNING\": 2,\n"
+ " \"FAILED\": 0,\n"
+ " \"CREATED\": 0,\n"
+ " \"FINISHED\": 0\n"
+ " },\n"
+ " \"metrics\": {\n"
+ " \"read-bytes\": 0,\n"
+ " \"read-bytes-complete\": false,\n"
+ " \"write-bytes\": 0,\n"
+ " \"write-bytes-complete\": true,\n"
+ " \"read-records\": 0,\n"
+ " \"read-records-complete\": false,\n"
+ " \"write-records\": 0,\n"
+ " \"write-records-complete\": true,\n"
+ " \"accumulated-backpressured-time\": 0,\n"
+ " \"accumulated-idle-time\": 0,\n"
+ " \"accumulated-busy-time\": 1\n"
+ " }\n"
+ " }\n"
+ " ],\n"
+ " \"status-counts\": {\n"
+ " \"DEPLOYING\": 0,\n"
+ " \"INITIALIZING\": 0,\n"
+ " \"SCHEDULED\": 0,\n"
+ " \"CANCELING\": 0,\n"
+ " \"CANCELED\": 0,\n"
+ " \"RECONCILING\": 0,\n"
+ " \"RUNNING\": 2,\n"
+ " \"FAILED\": 0,\n"
+ " \"CREATED\": 0,\n"
+ " \"FINISHED\": 0\n"
+ " },\n"
+ " \"plan\": {\n"
+ " \"jid\": \"bb8f15efbb37f2ce519f55cdc0e049bf\",\n"
+ " \"name\": \"State machine job\",\n"
+ " \"type\": \"STREAMING\",\n"
+ " \"nodes\": [\n"
+ " ]\n"
+ " }\n"
+ "}";
JobDetailsInfo jobDetailsInfo = new ObjectMapper().readValue(s, JobDetailsInfo.class);

var metricsCollector = new RestApiMetricsCollector<>();
var ex =
assertThrows(
NotReadyException.class,
() -> metricsCollector.getJobTopology(jobDetailsInfo));
assertTrue(
ex.getMessage().contains("incomplete IO metrics"),
"Exception message should indicate incomplete IO metrics, got: " + ex.getMessage());
}

@Test
public void testJobTopologyParsingFromJobDetailsWithSlotSharingGroup() throws Exception {
String s =
Expand Down