MartijnVisser commented on code in PR #29301:
URL: https://github.com/apache/flink/pull/29301#discussion_r4121575269


##########
flink-runtime/src/main/java/org/apache/flink/runtime/resourcemanager/slotmanager/FineGrainedSlotManager.java:
##########
@@ -286,6 +290,16 @@ public void suspend() {
             metricsUpdateFuture = null;
         }
 
+        // cancel pending delayed checks so they don't run against a 
shutting-down executor

Review Comment:
   In the fork logs the check is still running when the executor shuts down, 
with `close()` queued behind it. Pinned to one CPU, this branch crashes with 
exit 239 as often as `master`.



##########
flink-runtime/src/main/java/org/apache/flink/runtime/resourcemanager/slotmanager/FineGrainedSlotManager.java:
##########
@@ -286,6 +290,16 @@ public void suspend() {
             metricsUpdateFuture = null;
         }
 
+        // cancel pending delayed checks so they don't run against a 
shutting-down executor
+        if (requirementsCheckScheduledFuture != null) {
+            requirementsCheckScheduledFuture.cancel(false);

Review Comment:
   Should `requirementsCheckFuture` and `declareNeededResourceFuture` be reset 
here as well? Otherwise a `start()` after `suspend()` never schedules a delayed 
check again.



##########
flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/slotmanager/FineGrainedSlotManagerTest.java:
##########
@@ -1169,4 +1169,31 @@ void testMetricsUpdate() throws Exception {
                             .isEqualTo(DEFAULT_NUM_SLOTS_PER_WORKER);
                 });
     }
+
+    @Test
+    void testCloseCancelsPendingRequirementCheck() throws Exception {
+        final ManuallyTriggeredScheduledExecutor scheduledExecutor =
+                new ManuallyTriggeredScheduledExecutor();
+        new Context() {
+            {
+                setScheduledExecutor(scheduledExecutor);
+                
slotManagerConfigurationBuilder.setRequirementCheckDelay(Duration.ofMillis(50));
+                runTest(
+                        () -> {
+                            runInMainThreadAndWait(
+                                    () ->
+                                            getSlotManager()
+                                                    
.processResourceRequirements(
+                                                            
createResourceRequirementsForSingleSlot()));
+                            // a delayed (non-periodic) check is now pending, 
not yet run
+                            
assertThat(scheduledExecutor.getActiveNonPeriodicScheduledTask())
+                                    .isNotEmpty();
+                        });
+                // flush the main-thread executor so the enqueued close() has 
cancelled the
+                // pending check before we assert, instead of racing it
+                runInMainThreadAndWait(() -> {});

Review Comment:
   `runTest` never waits for `close()`, which is also why the first CI run 
failed on this test. Calling `closeFuture.get()` there gives me 0 crashes in 
320 class runs, against 8 in 32 on `master`.



-- 
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