Skip to content

Commit ecbcdbe

Browse files
committed
feat(sysadvisor): report numa memory_bandwidth to kcnr
1 parent 74e52fb commit ecbcdbe

5 files changed

Lines changed: 170 additions & 16 deletions

File tree

‎go.mod‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -176,6 +176,7 @@ require (
176176
)
177177

178178
replace (
179+
github.com/kubewharf/katalyst-api => github.com/syc4704413/katalyst-api v0.0.2-syc-mem-bw
179180
k8s.io/api => k8s.io/api v0.24.6
180181
k8s.io/apiextensions-apiserver => k8s.io/apiextensions-apiserver v0.24.6
181182
k8s.io/apimachinery => k8s.io/apimachinery v0.24.6

‎go.sum‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -574,8 +574,6 @@ github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ=
574574
github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI=
575575
github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY=
576576
github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE=
577-
github.com/kubewharf/katalyst-api v0.5.11-0.20260423040236-f1a2330d266e h1:/0LP1rzsXEI09g3eqUZCJcyv/XGU8MofdFyi7NBtZA0=
578-
github.com/kubewharf/katalyst-api v0.5.11-0.20260423040236-f1a2330d266e/go.mod h1:BZMVGVl3EP0eCn5xsDgV41/gjYkoh43abIYxrB10e3k=
579577
github.com/kubewharf/kubelet v1.24.6-kubewharf-pre.3 h1:oFwpVVeDESqgzkse3iSODzHRKX495VvUgslu2hhepCc=
580578
github.com/kubewharf/kubelet v1.24.6-kubewharf-pre.3/go.mod h1:MxbSZUx3wXztFneeelwWWlX7NAAStJ6expqq7gY2J3c=
581579
github.com/kyoh86/exportloopref v0.1.7/go.mod h1:h1rDl2Kdj97+Kwh4gdz3ujE7XHmH51Q0lUiZ1z4NLj8=
@@ -900,6 +898,8 @@ github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO
900898
github.com/stretchr/testify v1.8.3 h1:RP3t2pwF7cMEbC1dqtB6poj3niw/9gnV4Cjg5oW5gtY=
901899
github.com/stretchr/testify v1.8.3/go.mod h1:sz/lmYIOXD/1dqDmKjjqLyZ2RngseejIcXlSw2iwfAo=
902900
github.com/subosito/gotenv v1.2.0/go.mod h1:N0PQaV/YGNqwC0u51sEeR/aUtSLEXKX9iv69rRypqCw=
901+
github.com/syc4704413/katalyst-api v0.0.2-syc-mem-bw h1:z70YV8teGYYdXFTP19Ayrp5vpT2uxwF6ihjH3gNf4p0=
902+
github.com/syc4704413/katalyst-api v0.0.2-syc-mem-bw/go.mod h1:BZMVGVl3EP0eCn5xsDgV41/gjYkoh43abIYxrB10e3k=
903903
github.com/syndtr/gocapability v0.0.0-20200815063812-42c35b437635 h1:kdXcSzyDtseVEc4yCz2qF8ZrQvIDBJLl4S1c3GCXmoI=
904904
github.com/syndtr/gocapability v0.0.0-20200815063812-42c35b437635/go.mod h1:hkRG7XYTFWNJGYcbNJQlaLq0fg1yr4J4t/NcTQtrfww=
905905
github.com/tdakkota/asciicheck v0.0.0-20200416190851-d7f85be797a2/go.mod h1:yHp0ai0Z9gUljN3o0xMhYJnH/IcvkdTBOX2fmJ93JEM=

‎pkg/agent/sysadvisor/plugin/qosaware/reporter/nodemetric_reporter.go‎

Lines changed: 125 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,7 @@ import (
4343
"github.com/kubewharf/katalyst-core/pkg/config/generic"
4444
"github.com/kubewharf/katalyst-core/pkg/consts"
4545
"github.com/kubewharf/katalyst-core/pkg/metaserver"
46+
malachitetypes "github.com/kubewharf/katalyst-core/pkg/metaserver/agent/metric/provisioner/malachite/types"
4647
"github.com/kubewharf/katalyst-core/pkg/metrics"
4748
schedutil "github.com/kubewharf/katalyst-core/pkg/scheduler/util"
4849
"github.com/kubewharf/katalyst-core/pkg/util"
@@ -303,6 +304,13 @@ func (p *nodeMetricsReporterPlugin) getNodeMetricInfo() (*nodeapis.NodeMetricInf
303304
nmi.GenericUsage.CPU = cpuUsage
304305
}
305306

307+
memoryBandwidthUsage, err := p.getNodeMemoryBandwidthUsage()
308+
if err != nil {
309+
errList = append(errList, err)
310+
} else {
311+
nmi.GenericUsage.MemoryBandwidth = memoryBandwidthUsage
312+
}
313+
306314
for numaID := 0; numaID < p.metaServer.NumNUMANodes; numaID++ {
307315
numaUsage := nodeapis.NUMAMetricInfo{NUMAId: numaID, Usage: &nodeapis.ResourceMetric{}}
308316
memoryUsage, err := p.getNodeNUMAMemoryUsage(numaID)
@@ -319,6 +327,13 @@ func (p *nodeMetricsReporterPlugin) getNodeMetricInfo() (*nodeapis.NodeMetricInf
319327
numaUsage.Usage.CPU = numaCpuUsage
320328
}
321329

330+
numaMemoryBandwidthUsage, err := p.getNodeNUMAMemoryBandwidthUsage(numaID)
331+
if err != nil {
332+
errList = append(errList, err)
333+
} else {
334+
numaUsage.Usage.MemoryBandwidth = numaMemoryBandwidthUsage
335+
}
336+
322337
nmi.NUMAUsage = append(nmi.NUMAUsage, numaUsage)
323338
}
324339
return nmi, errors.NewAggregate(errList)
@@ -394,6 +409,7 @@ func (p *nodeMetricsReporterPlugin) getPodUsage(pod *v1.Pod) (v1.ResourceList, m
394409
}
395410
podCPUUsage := .0
396411
podMemUsage := .0
412+
podMemoryBandwidthUsage := resource.NewQuantity(0, resource.BinarySI)
397413
for _, container := range containers {
398414
if container.RampUp {
399415
rampUp = true
@@ -460,19 +476,46 @@ func (p *nodeMetricsReporterPlugin) getPodUsage(pod *v1.Pod) (v1.ResourceList, m
460476
usages[v1.ResourceMemory] = memUsage
461477
numaUsage[numaID] = usages
462478
}
479+
480+
containerNUMAMBWUsage, containerTotalMBW, err := p.getContainerNUMAMemoryBandwidthUsage(string(pod.UID), container.ContainerName)
481+
if err != nil {
482+
errList = append(errList, fmt.Errorf("failed to get container NUMA memory bandwidth usage, podUID=%v, containerName=%v, err=%v",
483+
pod.UID, container.ContainerName, err))
484+
} else {
485+
podMemoryBandwidthUsage.Add(*containerTotalMBW)
486+
for numaID, mbw := range containerNUMAMBWUsage {
487+
usages, ok := numaUsage[numaID]
488+
if !ok {
489+
usages = make(v1.ResourceList)
490+
}
491+
492+
mbwUsage, ok := usages[apiconsts.ResourceMemoryBandwidth]
493+
if !ok {
494+
mbwUsage = *resource.NewQuantity(0, resource.BinarySI)
495+
}
496+
mbwUsage.Add(mbw)
497+
usages[apiconsts.ResourceMemoryBandwidth] = mbwUsage
498+
numaUsage[numaID] = usages
499+
}
500+
}
463501
}
464502

465503
cpu := resource.NewMilliQuantity(int64(podCPUUsage*1000), resource.DecimalSI)
466504
memory := resource.NewQuantity(int64(podMemUsage), resource.BinarySI)
467505

468-
return v1.ResourceList{v1.ResourceMemory: *memory, v1.ResourceCPU: *cpu}, numaUsage, assignedNUMAs, rampUp, errors.NewAggregate(errList)
506+
return v1.ResourceList{
507+
v1.ResourceMemory: *memory,
508+
v1.ResourceCPU: *cpu,
509+
apiconsts.ResourceMemoryBandwidth: *podMemoryBandwidthUsage,
510+
}, numaUsage, assignedNUMAs, rampUp, errors.NewAggregate(errList)
469511
}
470512

471513
func (p *nodeMetricsReporterPlugin) getGroupUsage(pods []*v1.Pod, qosLevel string) (*nodeapis.ResourceMetric, []nodeapis.NUMAMetricInfo, []*v1.Pod, error) {
472514
var errList []error
473515

474516
cpu := resource.NewQuantity(0, resource.DecimalSI)
475517
memory := resource.NewQuantity(0, resource.BinarySI)
518+
memoryBandwidth := resource.NewQuantity(0, resource.BinarySI)
476519

477520
numaUsages := make(map[int]v1.ResourceList)
478521

@@ -494,14 +537,28 @@ func (p *nodeMetricsReporterPlugin) getGroupUsage(pods []*v1.Pod, qosLevel strin
494537
}
495538

496539
for numaID := range assignedNUMAs.ToSliceInt() {
497-
podNUMAUsage[numaID] = map[v1.ResourceName]resource.Quantity{
498-
v1.ResourceCPU: *resource.NewMilliQuantity(req.Cpu().MilliValue()/int64(assignedNUMAs.Size()), resource.DecimalSI),
499-
v1.ResourceMemory: *resource.NewQuantity(req.Memory().Value()/int64(assignedNUMAs.Size()), resource.BinarySI),
540+
usages, ok := podNUMAUsage[numaID]
541+
if !ok {
542+
usages = make(v1.ResourceList)
500543
}
544+
545+
usages[v1.ResourceCPU] = *resource.NewMilliQuantity(
546+
req.Cpu().MilliValue()/int64(assignedNUMAs.Size()),
547+
resource.DecimalSI,
548+
)
549+
usages[v1.ResourceMemory] = *resource.NewQuantity(
550+
req.Memory().Value()/int64(assignedNUMAs.Size()),
551+
resource.BinarySI,
552+
)
553+
554+
podNUMAUsage[numaID] = usages
501555
}
502556
}
503557
cpu.Add(*podUsage.Cpu())
504558
memory.Add(*podUsage.Memory())
559+
if mbw, ok := podUsage[apiconsts.ResourceMemoryBandwidth]; ok {
560+
memoryBandwidth.Add(mbw)
561+
}
505562

506563
for numaID, podUsages := range podNUMAUsage {
507564
usages, ok := numaUsages[numaID]
@@ -539,6 +596,13 @@ func (p *nodeMetricsReporterPlugin) getGroupUsage(pods []*v1.Pod, qosLevel strin
539596
resourceMetric.CPU = aggCPU
540597
}
541598

599+
aggMBW := p.getAggregatedMetric(memoryBandwidth, apiconsts.ResourceMemoryBandwidth, "getGroupUsage", qosLevel, "memory_bandwidth")
600+
if aggMBW == nil {
601+
errList = append(errList, fmt.Errorf("failed to get enough samples for group memory bandwidth, qosLevel=%v", qosLevel))
602+
} else {
603+
resourceMetric.MemoryBandwidth = aggMBW
604+
}
605+
542606
for numaID := 0; numaID < p.metaServer.NumNUMANodes; numaID++ {
543607
resourceUsages, ok := numaUsages[numaID]
544608
if !ok {
@@ -565,6 +629,14 @@ func (p *nodeMetricsReporterPlugin) getGroupUsage(pods []*v1.Pod, qosLevel strin
565629
resourceNUMAMetric.Memory = aggNUMAMem
566630
}
567631

632+
mbwUsage := resourceUsages[apiconsts.ResourceMemoryBandwidth]
633+
aggNUMAMBW := p.getAggregatedMetric(&mbwUsage, apiconsts.ResourceMemoryBandwidth, "getGroupNUMAUsage", qosLevel, "memory_bandwidth", strconv.Itoa(numaID))
634+
if aggNUMAMBW == nil {
635+
errList = append(errList, fmt.Errorf("failed to get enough samples for group numa memory bandwidth, qosLevel=%v, numa=%v", qosLevel, numaID))
636+
} else {
637+
resourceNUMAMetric.MemoryBandwidth = aggNUMAMBW
638+
}
639+
568640
resourceNUMAMetrics = append(resourceNUMAMetrics, nodeapis.NUMAMetricInfo{
569641
NUMAId: numaID,
570642
Usage: &resourceNUMAMetric,
@@ -660,3 +732,52 @@ func (p *nodeMetricsReporterPlugin) getAggregatedMetric(value *resource.Quantity
660732
}
661733
return aggregator.GetWindowedResources(*value)
662734
}
735+
736+
func (p *nodeMetricsReporterPlugin) getNodeNUMAMemoryBandwidthUsage(numaID int) (*resource.Quantity, error) {
737+
metricData, err := p.metaServer.GetNumaMetric(numaID, consts.MetricTotalPsMemBandwidthNuma)
738+
if err != nil {
739+
return nil, fmt.Errorf("failed to get %s for numa=%d: %w", consts.MetricTotalPsMemBandwidthNuma, numaID, err)
740+
}
741+
742+
return resource.NewQuantity(int64(metricData.Value), resource.BinarySI), nil
743+
}
744+
745+
func (p *nodeMetricsReporterPlugin) getNodeMemoryBandwidthUsage() (*resource.Quantity, error) {
746+
total := resource.NewQuantity(0, resource.BinarySI)
747+
748+
for numaID := 0; numaID < p.metaServer.NumNUMANodes; numaID++ {
749+
q, err := p.getNodeNUMAMemoryBandwidthUsage(numaID)
750+
if err != nil {
751+
return nil, err
752+
}
753+
total.Add(*q)
754+
}
755+
return total, nil
756+
}
757+
758+
func (p *nodeMetricsReporterPlugin) getContainerNUMAMemoryBandwidthUsage(podUID, containerName string) (map[int]resource.Quantity, *resource.Quantity, error) {
759+
key := fmt.Sprintf("%s/%s/%s", consts.MetricMbmTotalPsContainerL3, podUID, containerName)
760+
raw := p.metaServer.MetricsFetcher.GetByStringIndex(key)
761+
762+
statsMap, ok := raw.(map[int]malachitetypes.L3CacheBytesPS)
763+
if !ok || len(statsMap) == 0 {
764+
return map[int]resource.Quantity{}, resource.NewQuantity(0, resource.BinarySI), nil
765+
}
766+
767+
total := resource.NewQuantity(0, resource.BinarySI)
768+
numaUsage := make(map[int]resource.Quantity)
769+
770+
for _, l3Stat := range statsMap {
771+
q := *resource.NewQuantity(int64(l3Stat.MbmTotalBytesPS), resource.BinarySI)
772+
total.Add(q)
773+
774+
exist, ok := numaUsage[l3Stat.NumaID]
775+
if !ok {
776+
exist = *resource.NewQuantity(0, resource.BinarySI)
777+
}
778+
exist.Add(q)
779+
numaUsage[l3Stat.NumaID] = exist
780+
}
781+
782+
return numaUsage, total, nil
783+
}

‎pkg/consts/metric.go‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -345,6 +345,8 @@ const (
345345
MetricMbmlocalPsContainer = "mbm.local.ps.container"
346346
MetricMbmVictimPsContainer = "mbm.victim.ps.container"
347347
MetricResctrlDataContainer = "resctrl.data.container"
348+
349+
MetricMbmTotalPsContainerL3 = "mbm.total.ps.container.l3"
348350
)
349351

350352
// container blkio metrics

‎pkg/metaserver/agent/metric/provisioner/malachite/provisioner_calculate.go‎

Lines changed: 40 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -338,19 +338,26 @@ func (m *MalachiteMetricsProvisioner) setContainerRateMetric(podUID, containerNa
338338

339339
// setContainerMbmTotalMetric calcuate the total memory bandwidth usage of a container
340340
func (m *MalachiteMetricsProvisioner) setContainerMbmTotalMetric(podUID, containerName string, resctrlData types.MbmbandData, updateTime *time.Time) {
341+
resctrlKey := getContainerStringIndexKey(consts.MetricResctrlDataContainer, podUID, containerName)
342+
l3MetricKey := getContainerStringIndexKey(consts.MetricMbmTotalPsContainerL3, podUID, containerName)
343+
341344
// 1. get previous data
342-
oldResctrlData, ok := m.metricStore.GetByStringIndex(consts.MetricResctrlDataContainer).(types.MbmbandData)
345+
oldResctrlData, ok := m.metricStore.GetByStringIndex(resctrlKey).(types.MbmbandData)
343346
if !ok || len(oldResctrlData.Mbm) == 0 || len(resctrlData.Mbm) == 0 {
344347
// No previous or current data, store current for next round and return
345-
m.metricStore.SetByStringIndex(consts.MetricResctrlDataContainer, resctrlData)
348+
//m.metricStore.SetByStringIndex(consts.MetricResctrlDataContainer, resctrlData)
349+
m.metricStore.SetByStringIndex(resctrlKey, resctrlData)
350+
m.metricStore.SetByStringIndex(l3MetricKey, map[int]types.L3CacheBytesPS{})
346351
return
347352
}
348353

349354
// 2. calculate the time interval
350355
prevMetric, err := m.metricStore.GetContainerMetric(podUID, containerName, consts.MetricMbmTotalPsContainer)
351356
if err != nil || prevMetric.Time == nil {
352357
// if no previous metric, store current for next round and return
353-
m.metricStore.SetByStringIndex(consts.MetricResctrlDataContainer, resctrlData)
358+
//m.metricStore.SetByStringIndex(consts.MetricResctrlDataContainer, resctrlData)
359+
m.metricStore.SetByStringIndex(resctrlKey, resctrlData)
360+
m.metricStore.SetByStringIndex(l3MetricKey, map[int]types.L3CacheBytesPS{})
354361
m.metricStore.SetContainerMetric(podUID, containerName, consts.MetricMbmTotalPsContainer,
355362
metric.MetricData{Value: 0, Time: updateTime})
356363
m.metricStore.SetContainerMetric(podUID, containerName, consts.MetricMbmlocalPsContainer,
@@ -362,7 +369,15 @@ func (m *MalachiteMetricsProvisioner) setContainerMbmTotalMetric(podUID, contain
362369
if timeInterval <= 0 {
363370
return
364371
}
372+
373+
cpuCodeName := ""
374+
if val, ok := m.metricStore.GetByStringIndex(consts.MetricCPUCodeName).(string); ok {
375+
cpuCodeName = val
376+
}
377+
365378
var totalMbmBytesPS, totalLocalBytesPS float64
379+
l3CacheBandwidthStats := make(map[int]types.L3CacheBytesPS)
380+
366381
for _, l3Mon := range resctrlData.Mbm {
367382
l3CacheID := l3Mon.ID
368383

@@ -395,16 +410,27 @@ func (m *MalachiteMetricsProvisioner) setContainerMbmTotalMetric(podUID, contain
395410
continue
396411
}
397412
}
413+
414+
numaID, maxBytesPS := getNumaAndMaxBandwidth(int(l3CacheID), cpuCodeName)
415+
adjustedTotalBytesPS := totalBytesPS
416+
417+
if strings.Contains(cpuCodeName, consts.AMDGenoaArch) {
418+
// Notice: data adjustments needed due to genoa hardware bug
419+
adjustedTotalBytesPS += localBytesPS / 3 * 2
420+
}
421+
totalMbmBytesPS += adjustedTotalBytesPS
398422
totalLocalBytesPS += localBytesPS
399-
cpuCodeName := m.metricStore.GetByStringIndex(consts.MetricCPUCodeName)
400-
if codeName, ok := cpuCodeName.(string); ok {
401-
if strings.Contains(codeName, consts.AMDGenoaArch) {
402-
// Notice: data adjustments needed due to genoa hardware bug
403-
totalLocalBytesPS += localBytesPS / 3 * 2
404-
}
423+
l3CacheBandwidthStats[int(l3CacheID)] = types.L3CacheBytesPS{
424+
NumaID: numaID,
425+
MbmTotalBytesPS: uint64(adjustedTotalBytesPS),
426+
MbmLocalBytesPS: uint64(localBytesPS),
427+
MbmVictimBytesPS: 0,
428+
MBMMaxBytesPS: maxBytesPS,
405429
}
406430
}
407-
m.metricStore.SetByStringIndex(consts.MetricResctrlDataContainer, resctrlData)
431+
//m.metricStore.SetByStringIndex(consts.MetricResctrlDataContainer, resctrlData)
432+
m.metricStore.SetByStringIndex(resctrlKey, resctrlData)
433+
m.metricStore.SetByStringIndex(l3MetricKey, l3CacheBandwidthStats)
408434
m.metricStore.SetContainerMetric(podUID, containerName, consts.MetricMbmTotalPsContainer,
409435
metric.MetricData{Value: totalMbmBytesPS, Time: updateTime})
410436
m.metricStore.SetContainerMetric(podUID, containerName, consts.MetricMbmlocalPsContainer,
@@ -581,3 +607,7 @@ func aggregateNUMABytesPS(l3BytesPS map[int]types.L3CacheBytesPS, cpuCodeName st
581607
}
582608
return numaBytesPS
583609
}
610+
611+
func getContainerStringIndexKey(metricName, podUID, containerName string) string {
612+
return fmt.Sprintf("%s/%s/%s", metricName, podUID, containerName)
613+
}

0 commit comments

Comments
 (0)