This is an automated email from the ASF dual-hosted git repository.

zhongxjian pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/dubbo-kubernetes.git


The following commit(s) were added to refs/heads/master by this push:
     new 8225b73b [horus] Update reboot recovery (#448)
8225b73b is described below

commit 8225b73b136811719201378fe86c4cb937d86466
Author: mfordjody <[email protected]>
AuthorDate: Fri Oct 11 12:04:41 2024 +0800

    [horus] Update reboot recovery (#448)
---
 app/horus/base/db/db.go                 |  2 +-
 app/horus/core/horuser/node_downtime.go | 20 ++++++++----
 app/horus/core/horuser/node_recovery.go | 58 +++++++++++++++++++--------------
 app/horus/core/horuser/node_restart.go  |  4 +--
 manifests/horus/horus.yaml              | 12 +++----
 5 files changed, 56 insertions(+), 40 deletions(-)

diff --git a/app/horus/base/db/db.go b/app/horus/base/db/db.go
index b1afe246..0408b0ac 100644
--- a/app/horus/base/db/db.go
+++ b/app/horus/base/db/db.go
@@ -41,7 +41,7 @@ type NodeDataInfo struct {
        RecoveryMark         int64     `json:"recovery_mark" 
xorm:"recovery_mark"`
        RecoveryQL           string    `json:"recovery_ql" xorm:"recovery_ql"`
        DownTimeRecoveryMark int64     `json:"downtime_recovery_mark" 
xorm:"downtime_recovery_mark"`
-       DownTimeRecoveryQL   string    `json:"downtime_recovery_ql" 
xorm:"downtime_recovery_ql"`
+       DownTimeRecoveryQL   []string  `json:"downtime_recovery_ql" 
xorm:"downtime_recovery_ql"`
 }
 
 type PodDataInfo struct {
diff --git a/app/horus/core/horuser/node_downtime.go 
b/app/horus/core/horuser/node_downtime.go
index 9e523bc4..9883f5f9 100644
--- a/app/horus/core/horuser/node_downtime.go
+++ b/app/horus/core/horuser/node_downtime.go
@@ -91,18 +91,20 @@ func (h *Horuser) DownTimeNodes(clusterName, addr string) {
                                continue
                        }
                        nodeDownTimeRes[nodeName]++
+
                }
        }
 
        WithDownNodeIPs := make(map[string]string)
 
-       for nodeName, count := range nodeDownTimeRes {
+       for node, count := range nodeDownTimeRes {
+               klog.Infof("Count:%v", count)
                if count < aq {
                        klog.Error("downtimeNodes not reach threshold.")
-                       klog.Infof("clusterName:%v nodeName:%v threshold:%v 
count:%v", clusterName, nodeName, aq, count)
+                       klog.Infof("clusterName:%v nodeName:%v threshold:%v 
count:%v", clusterName, node, aq, count)
                        continue
                }
-               abnormalInfoSystemQL := 
fmt.Sprintf(h.cc.NodeDownTime.AbnormalInfoSystemQL, nodeName)
+               abnormalInfoSystemQL := 
fmt.Sprintf(h.cc.NodeDownTime.AbnormalInfoSystemQL, node)
 
                res, err := h.InstantQuery(addr, abnormalInfoSystemQL, 
clusterName, h.cc.NodeDownTime.PromQueryTimeSecond)
                if len(res) == 0 {
@@ -114,11 +116,11 @@ func (h *Horuser) DownTimeNodes(clusterName, addr string) 
{
                        klog.Infof("clusterName:%v\n AbnormalInfoSystemQL:%v, 
err:%v", clusterName, abnormalInfoSystemQL, err)
                        continue
                }
-               str := ""
+               instanceIP := ""
                for _, v := range res {
-                       str = string(v.Metric["instance"])
+                       instanceIP += string(v.Metric["instance"])
                }
-               WithDownNodeIPs[nodeName] = str
+               WithDownNodeIPs[node] = instanceIP
        }
 
        msg := fmt.Sprintf("\n【%s】\n【集群:%v】\n【已达到宕机临界点:%v】", 
h.cc.NodeDownTime.DingTalk.Title, clusterName, len(WithDownNodeIPs))
@@ -147,7 +149,11 @@ func (h *Horuser) DownTimeNodes(clusterName, addr string) {
                }()
 
                moduleName := 0
-               abnormalRecoveryQL := 
fmt.Sprintf(h.cc.NodeDownTime.AbnormalRecoveryQL[moduleName], nodeName)
+
+               abnormalRecoveryQL := []string{
+                       
fmt.Sprintf(h.cc.NodeDownTime.AbnormalRecoveryQL[moduleName], nodeName),
+               }
+
                write := db.NodeDataInfo{
                        NodeName:           nodeName,
                        NodeIP:             nodeIP,
diff --git a/app/horus/core/horuser/node_recovery.go 
b/app/horus/core/horuser/node_recovery.go
index 9c82dd02..614da50b 100644
--- a/app/horus/core/horuser/node_recovery.go
+++ b/app/horus/core/horuser/node_recovery.go
@@ -129,41 +129,51 @@ func (h *Horuser) downTimeRecoveryNodes(n 
db.NodeDataInfo) {
                return
        }
        rq := len(h.cc.NodeDownTime.AbnormalRecoveryQL)
-       vecs, err := h.InstantQuery(promAddr, n.DownTimeRecoveryQL, 
n.ClusterName, h.cc.NodeRecovery.PromQueryTimeSecond)
-       if err != nil {
-               klog.Errorf("downTimeRecoveryNodes InstantQuery err:%v", err)
-               klog.Infof("DownTimeRecoveryQL:%v", n.DownTimeRecoveryQL)
-               return
-       }
-       if len(vecs) != 1 {
-               klog.Infof("Expected 1 result, but got:%d", len(vecs))
-               return
-       }
-       if len(vecs) > rq {
-               klog.Error("downTimeRecoveryNodes not reach threshold")
-       }
-       if err != nil {
-               klog.Errorf("downTimeRecoveryNodes InstantQuery err:%v", err)
-               klog.Infof("DownTimeRecoveryQL:%v", n.DownTimeRecoveryQL)
-               return
-       }
-       klog.Info("recoveryNodes InstantQuery success.")
+       for _, ql := range n.DownTimeRecoveryQL { // 遍历每个 PromQL 查询
+               vecs, err := h.InstantQuery(promAddr, ql, n.ClusterName, 
h.cc.NodeRecovery.PromQueryTimeSecond)
+               if err != nil {
+                       klog.Errorf("downTimeRecoveryNodes InstantQuery err: 
%v", err)
+                       klog.Infof("DownTimeRecoveryQL: %v", ql) // 记录出错的查询
+                       continue                                 // 继续处理下一个查询
+               }
+
+               // 处理查询结果
+               if len(vecs) == 0 {
+                       klog.Infof("No results for query: %v", ql)
+                       continue
+               }
+
+               // 在这里处理查询成功的逻辑,比如检查结果并采取相应的操作
+               klog.Infof("Query successful for: %v", ql)
+
+               if len(vecs) != rq {
+                       klog.Infof("Expected %d results, but got: %d", rq, 
len(vecs))
+                       if len(vecs) > rq {
+                               klog.Error("downTimeRecoveryNodes did not reach 
threshold")
+                       }
+                       return
+               }
 
-       if len(vecs) == rq {
                err = h.UnCordon(n.NodeName, n.ClusterName)
                res := "Success"
                if err != nil {
-                       res = fmt.Sprintf("result failed:%v", err)
+                       res = fmt.Sprintf("result failed: %v", err)
                }
-               msg := fmt.Sprintf("\n【集群: %v】\n【封锁节点恢复调度】\n【已恢复调度节点: 
%v】\n【处理结果:%v】\n【日期: %v】\n", n.ClusterName, n.NodeName, res, n.CreateTime)
+
+               msg := fmt.Sprintf("\n【集群: %v】\n【封锁节点恢复调度】\n【已恢复调度节点: 
%v】\n【处理结果:%v】\n【日期: %v】\n",
+                       n.ClusterName, n.NodeName, res, n.CreateTime)
+
                alerter.DingTalkSend(h.cc.NodeDownTime.DingTalk, msg)
                alerter.SlackSend(h.cc.NodeDownTime.Slack, msg)
 
                success, err := n.DownTimeRecoveryMarker()
                if err != nil {
-                       klog.Errorf("DownTimeRecoveryMarker result failed 
err:%v", err)
+                       klog.Errorf("DownTimeRecoveryMarker result failed: %v", 
err)
                        return
                }
-               klog.Infof("DownTimeRecoveryMarker result success:%v", success)
+               klog.Infof("DownTimeRecoveryMarker result success: %v", success)
+
+               // 查询操作成功日志
+               klog.Info("recoveryNodes InstantQuery success.")
        }
 }
diff --git a/app/horus/core/horuser/node_restart.go 
b/app/horus/core/horuser/node_restart.go
index db6adcd9..5c81ed96 100644
--- a/app/horus/core/horuser/node_restart.go
+++ b/app/horus/core/horuser/node_restart.go
@@ -82,8 +82,8 @@ func (h *Horuser) TryRestart(node db.NodeDataInfo) {
                klog.Infof("Successfully restarted node %v.", node.NodeName)
 
                rq := len(h.cc.NodeDownTime.AbnormalRecoveryQL)
-               query := strings.Join(h.cc.NodeDownTime.AbnormalRecoveryQL, " 
and ")
-               vecs, err := 
h.InstantQuery(h.cc.PromMultiple[node.ClusterName], query, node.ClusterName, 
h.cc.NodeDownTime.PromQueryTimeSecond)
+               ql := strings.Join(h.cc.NodeDownTime.AbnormalRecoveryQL, " or ")
+               vecs, err := 
h.InstantQuery(h.cc.PromMultiple[node.ClusterName], ql, node.ClusterName, 
h.cc.NodeDownTime.PromQueryTimeSecond)
                if err != nil {
                        klog.Errorf("Failed to query Prometheus for recovery 
threshold after restart: %v", err)
                        return
diff --git a/manifests/horus/horus.yaml b/manifests/horus/horus.yaml
index 8c362540..21196cda 100644
--- a/manifests/horus/horus.yaml
+++ b/manifests/horus/horus.yaml
@@ -68,17 +68,17 @@ nodeDownTime:
   intervalSecond: 15
   promQueryTimeSecond: 60
   abnormalityQL:
-    - 100 - (avg by (node) (rate(node_cpu_seconds_total{mode="idle"}[5m])) * 
100) > 20
-    - (avg by (node) (node_memory_MemFree_bytes / node_memory_MemTotal_bytes 
)) * 100 > 50
-#    - node_filesystem_avail_bytes{mountpoint="/"} / 
node_filesystem_size_bytes{mountpoint="/"} * 100 < 15
+    - 100 - (avg by (node) (rate(node_cpu_seconds_total{mode="idle"}[5m])) * 
100) > 15
+#    - avg(node_memory_MemFree_bytes) by (node) < 20
+    - node_filesystem_avail_bytes{mountpoint="/"} / 
node_filesystem_size_bytes{mountpoint="/"} * 100 < 16
   abnormalInfoSystemQL:
     node_os_info{node="%s"}
   allSystemUser: "zxj"
   allSystemPassword: "1"
   abnormalRecoveryQL:
-    - 100 - (avg by (node) 
(rate(node_cpu_seconds_total{mode="idle",node="%s"}[5m])) * 100) < 20
-    - (avg by (node) (node_memory_MemFree_bytes{node="%s"} / 
node_memory_MemTotal_bytes{node="%s"} )) * 100 < 30
-    # - node_filesystem_avail_bytes{mountpoint="/"} / 
node_filesystem_size_bytes{mountpoint="/"} * 100 > 15
+    - 100 - (avg by (node) 
(rate(node_cpu_seconds_total{mode="idle",node="%s"}[5m])) * 100) < 15
+#    - (avg by (node) (node_memory_MemFree_bytes{node="%s"} / 
node_memory_MemTotal_bytes{node="%s"} )) * 100 > 50
+    - node_filesystem_avail_bytes{mountpoint="/",node="%s"} / 
node_filesystem_size_bytes{mountpoint="/",node="%s"} * 100 > 15
   kubeMultiple:
     cluster: config.1
   dingTalk:

Reply via email to