Просмотр исходного кода

move device usage computation

Signed-off-by: Sean Holcomb <seanholcomb@gmail.com>
Sean Holcomb 2 недель назад
Родитель
Сommit
efd95ea2b8

+ 41 - 0
core/pkg/compute/kubemodel/container.go

@@ -38,6 +38,9 @@ func (km *KubeModel) computeContainers(kms *kubemodel.KubeModelSet, start, end t
 	ramUsageAvgFuture := source.WithGroup(grp, metrics.QueryRAMUsageAvg(start, end))
 	ramUsageMaxFuture := source.WithGroup(grp, metrics.QueryRAMUsageMax(start, end))
 
+	deviceUsageAvgFuture := source.WithGroup(grp, metrics.QueryDCGMContainerUsageAvg(start, end))
+	deviceUsageMaxFuture := source.WithGroup(grp, metrics.QueryDCGMContainerUsageMax(start, end))
+
 	type containerKey struct {
 		podUID string
 		name   string
@@ -161,6 +164,44 @@ func (km *KubeModel) computeContainers(kms *kubemodel.KubeModelSet, start, end t
 		}
 	}
 
+	deviceUsageAvgResult, _ := deviceUsageAvgFuture.Await()
+	for _, res := range deviceUsageAvgResult {
+		if res.PodUID == "" || res.Container == "" {
+			continue
+		}
+		key := containerKey{podUID: res.PodUID, name: res.Container}
+		container, ok := containerMap[key]
+		if !ok {
+			log.Warnf("container %s/%s has not been initialized to add device usage avg", res.PodUID, res.Container)
+			continue
+		}
+		if container.DeviceUsages == nil {
+			container.DeviceUsages = make(map[string]kubemodel.DeviceUsage)
+		}
+		usage := container.DeviceUsages[res.UUID]
+		usage.UsageAvg = res.Value
+		container.DeviceUsages[res.UUID] = usage
+	}
+
+	deviceUsageMaxResult, _ := deviceUsageMaxFuture.Await()
+	for _, res := range deviceUsageMaxResult {
+		if res.PodUID == "" || res.Container == "" {
+			continue
+		}
+		key := containerKey{podUID: res.PodUID, name: res.Container}
+		container, ok := containerMap[key]
+		if !ok {
+			log.Warnf("container %s/%s has not been initialized to add device usage max", res.PodUID, res.Container)
+			continue
+		}
+		if container.DeviceUsages == nil {
+			container.DeviceUsages = make(map[string]kubemodel.DeviceUsage)
+		}
+		usage := container.DeviceUsages[res.UUID]
+		usage.UsageMax = res.Value
+		container.DeviceUsages[res.UUID] = usage
+	}
+
 	for _, container := range containerMap {
 		err := kms.RegisterContainer(container)
 		if err != nil {

+ 70 - 0
core/pkg/compute/kubemodel/container_test.go

@@ -161,6 +161,76 @@ func TestComputeContainers(t *testing.T) {
 				},
 			},
 		},
+		{
+			name: "device usage avg and max are populated",
+			overrides: map[string]any{
+				source.QueryContainerUptime: []*source.ContainerUptimeResult{
+					{UptimeResult: source.UptimeResult{UID: "pod-1", First: start, Last: end}, Container: "training"},
+				},
+				source.QueryDCGMContainerUsageAvg: []*source.DCGMDeviceContainerUsageResult{
+					{UUID: "GPU-abc123", PodUID: "pod-1", Container: "training", Value: 0.75},
+				},
+				source.QueryDCGMContainerUsageMax: []*source.DCGMDeviceContainerUsageResult{
+					{UUID: "GPU-abc123", PodUID: "pod-1", Container: "training", Value: 0.95},
+				},
+			},
+			want: map[string]*kubemodel.Container{
+				"pod-1/training": {
+					PodUID:           "pod-1",
+					Name:             "training",
+					Start:            start,
+					End:              end,
+					ResourceRequests: kubemodel.ResourceQuantities{},
+					ResourceLimits:   kubemodel.ResourceQuantities{},
+					DeviceUsages: map[string]kubemodel.DeviceUsage{
+						"GPU-abc123": {UsageAvg: 0.75, UsageMax: 0.95},
+					},
+				},
+			},
+		},
+		{
+			name: "device usage with empty pod uid or container is ignored",
+			overrides: map[string]any{
+				source.QueryContainerUptime: []*source.ContainerUptimeResult{
+					{UptimeResult: source.UptimeResult{UID: "pod-1", First: start, Last: end}, Container: "training"},
+				},
+				source.QueryDCGMContainerUsageAvg: []*source.DCGMDeviceContainerUsageResult{
+					{UUID: "GPU-abc123", PodUID: "", Container: "training", Value: 0.5},
+					{UUID: "GPU-abc123", PodUID: "pod-1", Container: "", Value: 0.5},
+				},
+			},
+			want: map[string]*kubemodel.Container{
+				"pod-1/training": {
+					PodUID:           "pod-1",
+					Name:             "training",
+					Start:            start,
+					End:              end,
+					ResourceRequests: kubemodel.ResourceQuantities{},
+					ResourceLimits:   kubemodel.ResourceQuantities{},
+				},
+			},
+		},
+		{
+			name: "device usage for unknown container is ignored",
+			overrides: map[string]any{
+				source.QueryContainerUptime: []*source.ContainerUptimeResult{
+					{UptimeResult: source.UptimeResult{UID: "pod-1", First: start, Last: end}, Container: "main"},
+				},
+				source.QueryDCGMContainerUsageAvg: []*source.DCGMDeviceContainerUsageResult{
+					{UUID: "GPU-abc123", PodUID: "pod-1", Container: "training", Value: 0.75},
+				},
+			},
+			want: map[string]*kubemodel.Container{
+				"pod-1/main": {
+					PodUID:           "pod-1",
+					Name:             "main",
+					Start:            start,
+					End:              end,
+					ResourceRequests: kubemodel.ResourceQuantities{},
+					ResourceLimits:   kubemodel.ResourceQuantities{},
+				},
+			},
+		},
 		{
 			name: "resource requests for unknown container are ignored",
 			overrides: map[string]any{

+ 0 - 29
core/pkg/compute/kubemodel/device.go

@@ -14,8 +14,6 @@ func (km *KubeModel) computeDevices(kms *kubemodel.KubeModelSet, start, end time
 
 	infoFuture := source.WithGroup(grp, metrics.QueryDCGMDeviceInfo(start, end))
 	uptimeFuture := source.WithGroup(grp, metrics.QueryDCGMDeviceUptime(start, end))
-	usageAvgFuture := source.WithGroup(grp, metrics.QueryDCGMContainerUsageAvg(start, end))
-	usageMaxFuture := source.WithGroup(grp, metrics.QueryDCGMContainerUsageMax(start, end))
 
 	deviceMap := make(map[string]*kubemodel.Device)
 
@@ -52,32 +50,5 @@ func (km *KubeModel) computeDevices(kms *kubemodel.KubeModelSet, start, end time
 		}
 	}
 
-	setUsage := func(res *source.DCGMDeviceContainerUsageResult, apply func(*kubemodel.DeviceUsage)) {
-		if res.PodUID == "" || res.Container == "" {
-			return
-		}
-		key := (&kubemodel.Container{PodUID: res.PodUID, Name: res.Container}).GetKey()
-		container, ok := kms.Containers[key]
-		if !ok {
-			return
-		}
-		if container.DeviceUsages == nil {
-			container.DeviceUsages = make(map[string]kubemodel.DeviceUsage)
-		}
-		usage := container.DeviceUsages[res.UUID]
-		apply(&usage)
-		container.DeviceUsages[res.UUID] = usage
-	}
-
-	usageAvgResult, _ := usageAvgFuture.Await()
-	for _, res := range usageAvgResult {
-		setUsage(res, func(u *kubemodel.DeviceUsage) { u.UsageAvg = res.Value })
-	}
-
-	usageMaxResult, _ := usageMaxFuture.Await()
-	for _, res := range usageMaxResult {
-		setUsage(res, func(u *kubemodel.DeviceUsage) { u.UsageMax = res.Value })
-	}
-
 	return nil
 }

+ 11 - 111
core/pkg/compute/kubemodel/device_test.go

@@ -16,17 +16,14 @@ func TestComputeDevices(t *testing.T) {
 	end := start.Add(time.Hour)
 
 	tests := []struct {
-		name         string
-		overrides    map[string]any
-		containers   map[string]*kubemodel.Container
-		wantDevices  map[string]*kubemodel.Device
-		wantUsageKey string
-		wantUsage    *kubemodel.DeviceUsage
+		name      string
+		overrides map[string]any
+		want      map[string]*kubemodel.Device
 	}{
 		{
-			name:        "no data returns empty device map",
-			overrides:   map[string]any{},
-			wantDevices: map[string]*kubemodel.Device{},
+			name:      "no data returns empty device map",
+			overrides: map[string]any{},
+			want:      map[string]*kubemodel.Device{},
 		},
 		{
 			name: "basic device info and uptime",
@@ -38,7 +35,7 @@ func TestComputeDevices(t *testing.T) {
 					{UUID: "GPU-abc123", First: start, Last: end},
 				},
 			},
-			wantDevices: map[string]*kubemodel.Device{
+			want: map[string]*kubemodel.Device{
 				"GPU-abc123": {
 					UUID:      "GPU-abc123",
 					Device:    "nvidia0",
@@ -55,7 +52,7 @@ func TestComputeDevices(t *testing.T) {
 					{UUID: "GPU-abc123", Device: "nvidia0", ModelName: "A100"},
 				},
 			},
-			wantDevices: map[string]*kubemodel.Device{},
+			want: map[string]*kubemodel.Device{},
 		},
 		{
 			name: "device with empty uuid is skipped",
@@ -67,7 +64,7 @@ func TestComputeDevices(t *testing.T) {
 					{UUID: "GPU-abc123", First: start, Last: end},
 				},
 			},
-			wantDevices: map[string]*kubemodel.Device{},
+			want: map[string]*kubemodel.Device{},
 		},
 		{
 			name: "duplicate device info entries use first occurrence",
@@ -80,90 +77,7 @@ func TestComputeDevices(t *testing.T) {
 					{UUID: "GPU-abc123", First: start, Last: end},
 				},
 			},
-			wantDevices: map[string]*kubemodel.Device{
-				"GPU-abc123": {
-					UUID:      "GPU-abc123",
-					Device:    "nvidia0",
-					ModelName: "A100",
-					Start:     start,
-					End:       end,
-				},
-			},
-		},
-		{
-			name: "container usage avg and max are applied to a registered container",
-			overrides: map[string]any{
-				source.QueryDCGMDeviceInfo: []*source.DCGMDeviceInfoResult{
-					{UUID: "GPU-abc123", Device: "nvidia0", ModelName: "A100"},
-				},
-				source.QueryDCGMDeviceUptime: []*source.DCGMDeviceUptimeResult{
-					{UUID: "GPU-abc123", First: start, Last: end},
-				},
-				source.QueryDCGMContainerUsageAvg: []*source.DCGMDeviceContainerUsageResult{
-					{UUID: "GPU-abc123", PodUID: "pod-1", Container: "training", Value: 0.75},
-				},
-				source.QueryDCGMContainerUsageMax: []*source.DCGMDeviceContainerUsageResult{
-					{UUID: "GPU-abc123", PodUID: "pod-1", Container: "training", Value: 0.95},
-				},
-			},
-			containers: map[string]*kubemodel.Container{
-				"pod-1/training": {PodUID: "pod-1", Name: "training"},
-			},
-			wantDevices: map[string]*kubemodel.Device{
-				"GPU-abc123": {
-					UUID:      "GPU-abc123",
-					Device:    "nvidia0",
-					ModelName: "A100",
-					Start:     start,
-					End:       end,
-				},
-			},
-			wantUsageKey: "pod-1/training",
-			wantUsage:    &kubemodel.DeviceUsage{UsageAvg: 0.75, UsageMax: 0.95},
-		},
-		{
-			name: "usage with empty pod uid or container is ignored",
-			overrides: map[string]any{
-				source.QueryDCGMDeviceInfo: []*source.DCGMDeviceInfoResult{
-					{UUID: "GPU-abc123", Device: "nvidia0", ModelName: "A100"},
-				},
-				source.QueryDCGMDeviceUptime: []*source.DCGMDeviceUptimeResult{
-					{UUID: "GPU-abc123", First: start, Last: end},
-				},
-				source.QueryDCGMContainerUsageAvg: []*source.DCGMDeviceContainerUsageResult{
-					{UUID: "GPU-abc123", PodUID: "", Container: "training", Value: 0.5},
-					{UUID: "GPU-abc123", PodUID: "pod-1", Container: "", Value: 0.5},
-				},
-			},
-			containers: map[string]*kubemodel.Container{
-				"pod-1/training": {PodUID: "pod-1", Name: "training"},
-			},
-			wantDevices: map[string]*kubemodel.Device{
-				"GPU-abc123": {
-					UUID:      "GPU-abc123",
-					Device:    "nvidia0",
-					ModelName: "A100",
-					Start:     start,
-					End:       end,
-				},
-			},
-			wantUsageKey: "pod-1/training",
-			wantUsage:    nil,
-		},
-		{
-			name: "usage for an unregistered container is ignored",
-			overrides: map[string]any{
-				source.QueryDCGMDeviceInfo: []*source.DCGMDeviceInfoResult{
-					{UUID: "GPU-abc123", Device: "nvidia0", ModelName: "A100"},
-				},
-				source.QueryDCGMDeviceUptime: []*source.DCGMDeviceUptimeResult{
-					{UUID: "GPU-abc123", First: start, Last: end},
-				},
-				source.QueryDCGMContainerUsageAvg: []*source.DCGMDeviceContainerUsageResult{
-					{UUID: "GPU-abc123", PodUID: "pod-1", Container: "training", Value: 0.75},
-				},
-			},
-			wantDevices: map[string]*kubemodel.Device{
+			want: map[string]*kubemodel.Device{
 				"GPU-abc123": {
 					UUID:      "GPU-abc123",
 					Device:    "nvidia0",
@@ -188,25 +102,11 @@ func TestComputeDevices(t *testing.T) {
 			require.NoError(t, err)
 
 			kms := kubemodel.NewKubeModelSet(start, end)
-			if tt.containers != nil {
-				kms.Containers = tt.containers
-			}
 
 			err = km.computeDevices(kms, start, end)
 			require.NoError(t, err)
 
-			assert.Equal(t, tt.wantDevices, kms.Devices)
-
-			if tt.wantUsageKey != "" {
-				c, ok := kms.Containers[tt.wantUsageKey]
-				require.True(t, ok)
-				if tt.wantUsage == nil {
-					assert.Empty(t, c.DeviceUsages)
-				} else {
-					require.NotNil(t, c.DeviceUsages)
-					assert.Equal(t, *tt.wantUsage, c.DeviceUsages["GPU-abc123"])
-				}
-			}
+			assert.Equal(t, tt.want, kms.Devices)
 		})
 	}
 }

+ 8 - 4
core/pkg/model/kubemodel/container.go

@@ -31,10 +31,6 @@ type DeviceUsage struct {
 	UsageMax float64 `json:"usageMax"`
 }
 
-func (c *Container) GetKey() string {
-	return fmt.Sprintf("%s/%s", c.PodUID, c.Name)
-}
-
 func (c *Container) ValidateContainer(window Window) error {
 	if c.PodUID == "" {
 		return fmt.Errorf("PodUID is missing for Container with name '%s'", c.Name)
@@ -66,3 +62,11 @@ func (kms *KubeModelSet) RegisterContainer(container *Container) error {
 
 	return nil
 }
+
+func (c *Container) GetKey() string {
+	return ContainerKey(c.PodUID, c.Name)
+}
+
+func ContainerKey(podUID, conatinerName string) string {
+	return fmt.Sprintf("%s/%s", podUID, conatinerName)
+}