From eba3014c59e0fefdfb2ef2fa3f6da793a0e28507 Mon Sep 17 00:00:00 2001 From: ning <710leo@gmail.com> Date: Thu, 9 May 2024 16:48:29 +0800 Subject: [PATCH] fix: alert engine rebuild hash --- alert/naming/hashring.go | 6 ++++++ alert/naming/heartbeat.go | 9 +++++++++ 2 files changed, 15 insertions(+) diff --git a/alert/naming/hashring.go b/alert/naming/hashring.go index f88dd21b..d1a3220f 100644 --- a/alert/naming/hashring.go +++ b/alert/naming/hashring.go @@ -67,6 +67,12 @@ func (chr *DatasourceHashRingType) Set(datasourceId string, r *consistent.Consis chr.Rings[datasourceId] = r } +func (chr *DatasourceHashRingType) Del(datasourceId string) { + chr.Lock() + defer chr.Unlock() + delete(chr.Rings, datasourceId) +} + func (chr *DatasourceHashRingType) Clear(engineName string) { chr.Lock() defer chr.Unlock() diff --git a/alert/naming/heartbeat.go b/alert/naming/heartbeat.go index 32387fa6..c2ecea6e 100644 --- a/alert/naming/heartbeat.go +++ b/alert/naming/heartbeat.go @@ -110,7 +110,9 @@ func (n *Naming) heartbeat() error { } } + newDatasource := make(map[int64]struct{}) for i := 0; i < len(datasourceIds); i++ { + newDatasource[datasourceIds[i]] = struct{}{} servers, err := n.ActiveServers(datasourceIds[i]) if err != nil { logger.Warningf("hearbeat %d get active server err:%v", datasourceIds[i], err) @@ -130,6 +132,13 @@ func (n *Naming) heartbeat() error { localss[datasourceIds[i]] = newss } + for dsId := range localss { + if _, exists := newDatasource[dsId]; !exists { + delete(localss, dsId) + DatasourceHashRing.Del(fmt.Sprintf("%d", dsId)) + } + } + // host 告警使用的是 hash ring err = models.AlertingEngineHeartbeatWithCluster(n.ctx, n.heartbeatConfig.Endpoint, n.heartbeatConfig.EngineName, HostDatasource) if err != nil {