Skip to content

Commit 05a7878

Browse files
committed
[FLINK-38719][runtime] Display the number of tasks on the TaskManagers page.
1 parent eb1e096 commit 05a7878

21 files changed

Lines changed: 109 additions & 2 deletions

File tree

docs/layouts/shortcodes/generated/rest_v1_dispatcher.html

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5240,6 +5240,9 @@
52405240
"freeSlots" : {
52415241
"type" : "integer"
52425242
},
5243+
"numberOfTasks" : {
5244+
"type" : "integer"
5245+
},
52435246
"hardware" : {
52445247
"type" : "object",
52455248
"id" : "urn:jsonschema:org:apache:flink:runtime:instance:HardwareDescription",
@@ -5464,6 +5467,9 @@
54645467
"freeSlots" : {
54655468
"type" : "integer"
54665469
},
5470+
"numberOfTasks" : {
5471+
"type" : "integer"
5472+
},
54675473
"hardware" : {
54685474
"type" : "object",
54695475
"id" : "urn:jsonschema:org:apache:flink:runtime:instance:HardwareDescription",

flink-runtime-web/src/test/resources/rest_api_v1.snapshot

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3931,6 +3931,9 @@
39313931
"freeSlots" : {
39323932
"type" : "integer"
39333933
},
3934+
"numberOfTasks" : {
3935+
"type" : "integer"
3936+
},
39343937
"totalResource" : {
39353938
"type" : "object",
39363939
"id" : "urn:jsonschema:org:apache:flink:runtime:rest:messages:ResourceProfileInfo",
@@ -4093,6 +4096,9 @@
40934096
"freeSlots" : {
40944097
"type" : "integer"
40954098
},
4099+
"numberOfTasks" : {
4100+
"type" : "integer"
4101+
},
40964102
"totalResource" : {
40974103
"type" : "object",
40984104
"id" : "urn:jsonschema:org:apache:flink:runtime:rest:messages:ResourceProfileInfo",

flink-runtime-web/web-dashboard/src/app/interfaces/task-manager.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ export interface TaskManagerDetail {
2727
timeSinceLastHeartbeat: number;
2828
slotsNumber: number;
2929
freeSlots: number;
30+
numberOfTasks: number;
3031
hardware: Hardware;
3132
metrics: Metrics;
3233
memoryConfiguration: MemoryConfiguration;
@@ -68,6 +69,7 @@ export interface TaskManagersItem {
6869
timeSinceLastHeartbeat: number;
6970
slotsNumber: number;
7071
freeSlots: number;
72+
numberOfTasks: number;
7173
hardware: Hardware;
7274
blocked?: boolean;
7375
}

flink-runtime-web/web-dashboard/src/app/pages/task-manager/list/task-manager-list.component.html

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@
3232
<th [nzSortFn]="sortHeartBeatFn" [nzWidth]="'160px'">Last Heartbeat</th>
3333
<th [nzSortFn]="sortSlotsNumberFn" [nzWidth]="'90px'">All Slots</th>
3434
<th [nzSortFn]="sortFreeSlotsFn" [nzWidth]="'100px'">Free Slots</th>
35+
<th [nzSortFn]="sortTasksFn" [nzWidth]="'100px'">Tasks</th>
3536
<th [nzSortFn]="sortCpuCoresFn" [nzWidth]="'110px'">CPU Cores</th>
3637
<th [nzSortFn]="sortPhysicalMemoryFn" [nzWidth]="'120px'">Physical MEM</th>
3738
<th [nzSortFn]="sortFreeMemoryFn" [nzWidth]="'130px'">JVM Heap Size</th>
@@ -53,6 +54,7 @@
5354
<td>{{ manager.timeSinceLastHeartbeat | date: 'yyyy-MM-dd HH:mm:ss' }}</td>
5455
<td>{{ manager.slotsNumber }}</td>
5556
<td>{{ manager.freeSlots }}</td>
57+
<td>{{ manager.numberOfTasks }}</td>
5658
<td>{{ manager.hardware.cpuCores }}</td>
5759
<td [attr.title]="manager.hardware.physicalMemory + ' bytes'">
5860
{{ manager.hardware.physicalMemory | humanizeBytes }}

flink-runtime-web/web-dashboard/src/app/pages/task-manager/list/task-manager-list.component.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,7 @@ export class TaskManagerListComponent implements OnInit, OnDestroy {
4949
public readonly sortHeartBeatFn = createSortFn(item => item.timeSinceLastHeartbeat);
5050
public readonly sortSlotsNumberFn = createSortFn(item => item.slotsNumber);
5151
public readonly sortFreeSlotsFn = createSortFn(item => item.freeSlots);
52+
public readonly sortTasksFn = createSortFn(item => item.numberOfTasks);
5253
public readonly sortCpuCoresFn = createSortFn(item => item.hardware?.cpuCores);
5354
public readonly sortPhysicalMemoryFn = createSortFn(item => item.hardware?.physicalMemory);
5455
public readonly sortFreeMemoryFn = createSortFn(item => item.hardware?.freeMemory);

flink-runtime-web/web-dashboard/src/app/pages/task-manager/status/task-manager-status.component.html

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,12 +22,15 @@
2222
<flink-blocked-badge *ngIf="taskManagerDetail?.blocked"></flink-blocked-badge>
2323
</div>
2424
<nz-descriptions *ngIf="taskManagerDetail" nzBordered nzSize="small">
25-
<nz-descriptions-item [nzSpan]="2" nzTitle="Path">
25+
<nz-descriptions-item [nzSpan]="1" nzTitle="Path">
2626
{{ taskManagerDetail.path }}
2727
</nz-descriptions-item>
2828
<nz-descriptions-item [nzSpan]="1" nzTitle="Free/All Slots">
2929
{{ taskManagerDetail.freeSlots }} / {{ taskManagerDetail.slotsNumber }}
3030
</nz-descriptions-item>
31+
<nz-descriptions-item [nzSpan]="1" nzTitle="Tasks">
32+
{{ taskManagerDetail.numberOfTasks }}
33+
</nz-descriptions-item>
3134
<nz-descriptions-item [nzSpan]="1" nzTitle="Last Heartbeat">
3235
{{ taskManagerDetail.timeSinceLastHeartbeat | date: 'yyyy-MM-dd HH:mm:ss' }}
3336
</nz-descriptions-item>

flink-runtime/src/main/java/org/apache/flink/runtime/resourcemanager/ResourceManager.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -687,6 +687,10 @@ public CompletableFuture<Collection<TaskManagerInfo>> requestTaskManagerInfo(Dur
687687
taskManagerHeartbeatManager.getLastHeartbeatFrom(resourceId),
688688
slotManager.getNumberRegisteredSlotsOf(taskExecutor.getInstanceID()),
689689
slotManager.getNumberFreeSlotsOf(taskExecutor.getInstanceID()),
690+
(int)
691+
slotManager
692+
.getLoadingWeightOf(taskExecutor.getInstanceID())
693+
.getLoading(),
690694
slotManager.getRegisteredResourceOf(taskExecutor.getInstanceID()),
691695
slotManager.getFreeResourceOf(taskExecutor.getInstanceID()),
692696
taskExecutor.getHardwareDescription(),
@@ -717,6 +721,7 @@ public CompletableFuture<TaskManagerInfoWithSlots> requestTaskManagerDetailsInfo
717721
taskManagerHeartbeatManager.getLastHeartbeatFrom(resourceId),
718722
slotManager.getNumberRegisteredSlotsOf(instanceId),
719723
slotManager.getNumberFreeSlotsOf(instanceId),
724+
(int) slotManager.getLoadingWeightOf(instanceId).getLoading(),
720725
slotManager.getRegisteredResourceOf(instanceId),
721726
slotManager.getFreeResourceOf(instanceId),
722727
taskExecutor.getHardwareDescription(),

flink-runtime/src/main/java/org/apache/flink/runtime/resourcemanager/slotmanager/ClusterResourceStatisticsProvider.java

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,10 +20,14 @@
2020

2121
import org.apache.flink.runtime.clusterframework.types.ResourceProfile;
2222
import org.apache.flink.runtime.instance.InstanceID;
23+
import org.apache.flink.runtime.scheduler.loading.LoadingWeight;
2324

2425
/** Provides statistics of cluster resources. */
2526
public interface ClusterResourceStatisticsProvider {
2627

28+
/** Get total loading weight of the current instance. */
29+
LoadingWeight getLoadingWeightOf(InstanceID instanceId);
30+
2731
/** Get total number of registered slots. */
2832
int getNumberRegisteredSlots();
2933

flink-runtime/src/main/java/org/apache/flink/runtime/resourcemanager/slotmanager/FineGrainedSlotManager.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@
3434
import org.apache.flink.runtime.resourcemanager.WorkerResourceSpec;
3535
import org.apache.flink.runtime.resourcemanager.registration.TaskExecutorConnection;
3636
import org.apache.flink.runtime.rest.messages.taskmanager.SlotInfo;
37+
import org.apache.flink.runtime.scheduler.loading.LoadingWeight;
3738
import org.apache.flink.runtime.slots.ResourceRequirement;
3839
import org.apache.flink.runtime.slots.ResourceRequirements;
3940
import org.apache.flink.runtime.taskexecutor.SlotReport;
@@ -755,6 +756,11 @@ private Set<PendingTaskManagerId> allocateTaskManagersAccordingTo(
755756
// Legacy APIs
756757
// ---------------------------------------------------------------------------------------------
757758

759+
@Override
760+
public LoadingWeight getLoadingWeightOf(InstanceID instanceId) {
761+
return taskManagerTracker.getLoadingWeightOf(instanceId);
762+
}
763+
758764
@Override
759765
public int getNumberRegisteredSlots() {
760766
return taskManagerTracker.getNumberRegisteredSlots();

flink-runtime/src/main/java/org/apache/flink/runtime/resourcemanager/slotmanager/FineGrainedTaskManagerRegistration.java

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,8 +22,12 @@
2222
import org.apache.flink.runtime.clusterframework.types.ResourceProfile;
2323
import org.apache.flink.runtime.instance.InstanceID;
2424
import org.apache.flink.runtime.resourcemanager.registration.TaskExecutorConnection;
25+
import org.apache.flink.runtime.scheduler.loading.DefaultLoadingWeight;
26+
import org.apache.flink.runtime.scheduler.loading.LoadingWeight;
2527
import org.apache.flink.util.Preconditions;
2628

29+
import javax.annotation.Nonnull;
30+
2731
import java.util.Collections;
2832
import java.util.HashMap;
2933
import java.util.Map;
@@ -163,4 +167,12 @@ public void notifyAllocation(
163167
slots.put(allocationId, taskManagerSlot);
164168
idleSince = Long.MAX_VALUE;
165169
}
170+
171+
@Nonnull
172+
@Override
173+
public LoadingWeight getLoading() {
174+
return slots.values().stream()
175+
.map(TaskManagerSlotInformation::getLoading)
176+
.reduce(DefaultLoadingWeight.EMPTY, LoadingWeight::merge);
177+
}
166178
}

0 commit comments

Comments
 (0)