etcd的服务发现使用

package service
 
import (
"context"
"encoding/json"
"errors"
"fmt"
"os"
"regexp"
 
"strconv"
"strings"
"sync"
"time"
 
//这里去掉了一部分
clientv3 "go.etcd.io/etcd/client/v3"
"go.etcd.io/etcd/client/v3/naming/endpoints"
)
 
type SeqProducer interface {
//GetServerAddrs() error
//全量更新节点信息
FullUpdateCache() error
//增量更新节点信息
IncrUpdateCache() error
//节点的数量统计
GenerateNo() (total int, current int)
}
 
type EtcdSeqGenerator struct {
Conf clientv3.Config
Handler *clientv3.Client
// cos/a
manager endpoints.Manager
SrvBasePath string
LocalAddr string
LocalPort int
LocalPid int
AddrCache sync.Map
Ttl int64 //租约时间
regCount int64
rwmutex sync.RWMutex
}
 
func NewEtcdSeqGer(ctx context.Context) (*EtcdSeqGenerator, error) {
endpoint := strings.Split(config.CosServerConfig.EtcdConf.EndPoints, ",")
conf := clientv3.Config{
Endpoints: endpoint,
DialTimeout: 2 * time.Second,
}
cli, err := clientv3.New(conf)
if err != nil {
return nil, err
}
//注意这里如果用容器的时候,获取ip是不固定的,最好在环境变量中安插或者用net.interfaces
serverIp, err := utils.LocalAddr("eth1")
if err != nil {
return nil, err
}
manager, err := endpoints.NewManager(cli, strings.TrimRight(config.CosServerConfig.EtcdConf.SrvBasePath, "/"))
if err != nil {
return nil, err
}
sqlGenerator := &EtcdSeqGenerator{
Conf: conf,
SrvBasePath: strings.TrimRight(config.CosServerConfig.EtcdConf.SrvBasePath, "/"),
Handler: cli,
LocalAddr: serverIp,
LocalPort: config.CosServerConfig.PbSet.Port,
LocalPid: os.Getpid(),
Ttl: config.CosServerConfig.EtcdConf.Ttl,
manager: manager,
//mutex: ,
//Val: ,
}
return sqlGenerator, nil
}
 
//refresh etcd hanlder
func (esg *EtcdSeqGenerator) refreshHandler() error {
esg.rwmutex.Lock()
defer esg.rwmutex.Unlock()
esg.regCount++
esg.Handler.Close()
esg.Handler = nil
cli, err := clientv3.New(esg.Conf)
if err != nil {
return err
}
esg.Handler = cli
manager, err := endpoints.NewManager(esg.Handler, esg.SrvBasePath+"/")
esg.manager = manager
return nil
}
 
//service register
func (esg *EtcdSeqGenerator) Register() (<-chan *clientv3.LeaseKeepAliveResponse, error) {
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
if esg.regCount > 0 {
if err := esg.refreshHandler(); err != nil {
return nil, err
}
}
leaseGrantResp, err := esg.Handler.Grant(ctx, esg.Ttl)
if err != nil {
return nil, err
}
 
key := fmt.Sprintf("%s/%s:%d:%d", esg.SrvBasePath, esg.LocalAddr, esg.LocalPort, esg.LocalPid)
update := endpoints.NewAddUpdateOpts(key, endpoints.Endpoint{Addr: fmt.Sprintf("%s:%d:%d", esg.LocalAddr, esg.LocalPort, esg.LocalPid)}, clientv3.WithLease(leaseGrantResp.ID))
if err := esg.manager.Update(context.Background(), []*endpoints.UpdateWithOpts{update}); err != nil {
return nil, err
}
 
return esg.Handler.KeepAlive(context.TODO(), leaseGrantResp.ID)
 
}
 
//get changes
func (esg *EtcdSeqGenerator) readUpdChange() (endpoints.WatchChannel, error) {
esg.rwmutex.RLock()
defer esg.rwmutex.RUnlock()
updkv, err := esg.manager.NewWatchChannel(context.TODO())
if err != nil {
return nil, err
}
return updkv, nil
}
 
func (esg *EtcdSeqGenerator) IncrUpdateCache() error {
updKv, err := esg.readUpdChange()
if err != nil {
return err
}
for {
select {
case updv, ok := <-updKv:
if !ok {
return fmt.Errorf("updKv quit unexpected")
}
for _, upd := range updv {
val, err := json.Marshal(upd.Endpoint)
if err != nil {
//log
}
switch upd.Op {
case endpoints.Add:
esg.AddrCache.Store(upd.Key, val)
case endpoints.Delete:
esg.AddrCache.Delete(upd.Key)
}
}
case <-time.After(2 * time.Second):
}
}
 
}
 
func (esg *EtcdSeqGenerator) readResp() (*clientv3.GetResponse, error) {
esg.rwmutex.RLock()
defer esg.rwmutex.RUnlock()
resp, err := esg.Handler.Get(context.Background(), esg.SrvBasePath, clientv3.WithPrefix())
if err != nil {
return nil, err
}
return resp, nil
}
 
//
func (esg *EtcdSeqGenerator) FullUpdateCache() error {
resp, err := esg.readResp()
if err != nil {
return err
}
esg.AddrCache.Range(func(k, v interface{}) bool {
delFlag := true
for _, kv := range resp.Kvs {
if _, ok := esg.AddrCache.Load(string(kv.Key)); ok {
/*
var val endpoints.Update
if err := json.Unmarshal(kv.Value, &val); err != nil {
//log
continue
}
*/
esg.AddrCache.Store(string(kv.Key), kv.Value)
delFlag = false
break
}
}
if delFlag {
esg.AddrCache.Delete(k.(string))
}
return true
})
 
return nil
}
func (esg *EtcdSeqGenerator) GenerateNo() (total int, current int, err error) {
 
var keys []string
esg.AddrCache.Range(
func(k, v interface{}) bool {
keys = append(keys, k.(string))
return true
})
 
myAddrNumber := utils.InetAtoN(esg.LocalAddr)
myPart := 0
reg := regexp.MustCompile("^[a-zA-Z_]+/" + cst.IpRegx + "(:[0-9]+){2}")
for _, key := range keys {
if !reg.MatchString(key) {
return total, current, fmt.Errorf("key:%s Format requirements are not met,need exam:cos/xx:xx:xx", key)
}
arr := strings.Split(key, "/")
uniqStrs := strings.Split(arr[len(arr)-1], ":")
port, err := strconv.Atoi(uniqStrs[1])
if err != nil {
return total, current, err
}
pid, err := strconv.Atoi(uniqStrs[2])
if err != nil {
return total, current, err
}
if myAddrNumber == utils.InetAtoN(uniqStrs[0]) && port == esg.LocalPort {
if pid != esg.LocalPid {
return total, current, fmt.Errorf("dulicate port for xxx-server")
}
myPart++
}
//secondVal := strings.Split("_", instance)[1]
if myAddrNumber > utils.InetAtoN(uniqStrs[0]) || myAddrNumber == utils.InetAtoN(uniqStrs[0]) && port <= esg.LocalPort {
current++
}
}
if myPart == 0 {
return total, current, errors.New("I'm not a partner")
}
return len(keys), current - 1, nil
}



请使用浏览器的分享功能分享到微信等