From 112bfee1c46a183beb4942f3a459aacf33d77d09 Mon Sep 17 00:00:00 2001
From: chenshijun <csj_sky@126.com>
Date: 星期五, 06 九月 2019 17:46:54 +0800
Subject: [PATCH] 服务器只能到服务器去拉取数据
---
agent.go | 122 +++++++++++++++++++++++++++-------------
1 files changed, 82 insertions(+), 40 deletions(-)
diff --git a/agent.go b/agent.go
index bc2585f..125f421 100644
--- a/agent.go
+++ b/agent.go
@@ -25,6 +25,7 @@
"io/ioutil"
"net"
"os"
+ "strconv"
"sync"
//"os"
@@ -40,6 +41,8 @@
const (
QueryEventGetDB = "GetDatabase"
QueryEventUpdateDBData = "UpdateDBData"
+ UserEventSyncSql = "SyncSql"
+ UserEventSyncDbTablePersonCache = "SyncCache"
)
// Agent warps the serf agent
@@ -110,6 +113,8 @@
go a.BroadcastMemberlist(BroadcastInterval * time.Second)
}
+var SyncDbTablePersonCacheChan = make(chan []byte,0)
+
// HandleEvent Handles serf.EventMemberJoin events,
// which will wait for members to join until the number of group members is equal to "groupExpect"
// when the startup mode is "ModeCluster",
@@ -118,18 +123,22 @@
switch ev := event.(type) {
case serf.UserEvent:
- //fmt.Println(string(ev.Payload))
- var sqlUe SqlUserEvent
- err := json.Unmarshal(ev.Payload, &sqlUe)
- if err !=nil {
- fmt.Println("sqlUe unmarshal err:",err)
- return
+ if ev.Name == UserEventSyncSql {
+ var sqlUe SqlUserEvent
+ err := json.Unmarshal(ev.Payload, &sqlUe)
+ if err !=nil {
+ fmt.Println("sqlUe unmarshal err:",err)
+ return
+ }
+ if sqlUe.Owner != a.conf.NodeName {
+ //results, err := ExecuteWriteSql(sqlArr)
+ flag, _ := ExecuteSqlByGorm(sqlUe.Sql)
+ fmt.Println("userEvent exec ",sqlUe.Sql,",Result:",flag)
+ }
+ } else if ev.Name == UserEventSyncDbTablePersonCache {
+ SyncDbTablePersonCacheChan <- ev.Payload
}
- if sqlUe.Owner != a.conf.NodeName {
- //results, err := ExecuteWriteSql(sqlArr)
- flag, _ := ExecuteSqlByGorm(sqlUe.Sql)
- fmt.Println("userEvent exec ",sqlUe.Sql,",Result:",flag)
- }
+
case *serf.Query:
@@ -200,6 +209,18 @@
//var res []*Rows
//json.Unmarshal(bytesReturn, &res)
}
+ case serf.MemberEvent:
+ if event.EventType() == serf.EventMemberLeave {
+ if ev.Members !=nil && len(ev.Members) ==1 {
+ leaveMember := ev.Members[0]
+ leaveSql := "delete from cluster_node where node_id='"+leaveMember.Name+"'"
+ ExecuteSqlByGorm([]string{ leaveSql })
+
+ fmt.Println("EventMemberLeave,current Members:",ev.Members)
+ }
+ return
+ }
+
default:
fmt.Printf("Unknown event type: %s\n", ev.EventType().String())
@@ -428,9 +449,16 @@
var specmembername string
for _, m := range mbs {
fmt.Println("m",m)
- if m.Name != a.conf.NodeName {
- specmembername = m.Name
- break
+ if m.Name != a.conf.NodeName { //鍓嶇紑锛欴SVAD:鍒嗘瀽鏈嶅姟鍣� DSPAD:杩涘嚭鍏ad
+ if strings.HasPrefix(a.conf.NodeName, "DSVAD"){
+ if strings.HasPrefix(m.Name, "DSVAD") {
+ specmembername = m.Name
+ break
+ }
+ }else{
+ specmembername = m.Name
+ break
+ }
}
}
fmt.Println("mbs:",mbs,"a.conf.BindAddr:",a.conf.BindAddr,"specmembername:",specmembername)
@@ -501,21 +529,29 @@
fmt.Println("sqlUE marshal err:",err)
return
}
- err = a.UserEvent("SyncSql", ueB, false)
+ err = a.UserEvent(UserEventSyncSql, ueB, false)
if err == nil || !strings.Contains(err.Error(), "cannot contain") {
fmt.Println("err: ", err)
}
}
+//鏇存柊鍚屾搴撶殑姣斿缂撳瓨
+func (a *Agent) SyncDbTablePersonCache(b []byte) {
+ err := a.UserEvent(UserEventSyncDbTablePersonCache, b, false)
+ if err !=nil{
+ fmt.Println("UserEventSyncDbTablePersonCache err:",err)
+ }
+}
+
//Init serf Init
-func Init(clusterID string, password string, nodeID string, ips []string) (*Agent, error) {
+func Init(clusterID string, password string, nodeID string, addrs []string) (*Agent, error) {
agent, err := InitNode(clusterID, password, nodeID)
if err != nil {
fmt.Printf("InitNode failed, error: %s", err)
return agent, err
}
- err = agent.JoinByNodeIP(ips)
+ err = agent.JoinByNodeAddrs(addrs)
if err != nil {
fmt.Printf("JoinByNodeIP failed, error: %s", err)
return agent, err
@@ -560,43 +596,49 @@
return agent, nil
}
-func (a *Agent) JoinByNodeIP(ips []string) error {
+func (a *Agent) JoinByNodeAddrs(addrs []string) error {
var nodes []string
- if len(ips) == 0 {
+ if len(addrs) == 0 {
return fmt.Errorf("No Nodes To Join!")
}
- for _, ip := range ips {
- node := fmt.Sprintf("%s:%d", ip, DefaultBindPort)
- nodes = append(nodes, node)
+ for _, addr := range addrs {
+ nodes = append(nodes, addr)
}
- n, err := a.Agent.Join(nodes, true)
- if err != nil || n == 0 {
- //a.Stop()
- //fmt.Println("Stop node")
- return fmt.Errorf("Error Encrypt Key!")
- }
+ a.Agent.Join(nodes, true)
- return err
+ return nil
}
-type Node struct {
- clusterID string
- NodeID string
- IP string
- isAlive int //StatusNone:0, StatusAlive:1, StatusLeaving:2, StatusLeft:3, StatusFailed:4
-}
+//func (a *Agent) JoinByNodeIP(ips []string) error {
+// var nodes []string
+//
+// if len(ips) == 0 {
+// return fmt.Errorf("No Nodes To Join!")
+// }
+// for _, ip := range ips {
+// node := fmt.Sprintf("%s:%d", ip, DefaultBindPort)
+// nodes = append(nodes, node)
+// }
+//
+// n, err := a.Agent.Join(nodes, true)
+// if err != nil || n == 0 {
+// return fmt.Errorf("Error Encrypt Key!")
+// }
+//
+// return err
+//}
-func (a *Agent) GetNodes() (nodes []Node) {
- var node Node
+func (a *Agent) GetNodes() (nodes []NodeInfo) {
+ var node NodeInfo
fmt.Println("a.conf.ClusterID:", a.conf.ClusterID)
mbs := a.GroupMembers(a.conf.ClusterID)
for _, mb := range mbs {
node.NodeID = mb.Name
- node.IP = mb.Addr.String()
- node.isAlive = int(mb.Status)
- node.clusterID = mb.Tags[tagKeyClusterID]
+ node.NodeAddress = mb.Addr.String() + ":" + strconv.Itoa(int(mb.Port))
+ node.IsAlive = int(mb.Status)
+ node.ClusterID = mb.Tags[tagKeyClusterID]
nodes = append(nodes, node)
}
--
Gitblit v1.8.0