reswqa commented on code in PR #22196:
URL: https://github.com/apache/flink/pull/22196#discussion_r1143047050
##########
flink-runtime/src/main/java/org/apache/flink/runtime/resourcemanager/slotmanager/FineGrainedTaskManagerTracker.java:
##########
@@ -387,14 +520,187 @@ public ResourceProfile getPendingResource() {
return totalPendingResource;
}
- @Override
- public void clear() {
- slots.clear();
- taskManagerRegistrations.clear();
- totalRegisteredResource = ResourceProfile.ZERO;
- pendingTaskManagers.clear();
- totalPendingResource = ResourceProfile.ZERO;
- pendingSlotAllocationRecords.clear();
- unWantedTaskManagers.clear();
+ //
---------------------------------------------------------------------------------------------
+ // Resource allocations.
+ //
---------------------------------------------------------------------------------------------
+
+ private Set<PendingTaskManagerId> allocateTaskManagersAccordingTo(
+ List<PendingTaskManager> pendingTaskManagers) {
+ final Set<PendingTaskManagerId> failedAllocations = new HashSet<>();
+ for (PendingTaskManager pendingTaskManager : pendingTaskManagers) {
+ if (!allocateResource(pendingTaskManager)) {
+
failedAllocations.add(pendingTaskManager.getPendingTaskManagerId());
+ }
+ }
+ return failedAllocations;
+ }
+
+ private boolean allocateResource(PendingTaskManager pendingTaskManager) {
+ checkInit();
+ Preconditions.checkState(resourceAllocator.isSupported());
+ if
(isMaxTotalResourceExceededAfterAdding(pendingTaskManager.getTotalResourceProfile()))
{
+ LOG.info(
+ "Could not allocate {}. Max total resource limitation <{},
{}> is reached.",
+ pendingTaskManager,
+ maxTotalCpu,
+ maxTotalMem.toHumanReadableString());
+ return false;
+ }
+
+ addPendingTaskManager(pendingTaskManager);
+ return true;
+ }
+
+ //
---------------------------------------------------------------------------------------------
+ // Internal periodic check methods
+ //
---------------------------------------------------------------------------------------------
+
+ private void checkTaskManagerTimeouts() {
+ for (TaskManagerInfo timeoutTaskManager : getTimeOutTaskManagers()) {
+ if (waitResultConsumedBeforeRelease) {
+ releaseIdleTaskExecutorIfPossible(timeoutTaskManager);
+ } else {
+ releaseIdleTaskExecutor(timeoutTaskManager.getInstanceId());
+ }
+ }
+ }
+
+ private Collection<TaskManagerInfo> getTimeOutTaskManagers() {
+ long currentTime = System.currentTimeMillis();
+ return getRegisteredTaskManagers().stream()
+ .filter(
+ taskManager ->
+ taskManager.isIdle()
+ && currentTime -
taskManager.getIdleSince()
+ >=
taskManagerTimeout.toMilliseconds())
+ .collect(Collectors.toList());
+ }
+
+ private void releaseIdleTaskExecutorIfPossible(TaskManagerInfo
taskManagerInfo) {
+ final long idleSince = taskManagerInfo.getIdleSince();
+ taskManagerInfo
+ .getTaskExecutorConnection()
+ .getTaskExecutorGateway()
+ .canBeReleased()
+ .thenAcceptAsync(
+ canBeReleased -> {
+ boolean stillIdle = idleSince ==
taskManagerInfo.getIdleSince();
+ if (stillIdle && canBeReleased) {
+
releaseIdleTaskExecutor(taskManagerInfo.getInstanceId());
+ }
+ },
+ mainThreadExecutor);
+ }
+
+ private void releaseIdleTaskExecutor(InstanceID timedOutTaskManagerId) {
+ checkInit();
+ if (resourceAllocator.isSupported()) {
+ addUnWantedTaskManager(timedOutTaskManagerId);
+ declareNeededResourcesWithDelay();
+ }
+ }
+
+ private void addUnWantedTaskManager(InstanceID instanceId) {
+ final FineGrainedTaskManagerRegistration taskManager =
+ taskManagerRegistrations.get(instanceId);
+ if (taskManager != null) {
+ unWantedTaskManagers.put(
+ instanceId,
+ WorkerResourceSpec.fromTotalResourceProfile(
+ taskManager.getTotalResource(),
+ SlotManagerUtils.calculateDefaultNumSlots(
+ taskManager.getTotalResource(),
+
taskManager.getDefaultSlotResourceProfile())));
+ } else {
+ LOG.debug("Unwanted task manager {} does not exists.", instanceId);
+ }
+ }
+
+ @VisibleForTesting
+ public Map<InstanceID, WorkerResourceSpec> getUnWantedTaskManager() {
+ return unWantedTaskManagers;
+ }
+
+ void declareNeededResourcesWithDelay() {
+ Preconditions.checkState(resourceAllocator.isSupported());
+
+ if (declareNeededResourceDelay.toMillis() <= 0) {
+ declareNeededResources();
+ } else {
+ if (declareNeededResourceFuture == null ||
declareNeededResourceFuture.isDone()) {
+ declareNeededResourceFuture = new CompletableFuture<>();
+ scheduledExecutor.schedule(
+ () ->
+ mainThreadExecutor.execute(
+ () -> {
+ declareNeededResources();
+
Preconditions.checkNotNull(declareNeededResourceFuture)
+ .complete(null);
+ }),
+ declareNeededResourceDelay.toMillis(),
+ TimeUnit.MILLISECONDS);
+ }
+ }
+ }
+
+ /** DO NOT call this method directly. Use {@link
#declareNeededResourcesWithDelay()} instead. */
+ private void declareNeededResources() {
+ Map<InstanceID, WorkerResourceSpec> unWantedTaskManagers =
this.getUnWantedTaskManager();
+ Map<WorkerResourceSpec, Set<InstanceID>> unWantedTaskManagerBySpec =
+ unWantedTaskManagers.entrySet().stream()
+ .collect(
+ Collectors.groupingBy(
+ Map.Entry::getValue,
+ Collectors.mapping(Map.Entry::getKey,
Collectors.toSet())));
+
+ // registered TaskManagers except unwanted worker.
+ Stream<WorkerResourceSpec> registeredTaskManagerStream =
+ this.getRegisteredTaskManagers().stream()
Review Comment:
Maybe some `this` pointer for method calls in
`FineGrainedTaskManagerTracker` should not be needed.
##########
flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/slotmanager/FineGrainedTaskManagerTrackerTest.java:
##########
@@ -319,14 +435,118 @@ public void testGetStatistics() {
TASK_EXECUTOR_CONNECTION.getInstanceID(),
defaultSlotResource,
SlotState.ALLOCATED);
- taskManagerTracker.addPendingTaskManager(
- new PendingTaskManager(ResourceProfile.fromResources(4, 200),
1));
- assertThat(taskManagerTracker.getFreeResource(),
is(ResourceProfile.fromResources(6, 700)));
- assertThat(taskManagerTracker.getRegisteredResource(),
is(totalResource));
- assertThat(taskManagerTracker.getNumberRegisteredSlots(), is(10));
- assertThat(taskManagerTracker.getNumberFreeSlots(), is(8));
+ PendingTaskManager pendingTaskManager =
+ new PendingTaskManager(ResourceProfile.fromResources(4, 200),
1);
+ taskManagerTracker.allocateTaskManagersAccordingTo(
+ new ResourceAllocationResult.Builder()
+ .addPendingTaskManagerAllocate(pendingTaskManager)
+ .addAllocationOnPendingResource(
+ jobId,
+ pendingTaskManager.getPendingTaskManagerId(),
+ ResourceProfile.fromResources(4, 200))
+ .build());
+
+ assertThat(taskManagerTracker.getFreeResource())
+ .isEqualTo(ResourceProfile.fromResources(6, 700));
+
assertThat(taskManagerTracker.getRegisteredResource()).isEqualTo(totalResource);
+
assertThat(taskManagerTracker.getNumberRegisteredSlots()).isEqualTo(10);
+ assertThat(taskManagerTracker.getNumberFreeSlots()).isEqualTo(8);
+ assertThat(taskManagerTracker.getPendingResource())
+ .isEqualTo(ResourceProfile.fromResources(4, 200));
+ }
+
+ @Test
+ void testTimeoutForUnusedTaskManager() throws Exception {
+ final Time taskManagerTimeout = Time.milliseconds(50L);
+
+ final CompletableFuture<InstanceID> releaseResourceFuture = new
CompletableFuture<>();
+
+ ResourceAllocator resourceAllocator =
+ new TestingResourceAllocatorBuilder()
+ .setDeclareResourceNeededConsumer(
+ (resourceDeclarations) -> {
+
assertThat(resourceDeclarations.size()).isEqualTo(1);
+ ResourceDeclaration resourceDeclaration =
+
resourceDeclarations.iterator().next();
+
assertThat(resourceDeclaration.getNumNeeded()).isEqualTo(0);
+
assertThat(resourceDeclaration.getUnwantedWorkers().size())
+ .isEqualTo(1);
+
+ releaseResourceFuture.complete(
+ resourceDeclaration
+ .getUnwantedWorkers()
+ .iterator()
+ .next());
+ })
+ .build();
+
+ FineGrainedTaskManagerTracker taskManagerTracker =
+ new FineGrainedTaskManagerTrackerBuilder(
+ new ScheduledExecutorServiceAdapter(
+ EXECUTOR_RESOURCE.getExecutor()))
+ .setTaskManagerTimeout(taskManagerTimeout)
+ .build();
+ taskManagerTracker.initialize(resourceAllocator,
EXECUTOR_RESOURCE.getExecutor());
+
+ taskManagerTracker.registerTaskManager(
+ TASK_EXECUTOR_CONNECTION,
+ DEFAULT_TOTAL_RESOURCE_PROFILE,
+ DEFAULT_SLOT_RESOURCE_PROFILE,
+ null);
+
+ AllocationID allocationID = new AllocationID();
+ JobID jobID = new JobID();
+ taskManagerTracker.notifySlotStatus(
+ allocationID,
+ jobID,
+ TASK_EXECUTOR_CONNECTION.getInstanceID(),
+ ResourceProfile.fromResources(3, 200),
+ SlotState.ALLOCATED);
+
+ assertThat(
+ taskManagerTracker.getRegisteredTaskManager(
+ TASK_EXECUTOR_CONNECTION.getInstanceID()))
+ .hasValueSatisfying(
+ taskManagerInfo ->
+ assertThat(taskManagerInfo.getIdleSince())
+ .isEqualTo(Long.MAX_VALUE));
+
+ taskManagerTracker.notifySlotStatus(
+ allocationID,
+ jobID,
+ TASK_EXECUTOR_CONNECTION.getInstanceID(),
+ ResourceProfile.fromResources(3, 200),
+ SlotState.FREE);
+
assertThat(
- taskManagerTracker.getPendingResource(),
is(ResourceProfile.fromResources(4, 200)));
+ taskManagerTracker.getRegisteredTaskManager(
+ TASK_EXECUTOR_CONNECTION.getInstanceID()))
+ .hasValueSatisfying(
+ taskManagerInfo ->
+ assertThat(taskManagerInfo.getIdleSince())
+ .isNotEqualTo(Long.MAX_VALUE));
+
+ assertThat(releaseResourceFuture.get(1, TimeUnit.SECONDS))
Review Comment:
We'd better use
`FlinkCompletableFutureAssert.assertThatFuture(releaseResourceFuture).eventuallySucceeds()`
to get rid of the notorious timeout here.
##########
flink-runtime/src/main/java/org/apache/flink/runtime/resourcemanager/slotmanager/FineGrainedTaskManagerTracker.java:
##########
@@ -77,30 +130,129 @@ public FineGrainedTaskManagerTracker() {
}
@Override
- public void replaceAllPendingAllocations(
+ public void initialize(ResourceAllocator resourceAllocator, Executor
mainThreadExecutor) {
+ this.resourceAllocator = resourceAllocator;
+ this.mainThreadExecutor = mainThreadExecutor;
+ this.started = true;
+
+ taskManagerTimeoutsCheck =
+ scheduledExecutor.scheduleWithFixedDelay(
+ () ->
mainThreadExecutor.execute(this::checkTaskManagerTimeouts),
+ 0L,
+ taskManagerTimeout.toMilliseconds(),
+ TimeUnit.MILLISECONDS);
+ }
+
+ @Override
+ public void close() {
+ // stop the timeout checks for the TaskManagers
+ if (taskManagerTimeoutsCheck != null) {
+ taskManagerTimeoutsCheck.cancel(false);
+ taskManagerTimeoutsCheck = null;
+ }
+
+ slots.clear();
+ taskManagerRegistrations.clear();
+ totalRegisteredResource = ResourceProfile.ZERO;
+ pendingTaskManagers.clear();
+ totalPendingResource = ResourceProfile.ZERO;
+ pendingSlotAllocationRecords.clear();
+ unWantedTaskManagers.clear();
+ mainThreadExecutor = null;
+ resourceAllocator = null;
+ started = false;
+ }
+
+ private void replaceAllPendingAllocations(
Map<PendingTaskManagerId, Map<JobID, ResourceCounter>>
pendingSlotAllocations) {
Preconditions.checkNotNull(pendingSlotAllocations);
+ Preconditions.checkState(resourceAllocator.isSupported());
LOG.trace("Record the pending allocations {}.",
pendingSlotAllocations);
pendingSlotAllocationRecords.clear();
pendingSlotAllocationRecords.putAll(pendingSlotAllocations);
removeUnusedPendingTaskManagers();
Review Comment:
How about moving `declareNeededResourcesWithDelay` to the end of this method
to align with `clearAllPendingAllocations`.
##########
flink-runtime/src/main/java/org/apache/flink/runtime/resourcemanager/slotmanager/FineGrainedTaskManagerTracker.java:
##########
@@ -387,14 +520,187 @@ public ResourceProfile getPendingResource() {
return totalPendingResource;
}
- @Override
- public void clear() {
- slots.clear();
- taskManagerRegistrations.clear();
- totalRegisteredResource = ResourceProfile.ZERO;
- pendingTaskManagers.clear();
- totalPendingResource = ResourceProfile.ZERO;
- pendingSlotAllocationRecords.clear();
- unWantedTaskManagers.clear();
+ //
---------------------------------------------------------------------------------------------
+ // Resource allocations.
+ //
---------------------------------------------------------------------------------------------
+
+ private Set<PendingTaskManagerId> allocateTaskManagersAccordingTo(
+ List<PendingTaskManager> pendingTaskManagers) {
+ final Set<PendingTaskManagerId> failedAllocations = new HashSet<>();
+ for (PendingTaskManager pendingTaskManager : pendingTaskManagers) {
+ if (!allocateResource(pendingTaskManager)) {
+
failedAllocations.add(pendingTaskManager.getPendingTaskManagerId());
+ }
+ }
+ return failedAllocations;
+ }
+
+ private boolean allocateResource(PendingTaskManager pendingTaskManager) {
+ checkInit();
+ Preconditions.checkState(resourceAllocator.isSupported());
+ if
(isMaxTotalResourceExceededAfterAdding(pendingTaskManager.getTotalResourceProfile()))
{
+ LOG.info(
+ "Could not allocate {}. Max total resource limitation <{},
{}> is reached.",
+ pendingTaskManager,
+ maxTotalCpu,
+ maxTotalMem.toHumanReadableString());
+ return false;
+ }
+
+ addPendingTaskManager(pendingTaskManager);
+ return true;
+ }
+
+ //
---------------------------------------------------------------------------------------------
+ // Internal periodic check methods
+ //
---------------------------------------------------------------------------------------------
+
+ private void checkTaskManagerTimeouts() {
+ for (TaskManagerInfo timeoutTaskManager : getTimeOutTaskManagers()) {
+ if (waitResultConsumedBeforeRelease) {
+ releaseIdleTaskExecutorIfPossible(timeoutTaskManager);
+ } else {
+ releaseIdleTaskExecutor(timeoutTaskManager.getInstanceId());
+ }
+ }
+ }
+
+ private Collection<TaskManagerInfo> getTimeOutTaskManagers() {
+ long currentTime = System.currentTimeMillis();
+ return getRegisteredTaskManagers().stream()
+ .filter(
+ taskManager ->
+ taskManager.isIdle()
+ && currentTime -
taskManager.getIdleSince()
+ >=
taskManagerTimeout.toMilliseconds())
+ .collect(Collectors.toList());
+ }
+
+ private void releaseIdleTaskExecutorIfPossible(TaskManagerInfo
taskManagerInfo) {
+ final long idleSince = taskManagerInfo.getIdleSince();
+ taskManagerInfo
+ .getTaskExecutorConnection()
+ .getTaskExecutorGateway()
+ .canBeReleased()
+ .thenAcceptAsync(
+ canBeReleased -> {
+ boolean stillIdle = idleSince ==
taskManagerInfo.getIdleSince();
+ if (stillIdle && canBeReleased) {
+
releaseIdleTaskExecutor(taskManagerInfo.getInstanceId());
+ }
+ },
+ mainThreadExecutor);
+ }
+
+ private void releaseIdleTaskExecutor(InstanceID timedOutTaskManagerId) {
+ checkInit();
+ if (resourceAllocator.isSupported()) {
+ addUnWantedTaskManager(timedOutTaskManagerId);
+ declareNeededResourcesWithDelay();
+ }
+ }
+
+ private void addUnWantedTaskManager(InstanceID instanceId) {
+ final FineGrainedTaskManagerRegistration taskManager =
+ taskManagerRegistrations.get(instanceId);
+ if (taskManager != null) {
+ unWantedTaskManagers.put(
+ instanceId,
+ WorkerResourceSpec.fromTotalResourceProfile(
+ taskManager.getTotalResource(),
+ SlotManagerUtils.calculateDefaultNumSlots(
+ taskManager.getTotalResource(),
+
taskManager.getDefaultSlotResourceProfile())));
+ } else {
+ LOG.debug("Unwanted task manager {} does not exists.", instanceId);
+ }
+ }
+
+ @VisibleForTesting
+ public Map<InstanceID, WorkerResourceSpec> getUnWantedTaskManager() {
+ return unWantedTaskManagers;
+ }
+
+ void declareNeededResourcesWithDelay() {
+ Preconditions.checkState(resourceAllocator.isSupported());
+
+ if (declareNeededResourceDelay.toMillis() <= 0) {
+ declareNeededResources();
+ } else {
+ if (declareNeededResourceFuture == null ||
declareNeededResourceFuture.isDone()) {
+ declareNeededResourceFuture = new CompletableFuture<>();
+ scheduledExecutor.schedule(
+ () ->
+ mainThreadExecutor.execute(
+ () -> {
+ declareNeededResources();
+
Preconditions.checkNotNull(declareNeededResourceFuture)
+ .complete(null);
+ }),
+ declareNeededResourceDelay.toMillis(),
+ TimeUnit.MILLISECONDS);
+ }
+ }
+ }
+
+ /** DO NOT call this method directly. Use {@link
#declareNeededResourcesWithDelay()} instead. */
+ private void declareNeededResources() {
Review Comment:
Maybe we also need to `checkInit` in this method as `resourceAllocator` is
used here.
Because this method was submitted to executor, even though `checkInit` at
the caller does not guarantee that there is no problem.
##########
flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/slotmanager/FineGrainedTaskManagerTrackerTest.java:
##########
@@ -319,14 +435,118 @@ public void testGetStatistics() {
TASK_EXECUTOR_CONNECTION.getInstanceID(),
defaultSlotResource,
SlotState.ALLOCATED);
- taskManagerTracker.addPendingTaskManager(
- new PendingTaskManager(ResourceProfile.fromResources(4, 200),
1));
- assertThat(taskManagerTracker.getFreeResource(),
is(ResourceProfile.fromResources(6, 700)));
- assertThat(taskManagerTracker.getRegisteredResource(),
is(totalResource));
- assertThat(taskManagerTracker.getNumberRegisteredSlots(), is(10));
- assertThat(taskManagerTracker.getNumberFreeSlots(), is(8));
+ PendingTaskManager pendingTaskManager =
+ new PendingTaskManager(ResourceProfile.fromResources(4, 200),
1);
+ taskManagerTracker.allocateTaskManagersAccordingTo(
+ new ResourceAllocationResult.Builder()
+ .addPendingTaskManagerAllocate(pendingTaskManager)
+ .addAllocationOnPendingResource(
+ jobId,
+ pendingTaskManager.getPendingTaskManagerId(),
+ ResourceProfile.fromResources(4, 200))
+ .build());
+
+ assertThat(taskManagerTracker.getFreeResource())
+ .isEqualTo(ResourceProfile.fromResources(6, 700));
+
assertThat(taskManagerTracker.getRegisteredResource()).isEqualTo(totalResource);
+
assertThat(taskManagerTracker.getNumberRegisteredSlots()).isEqualTo(10);
+ assertThat(taskManagerTracker.getNumberFreeSlots()).isEqualTo(8);
+ assertThat(taskManagerTracker.getPendingResource())
+ .isEqualTo(ResourceProfile.fromResources(4, 200));
+ }
+
+ @Test
+ void testTimeoutForUnusedTaskManager() throws Exception {
Review Comment:
It seems that `testTimeoutForUnusedTaskManager` in
`FineGrainedSlotManagerTest` is not needed.
##########
flink-runtime/src/main/java/org/apache/flink/runtime/resourcemanager/slotmanager/FineGrainedTaskManagerTracker.java:
##########
@@ -387,14 +520,187 @@ public ResourceProfile getPendingResource() {
return totalPendingResource;
}
- @Override
- public void clear() {
- slots.clear();
- taskManagerRegistrations.clear();
- totalRegisteredResource = ResourceProfile.ZERO;
- pendingTaskManagers.clear();
- totalPendingResource = ResourceProfile.ZERO;
- pendingSlotAllocationRecords.clear();
- unWantedTaskManagers.clear();
+ //
---------------------------------------------------------------------------------------------
+ // Resource allocations.
+ //
---------------------------------------------------------------------------------------------
+
+ private Set<PendingTaskManagerId> allocateTaskManagersAccordingTo(
+ List<PendingTaskManager> pendingTaskManagers) {
+ final Set<PendingTaskManagerId> failedAllocations = new HashSet<>();
+ for (PendingTaskManager pendingTaskManager : pendingTaskManagers) {
+ if (!allocateResource(pendingTaskManager)) {
+
failedAllocations.add(pendingTaskManager.getPendingTaskManagerId());
+ }
+ }
+ return failedAllocations;
+ }
+
+ private boolean allocateResource(PendingTaskManager pendingTaskManager) {
+ checkInit();
+ Preconditions.checkState(resourceAllocator.isSupported());
+ if
(isMaxTotalResourceExceededAfterAdding(pendingTaskManager.getTotalResourceProfile()))
{
+ LOG.info(
+ "Could not allocate {}. Max total resource limitation <{},
{}> is reached.",
+ pendingTaskManager,
+ maxTotalCpu,
+ maxTotalMem.toHumanReadableString());
+ return false;
+ }
+
+ addPendingTaskManager(pendingTaskManager);
+ return true;
+ }
+
+ //
---------------------------------------------------------------------------------------------
+ // Internal periodic check methods
+ //
---------------------------------------------------------------------------------------------
+
+ private void checkTaskManagerTimeouts() {
+ for (TaskManagerInfo timeoutTaskManager : getTimeOutTaskManagers()) {
+ if (waitResultConsumedBeforeRelease) {
+ releaseIdleTaskExecutorIfPossible(timeoutTaskManager);
+ } else {
+ releaseIdleTaskExecutor(timeoutTaskManager.getInstanceId());
+ }
+ }
+ }
+
+ private Collection<TaskManagerInfo> getTimeOutTaskManagers() {
+ long currentTime = System.currentTimeMillis();
+ return getRegisteredTaskManagers().stream()
+ .filter(
+ taskManager ->
+ taskManager.isIdle()
+ && currentTime -
taskManager.getIdleSince()
+ >=
taskManagerTimeout.toMilliseconds())
+ .collect(Collectors.toList());
+ }
+
+ private void releaseIdleTaskExecutorIfPossible(TaskManagerInfo
taskManagerInfo) {
+ final long idleSince = taskManagerInfo.getIdleSince();
+ taskManagerInfo
+ .getTaskExecutorConnection()
+ .getTaskExecutorGateway()
+ .canBeReleased()
+ .thenAcceptAsync(
+ canBeReleased -> {
+ boolean stillIdle = idleSince ==
taskManagerInfo.getIdleSince();
+ if (stillIdle && canBeReleased) {
+
releaseIdleTaskExecutor(taskManagerInfo.getInstanceId());
+ }
+ },
+ mainThreadExecutor);
+ }
+
+ private void releaseIdleTaskExecutor(InstanceID timedOutTaskManagerId) {
+ checkInit();
+ if (resourceAllocator.isSupported()) {
+ addUnWantedTaskManager(timedOutTaskManagerId);
+ declareNeededResourcesWithDelay();
+ }
+ }
+
+ private void addUnWantedTaskManager(InstanceID instanceId) {
+ final FineGrainedTaskManagerRegistration taskManager =
+ taskManagerRegistrations.get(instanceId);
+ if (taskManager != null) {
+ unWantedTaskManagers.put(
+ instanceId,
+ WorkerResourceSpec.fromTotalResourceProfile(
+ taskManager.getTotalResource(),
+ SlotManagerUtils.calculateDefaultNumSlots(
+ taskManager.getTotalResource(),
+
taskManager.getDefaultSlotResourceProfile())));
+ } else {
+ LOG.debug("Unwanted task manager {} does not exists.", instanceId);
+ }
+ }
+
+ @VisibleForTesting
+ public Map<InstanceID, WorkerResourceSpec> getUnWantedTaskManager() {
Review Comment:
```suggestion
private Map<InstanceID, WorkerResourceSpec> getUnWantedTaskManager() {
```
I believe this can be private if we do not depend on this method in test.
##########
flink-runtime/src/main/java/org/apache/flink/runtime/resourcemanager/slotmanager/ResourceDeclaration.java:
##########
@@ -51,7 +51,7 @@ public int getNumNeeded() {
return numNeeded;
}
- public Collection<InstanceID> getUnwantedWorkers() {
+ public Set<InstanceID> getUnwantedWorkers() {
Review Comment:
What is the purpose of this change? It is best to add some description in
the commit message.
##########
flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/slotmanager/FineGrainedTaskManagerTrackerTest.java:
##########
@@ -319,14 +435,118 @@ public void testGetStatistics() {
TASK_EXECUTOR_CONNECTION.getInstanceID(),
defaultSlotResource,
SlotState.ALLOCATED);
- taskManagerTracker.addPendingTaskManager(
- new PendingTaskManager(ResourceProfile.fromResources(4, 200),
1));
- assertThat(taskManagerTracker.getFreeResource(),
is(ResourceProfile.fromResources(6, 700)));
- assertThat(taskManagerTracker.getRegisteredResource(),
is(totalResource));
- assertThat(taskManagerTracker.getNumberRegisteredSlots(), is(10));
- assertThat(taskManagerTracker.getNumberFreeSlots(), is(8));
+ PendingTaskManager pendingTaskManager =
+ new PendingTaskManager(ResourceProfile.fromResources(4, 200),
1);
+ taskManagerTracker.allocateTaskManagersAccordingTo(
+ new ResourceAllocationResult.Builder()
+ .addPendingTaskManagerAllocate(pendingTaskManager)
+ .addAllocationOnPendingResource(
+ jobId,
+ pendingTaskManager.getPendingTaskManagerId(),
+ ResourceProfile.fromResources(4, 200))
+ .build());
+
+ assertThat(taskManagerTracker.getFreeResource())
+ .isEqualTo(ResourceProfile.fromResources(6, 700));
+
assertThat(taskManagerTracker.getRegisteredResource()).isEqualTo(totalResource);
+
assertThat(taskManagerTracker.getNumberRegisteredSlots()).isEqualTo(10);
+ assertThat(taskManagerTracker.getNumberFreeSlots()).isEqualTo(8);
+ assertThat(taskManagerTracker.getPendingResource())
+ .isEqualTo(ResourceProfile.fromResources(4, 200));
+ }
+
+ @Test
+ void testTimeoutForUnusedTaskManager() throws Exception {
+ final Time taskManagerTimeout = Time.milliseconds(50L);
+
+ final CompletableFuture<InstanceID> releaseResourceFuture = new
CompletableFuture<>();
+
+ ResourceAllocator resourceAllocator =
+ new TestingResourceAllocatorBuilder()
+ .setDeclareResourceNeededConsumer(
+ (resourceDeclarations) -> {
+
assertThat(resourceDeclarations.size()).isEqualTo(1);
+ ResourceDeclaration resourceDeclaration =
+
resourceDeclarations.iterator().next();
+
assertThat(resourceDeclaration.getNumNeeded()).isEqualTo(0);
+
assertThat(resourceDeclaration.getUnwantedWorkers().size())
+ .isEqualTo(1);
+
+ releaseResourceFuture.complete(
+ resourceDeclaration
+ .getUnwantedWorkers()
+ .iterator()
+ .next());
+ })
+ .build();
+
+ FineGrainedTaskManagerTracker taskManagerTracker =
+ new FineGrainedTaskManagerTrackerBuilder(
+ new ScheduledExecutorServiceAdapter(
+ EXECUTOR_RESOURCE.getExecutor()))
+ .setTaskManagerTimeout(taskManagerTimeout)
+ .build();
+ taskManagerTracker.initialize(resourceAllocator,
EXECUTOR_RESOURCE.getExecutor());
+
+ taskManagerTracker.registerTaskManager(
+ TASK_EXECUTOR_CONNECTION,
+ DEFAULT_TOTAL_RESOURCE_PROFILE,
+ DEFAULT_SLOT_RESOURCE_PROFILE,
+ null);
+
+ AllocationID allocationID = new AllocationID();
+ JobID jobID = new JobID();
+ taskManagerTracker.notifySlotStatus(
+ allocationID,
+ jobID,
+ TASK_EXECUTOR_CONNECTION.getInstanceID(),
+ ResourceProfile.fromResources(3, 200),
+ SlotState.ALLOCATED);
+
+ assertThat(
+ taskManagerTracker.getRegisteredTaskManager(
+ TASK_EXECUTOR_CONNECTION.getInstanceID()))
+ .hasValueSatisfying(
+ taskManagerInfo ->
+ assertThat(taskManagerInfo.getIdleSince())
+ .isEqualTo(Long.MAX_VALUE));
+
+ taskManagerTracker.notifySlotStatus(
+ allocationID,
+ jobID,
+ TASK_EXECUTOR_CONNECTION.getInstanceID(),
+ ResourceProfile.fromResources(3, 200),
+ SlotState.FREE);
+
assertThat(
- taskManagerTracker.getPendingResource(),
is(ResourceProfile.fromResources(4, 200)));
+ taskManagerTracker.getRegisteredTaskManager(
+ TASK_EXECUTOR_CONNECTION.getInstanceID()))
+ .hasValueSatisfying(
+ taskManagerInfo ->
+ assertThat(taskManagerInfo.getIdleSince())
+ .isNotEqualTo(Long.MAX_VALUE));
+
+ assertThat(releaseResourceFuture.get(1, TimeUnit.SECONDS))
+ .isEqualTo(TASK_EXECUTOR_CONNECTION.getInstanceID());
+ assertThat(taskManagerTracker.getUnWantedTaskManager())
Review Comment:
IMO, We should try to keep our test is a black box, and if there are other
approach that can verify the expected behavior, try not to introduce
`VisibleForTesting`.
As this case, the following assertion is enough to check the
`unwantedTaskManager` is expected, right?
```
assertThat(releaseResourceFuture.get(1, TimeUnit.SECONDS))
.isEqualTo(TASK_EXECUTOR_CONNECTION.getInstanceID());
```
--
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]