From b109ffe2f10db323938629dce20cd4a4fe0c0488 Mon Sep 17 00:00:00 2001 From: zijiren <84728412+zijiren233@users.noreply.github.com> Date: Tue, 2 Jul 2024 16:21:20 +0800 Subject: [PATCH] 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 --- .../controllers/monitor_controller.go | 189 ++++++++++++------ controllers/resources/main.go | 11 +- 2 files changed, 133 insertions(+), 67 deletions(-) diff --git a/controllers/resources/controllers/monitor_controller.go b/controllers/resources/controllers/monitor_controller.go index ea79af6dd..456bb007e 100644 --- a/controllers/resources/controllers/monitor_controller.go +++ b/controllers/resources/controllers/monitor_controller.go @@ -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 { diff --git a/controllers/resources/main.go b/controllers/resources/main.go index a32a2771f..fe1cec343 100644 --- a/controllers/resources/main.go +++ b/controllers/resources/main.go @@ -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")