feat: optimize resource monitoring logic, performance and fix nil pointer issue (#4827)

* feat: optimize resource monitoring logic and performance

This commit introduces several optimizations to the resource monitoring logic and improves its overall performance.

**Code Structure Improvements:**

- Refactored resource usage monitoring logic into dedicated functions for Pods, PVCs, backups, and Services, enhancing code readability and maintainability.

**Performance Enhancements:**

- Implemented `FieldSelector` in `monitorPodResourceUsage`, `monitorPVCResourceUsage`, and `monitorServiceResourceUsage` to filter resources efficiently, minimizing API request overhead.
- Introduced a `gpuMutex` to safeguard the `NvidiaGpu` map, preventing concurrent read-write issues and potential race conditions.

* Fix: ObjStorageClient nil pointer

link https://github.com/labring/sealos/pull/4780/files#diff-ed55e3d94d63443dc853d1dd28cc8899b61633a6e290bc2c5d6925eb0cd41d6dR275

* Fix: non-exact field matches are not supported by the cache

* Fix: create index

* Fix: lint

* Fix: InitIndexField
This commit is contained in:
zijiren
2024-07-02 16:21:20 +08:00
committed by GitHub
parent c6e0376706
commit b109ffe2f1
2 changed files with 133 additions and 67 deletions
@@ -53,6 +53,7 @@ import (
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/resource"
"k8s.io/apimachinery/pkg/fields"
"k8s.io/apimachinery/pkg/runtime"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
@@ -68,6 +69,7 @@ type MonitorReconciler struct {
wg sync.WaitGroup
periodicReconcile time.Duration
NvidiaGpu map[string]gpu.NvidiaGPU
gpuMutex sync.Mutex
DBClient database.Interface
TrafficClient database.Interface
Properties *resources.PropertyTypeLS
@@ -116,6 +118,7 @@ func NewMonitorReconciler(mgr ctrl.Manager) (*MonitorReconciler, error) {
periodicReconcile: 1 * time.Minute,
PromURL: os.Getenv(PrometheusURL),
ObjectStorageInstance: os.Getenv(ObjectStorageInstance),
NvidiaGpu: make(map[string]gpu.NvidiaGPU),
}
concurrentLimit = env.GetInt64EnvWithDefault(ConcurrentLimit, DefaultConcurrencyLimit)
var err error
@@ -133,6 +136,19 @@ func NewMonitorReconciler(mgr ctrl.Manager) (*MonitorReconciler, error) {
return r, nil
}
func InitIndexField(mgr ctrl.Manager) error {
if err := mgr.GetFieldIndexer().IndexField(context.Background(), &corev1.PersistentVolumeClaim{}, "status.phase", func(rawObj client.Object) []string {
pvc := rawObj.(*corev1.PersistentVolumeClaim)
return []string{string(pvc.Status.Phase)}
}); err != nil {
return err
}
return mgr.GetFieldIndexer().IndexField(context.Background(), &corev1.Service{}, "spec.type", func(rawObj client.Object) []string {
svc := rawObj.(*corev1.Service)
return []string{string(svc.Spec.Type)}
})
}
func (r *MonitorReconciler) StartReconciler(ctx context.Context) error {
r.startPeriodicReconcile()
if r.TrafficClient != nil || r.ObjStorageClient != nil {
@@ -272,24 +288,70 @@ func (r *MonitorReconciler) preMonitorResourceUsage() error {
r.currentObjectMetrics = metrics
logger.Info("success query object storage resource usage", "time", time.Now().Format("2006-01-02 15:04:05"))
}
r.ObjStorageUserBackupSize = objstorage.GetUserBakFileSize(r.ObjStorageClient)
fmt.Println("ObjStorageUserBackupSize", r.ObjStorageUserBackupSize)
if r.ObjStorageClient != nil {
r.ObjStorageUserBackupSize = objstorage.GetUserBakFileSize(r.ObjStorageClient)
logger.Info("success query object storage backup size", "time", time.Now().Format("2006-01-02 15:04:05"))
}
return nil
}
func (r *MonitorReconciler) monitorResourceUsage(namespace *corev1.Namespace) error {
timeStamp := time.Now().UTC()
podList := corev1.PodList{}
resUsed := map[string]map[corev1.ResourceName]*quantity{}
resNamed := make(map[string]*resources.ResourceNamed)
if err := r.List(context.Background(), &podList, &client.ListOptions{Namespace: namespace.Name}); err != nil {
return err
if err := r.monitorPodResourceUsage(namespace.Name, resUsed, resNamed); err != nil {
return fmt.Errorf("failed to monitor pod resource usage: %v", err)
}
for _, pod := range podList.Items {
if pod.Spec.NodeName == "" || (pod.Status.Phase == corev1.PodSucceeded && time.Since(pod.Status.StartTime.Time) > 1*time.Minute) {
if err := r.monitorPVCResourceUsage(namespace.Name, resUsed, resNamed); err != nil {
return fmt.Errorf("failed to monitor PVC resource usage: %v", err)
}
if err := r.monitorDatabaseBackupUsage(namespace.Name, resUsed, resNamed); err != nil {
return fmt.Errorf("failed to monitor backup resource usage: %v", err)
}
if err := r.monitorServiceResourceUsage(namespace.Name, resUsed, resNamed); err != nil {
return fmt.Errorf("failed to monitor service resource usage: %v", err)
}
if err := r.monitorObjectStorageUsage(namespace.Name, resUsed, resNamed); err != nil {
return fmt.Errorf("failed to get object storage resource usage: %v", err)
}
var monitors []*resources.Monitor
for name, podResource := range resUsed {
isEmpty, used := r.getResourceUsed(podResource)
if isEmpty {
continue
}
podResNamed := resources.NewResourceNamed(&pod)
monitors = append(monitors, &resources.Monitor{
Category: namespace.Name,
Used: used,
Time: timeStamp,
Type: resNamed[name].Type(),
Name: resNamed[name].Name(),
})
}
return r.DBClient.InsertMonitor(context.Background(), monitors...)
}
func (r *MonitorReconciler) monitorPodResourceUsage(namespace string, resUsed map[string]map[corev1.ResourceName]*quantity, resNamed map[string]*resources.ResourceNamed) error {
podList := &corev1.PodList{}
if err := r.List(context.Background(), podList, &client.ListOptions{
Namespace: namespace,
}); err != nil {
return fmt.Errorf("failed to list pods: %v", err)
}
for i := range podList.Items {
pod := &podList.Items[i]
if pod.Spec.NodeName == "" || pod.Status.Phase == corev1.PodSucceeded && time.Since(pod.Status.StartTime.Time) > 1*time.Minute {
continue
}
podResNamed := resources.NewResourceNamed(pod)
resNamed[podResNamed.String()] = podResNamed
if resUsed[podResNamed.String()] == nil {
resUsed[podResNamed.String()] = initResources()
@@ -299,8 +361,7 @@ func (r *MonitorReconciler) monitorResourceUsage(namespace *corev1.Namespace) er
for _, container := range pod.Spec.Containers {
// gpu only use limit and not ignore pod pending status
if gpuRequest, ok := container.Resources.Limits[gpu.NvidiaGpuKey]; ok {
err := r.getGPUResourceUsage(pod, gpuRequest, resUsed[podResNamed.String()])
if err != nil {
if err := r.getGPUResourceUsage(pod, gpuRequest, resUsed[podResNamed.String()]); err != nil {
r.Logger.Error(err, "get gpu resource usage failed", "pod", pod.Name)
}
}
@@ -319,48 +380,68 @@ func (r *MonitorReconciler) monitorResourceUsage(namespace *corev1.Namespace) er
}
}
}
return nil
}
//logger.Info("mid", "namespace", namespace.Name, "time", timeStamp.Format("2006-01-02 15:04:05"), "resourceMap", resourceMap, "podsRes", podsRes)
pvcList := corev1.PersistentVolumeClaimList{}
if err := r.List(context.Background(), &pvcList, &client.ListOptions{Namespace: namespace.Name}); err != nil {
func (r *MonitorReconciler) monitorPVCResourceUsage(namespace string, resUsed map[string]map[corev1.ResourceName]*quantity, resNamed map[string]*resources.ResourceNamed) error {
pvcList := &corev1.PersistentVolumeClaimList{}
if err := r.List(context.Background(), pvcList, &client.ListOptions{
Namespace: namespace,
FieldSelector: fields.OneTermEqualSelector("status.phase", string(corev1.ClaimBound)),
}); err != nil {
return fmt.Errorf("failed to list pvc: %v", err)
}
for _, pvc := range pvcList.Items {
if pvc.Status.Phase != corev1.ClaimBound {
continue
}
for i := range pvcList.Items {
pvc := &pvcList.Items[i]
if len(pvc.OwnerReferences) > 0 && pvc.OwnerReferences[0].Kind == "BackupRepo" {
continue
}
pvcRes := resources.NewResourceNamed(&pvc)
pvcRes := resources.NewResourceNamed(pvc)
if resUsed[pvcRes.String()] == nil {
resNamed[pvcRes.String()] = pvcRes
resUsed[pvcRes.String()] = initResources()
}
resUsed[pvcRes.String()][corev1.ResourceStorage].Add(pvc.Spec.Resources.Requests[corev1.ResourceStorage])
}
if backupSize := r.ObjStorageUserBackupSize[getBackupObjectStorageName(namespace.Name)]; backupSize > 0 {
backupRes := resources.NewObjStorageResourceNamed("DB-BACKUP")
if resUsed[backupRes.String()] == nil {
resNamed[backupRes.String()] = backupRes
resUsed[backupRes.String()] = initResources()
}
resUsed[backupRes.String()][corev1.ResourceStorage].Add(*resource.NewQuantity(backupSize, resource.BinarySI))
return nil
}
func (r *MonitorReconciler) monitorDatabaseBackupUsage(namespace string, resUsed map[string]map[corev1.ResourceName]*quantity, resNamed map[string]*resources.ResourceNamed) error {
if r.ObjStorageUserBackupSize == nil {
return nil
}
svcList := corev1.ServiceList{}
if err := r.List(context.Background(), &svcList, &client.ListOptions{Namespace: namespace.Name}); err != nil {
backupSize := r.ObjStorageUserBackupSize[getBackupObjectStorageName(namespace)]
if backupSize <= 0 {
return nil
}
backupRes := resources.NewObjStorageResourceNamed("DB-BACKUP")
if resUsed[backupRes.String()] == nil {
resNamed[backupRes.String()] = backupRes
resUsed[backupRes.String()] = initResources()
}
resUsed[backupRes.String()][corev1.ResourceStorage].Add(*resource.NewQuantity(backupSize, resource.BinarySI))
return nil
}
func (r *MonitorReconciler) monitorServiceResourceUsage(namespace string, resUsed map[string]map[corev1.ResourceName]*quantity, resNamed map[string]*resources.ResourceNamed) error {
svcList := &corev1.ServiceList{}
if err := r.List(context.Background(), svcList, &client.ListOptions{
Namespace: namespace,
FieldSelector: fields.OneTermEqualSelector("spec.type", string(corev1.ServiceTypeNodePort)),
}); err != nil {
return fmt.Errorf("failed to list svc: %v", err)
}
for _, svc := range svcList.Items {
if svc.Spec.Type != corev1.ServiceTypeNodePort || len(svc.Spec.Ports) == 0 {
for i := range svcList.Items {
svc := &svcList.Items[i]
if len(svc.Spec.Ports) == 0 {
continue
}
port := make(map[int32]struct{})
for i := range svc.Spec.Ports {
port[svc.Spec.Ports[i].NodePort] = struct{}{}
for _, svcPort := range svc.Spec.Ports {
port[svcPort.NodePort] = struct{}{}
}
svcRes := resources.NewResourceNamed(&svc)
svcRes := resources.NewResourceNamed(svc)
if resUsed[svcRes.String()] == nil {
resNamed[svcRes.String()] = svcRes
resUsed[svcRes.String()] = initResources()
@@ -368,28 +449,7 @@ func (r *MonitorReconciler) monitorResourceUsage(namespace *corev1.Namespace) er
// nodeport 1:1000, the measurement is quantity 1000
resUsed[svcRes.String()][corev1.ResourceServicesNodePorts].Add(*resource.NewQuantity(int64(1000*len(port)), resource.BinarySI))
}
var monitors []*resources.Monitor
if username := config.GetUserNameByNamespace(namespace.Name); r.ObjStorageClient != nil {
if err := r.getObjStorageUsed(username, &resNamed, &resUsed); err != nil {
r.Logger.Error(err, "failed to get object storage used", "username", username)
}
}
for name, podResource := range resUsed {
isEmpty, used := r.getResourceUsed(podResource)
if isEmpty {
continue
}
monitors = append(monitors, &resources.Monitor{
Category: namespace.Name,
Used: used,
Time: timeStamp,
Type: resNamed[name].Type(),
Name: resNamed[name].Name(),
})
}
return r.DBClient.InsertMonitor(context.Background(), monitors...)
return nil
}
func getBackupObjectStorageName(namespace string) string {
@@ -413,20 +473,21 @@ func (r *MonitorReconciler) getResourceUsed(podResource map[corev1.ResourceName]
return isEmpty, used
}
func (r *MonitorReconciler) getObjStorageUsed(user string, namedMap *map[string]*resources.ResourceNamed, resMap *map[string]map[corev1.ResourceName]*quantity) error {
if r.currentObjectMetrics == nil || r.currentObjectMetrics[user].Usage == nil {
func (r *MonitorReconciler) monitorObjectStorageUsage(namespace string, resMap map[string]map[corev1.ResourceName]*quantity, namedMap map[string]*resources.ResourceNamed) error {
username := config.GetUserNameByNamespace(namespace)
if r.currentObjectMetrics == nil || r.currentObjectMetrics[username].Usage == nil {
return nil
}
for bucket, usage := range r.currentObjectMetrics[user].Usage {
for bucket, usage := range r.currentObjectMetrics[username].Usage {
if bucket == "" || usage <= 0 {
continue
}
objStorageNamed := resources.NewObjStorageResourceNamed(bucket)
(*namedMap)[objStorageNamed.String()] = objStorageNamed
if _, ok := (*resMap)[objStorageNamed.String()]; !ok {
(*resMap)[objStorageNamed.String()] = initResources()
namedMap[objStorageNamed.String()] = objStorageNamed
if _, ok := resMap[objStorageNamed.String()]; !ok {
resMap[objStorageNamed.String()] = initResources()
}
(*resMap)[objStorageNamed.String()][corev1.ResourceStorage].Add(*resource.NewQuantity(usage, resource.BinarySI))
resMap[objStorageNamed.String()][corev1.ResourceStorage].Add(*resource.NewQuantity(usage, resource.BinarySI))
}
return nil
}
@@ -538,8 +599,10 @@ func (r *MonitorReconciler) handlerTrafficUsed(startTime, endTime time.Time, mon
return nil
}
func (r *MonitorReconciler) getGPUResourceUsage(pod corev1.Pod, gpuReq resource.Quantity, rs map[corev1.ResourceName]*quantity) (err error) {
func (r *MonitorReconciler) getGPUResourceUsage(pod *corev1.Pod, gpuReq resource.Quantity, rs map[corev1.ResourceName]*quantity) (err error) {
nodeName := pod.Spec.NodeName
r.gpuMutex.Lock()
defer r.gpuMutex.Unlock()
gpuModel, exist := r.NvidiaGpu[nodeName]
if !exist {
if r.NvidiaGpu, err = gpu.GetNodeGpuModel(r.Client); err != nil {
+7 -4
View File
@@ -102,10 +102,13 @@ func main() {
}
setupLog.Info("starting manager")
//if err := mgr.Start(ctrl.SetupSignalHandler()); err != nil {
// setupLog.Error(err, "problem running manager")
// os.Exit(1)
//}
err = controllers.InitIndexField(mgr)
if err != nil {
setupLog.Error(err, "failed to init index field")
os.Exit(1)
}
go func() {
if err := mgr.Start(ctrl.SetupSignalHandler()); err != nil {
setupLog.Error(err, "problem running manager")