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]