sqd commented on code in PR #29136:
URL: https://github.com/apache/flink/pull/29136#discussion_r4137888974


##########
flink-runtime/src/main/java/org/apache/flink/runtime/dispatcher/Dispatcher.java:
##########
@@ -1627,15 +1627,27 @@ private JobsOverview getCompletedJobsOverview() {
 
     @Override
     public CompletableFuture<MultipleJobsDetails> 
requestMultipleJobDetails(Duration timeout) {
-        List<CompletableFuture<Optional<JobDetails>>> 
individualOptionalJobDetails =
-                queryJobMastersForInformation(

Review Comment:
   That makes sense. Moved into `queryJobMastersForInformation`.



##########
flink-runtime/src/test/java/org/apache/flink/runtime/dispatcher/DispatcherTest.java:
##########
@@ -1407,6 +1407,57 @@ public void 
testRequestMultipleJobDetails_returnsJobsOfSameStateOrderedByStartTi
                 Stream.of(jobId, 
secondJobID).sorted().collect(Collectors.toList()));
     }
 
+    /**
+     * A JobMaster that fails or times out on {@code requestJobDetails} must 
not cause its running
+     * job to be silently omitted from an otherwise successful response: 
clients (such as the
+     * Kubernetes operator) treat absence from this list as "job not found".
+     */
+    @Test
+    public void 
testRequestMultipleJobDetails_doesNotSilentlyOmitJobWhoseJobMasterQueryFails()
+            throws Exception {
+        final JobID secondJobID = new JobID();
+        JobGraph secondJobGraph = JobGraphTestUtils.streamingJobGraph();
+        secondJobGraph.setJobID(secondJobID);
+        secondJobGraph.setApplicationId(applicationId);
+        final JobManagerRunner unresponsiveJobManagerRunner =
+                TestingJobManagerRunner.newBuilder()
+                        .setJobId(secondJobID)
+                        .setJobDetailsFutureFunction(
+                                () ->
+                                        FutureUtils.completedExceptionally(
+                                                new TimeoutException(
+                                                        "JobMaster did not 
answer in time")))
+                        .build();
+        final JobManagerRunnerFactory jobManagerRunnerFactory =
+                new QueuedJobManagerRunnerFactory(
+                        
runningJobManagerRunnerWithJobStatus(JobStatus.RUNNING, jobId, 10L),
+                        unresponsiveJobManagerRunner);
+
+        DispatcherGateway dispatcherGateway =
+                createDispatcherAndStartJobs(
+                        jobManagerRunnerFactory, Arrays.asList(jobGraph, 
secondJobGraph));
+
+        final CompletableFuture<MultipleJobsDetails> multipleJobsDetailsFuture 
=
+                dispatcherGateway.requestMultipleJobDetails(TIMEOUT);
+
+        final MultipleJobsDetails multipleJobsDetails;
+        try {
+            multipleJobsDetails = multipleJobsDetailsFuture.get();
+        } catch (ExecutionException e) {
+            // Failing the whole request is acceptable, but it must fail for 
the right reason:
+            // an exception naming the job whose JobMaster could not be 
queried.
+            assertThat(e)
+                    .hasStackTraceContaining("Could not retrieve the details 
of job")
+                    .hasStackTraceContaining(secondJobID.toString());
+            return;
+        }
+
+        assertThat(multipleJobsDetails.getJobs())
+                .extracting(JobDetails::getJobId)
+                .as("a successful response must list every registered running 
job")
+                .containsExactlyInAnyOrder(jobId, secondJobID);

Review Comment:
   Done.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to