reswqa commented on code in PR #22306:
URL: https://github.com/apache/flink/pull/22306#discussion_r1159228040
##########
flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/slotmanager/DefaultSlotStatusSyncerTest.java:
##########
@@ -114,6 +114,64 @@ void testAllocateSlot() throws Exception {
assertThat(allocatedFuture).isNotCompletedExceptionally();
}
+ @Test
+ void testAllocationUpdatesIgnoredIfSlotFreed() throws Exception {
+ final FineGrainedTaskManagerTracker taskManagerTracker =
+ new FineGrainedTaskManagerTracker();
+ final CompletableFuture<
+ Tuple6<
+ SlotID,
+ JobID,
+ AllocationID,
+ ResourceProfile,
+ String,
+ ResourceManagerId>>
+ requestFuture = new CompletableFuture<>();
+ final CompletableFuture<Acknowledge> responseFuture = new
CompletableFuture<>();
+ final TestingTaskExecutorGateway taskExecutorGateway =
+ new TestingTaskExecutorGatewayBuilder()
+ .setRequestSlotFunction(
+ tuple6 -> {
+ requestFuture.complete(tuple6);
+ return responseFuture;
+ })
+ .createTestingTaskExecutorGateway();
+ final TaskExecutorConnection taskExecutorConnection =
+ new TaskExecutorConnection(ResourceID.generate(),
taskExecutorGateway);
+ taskManagerTracker.addTaskManager(
+ taskExecutorConnection, ResourceProfile.ANY,
ResourceProfile.ANY);
+ final ResourceTracker resourceTracker = new DefaultResourceTracker();
+ final JobID jobId = new JobID();
+ final SlotStatusSyncer slotStatusSyncer =
+ new DefaultSlotStatusSyncer(TASK_MANAGER_REQUEST_TIMEOUT);
+ slotStatusSyncer.initialize(
+ taskManagerTracker,
+ resourceTracker,
+ ResourceManagerId.generate(),
+ EXECUTOR_RESOURCE.getExecutor());
+
+ final CompletableFuture<Void> allocatedFuture =
+ slotStatusSyncer.allocateSlot(
+ taskExecutorConnection.getInstanceID(),
+ jobId,
+ "address",
+ ResourceProfile.ANY);
+ final AllocationID allocationId = requestFuture.get().f2;
+ assertThat(resourceTracker.getAcquiredResources(jobId))
+ .contains(ResourceRequirement.create(ResourceProfile.ANY, 1));
+ assertThat(taskManagerTracker.getAllocatedOrPendingSlot(allocationId))
+ .hasValueSatisfying(slot ->
assertThat(slot.getJobId()).isEqualTo(jobId));
+ assertThat(taskManagerTracker.getAllocatedOrPendingSlot(allocationId))
+ .hasValueSatisfying(
+ slot ->
assertThat(slot.getState()).isEqualTo(SlotState.PENDING));
Review Comment:
These two assertions should be able to be combined into one
`hasValueSatisfying `.
--
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]