Go 服務註冊與發現筆記

概述

朋友們好啊,這篇筆記我們圍繞 Go 來記錄一下服務註冊與發現概述與流程,如註冊中心的機制 (心跳、時間間隔)、gRPC 接入註冊中心 (resolver 實現、etcd 租約、消費者服務發現等等),以etcd作爲註冊中心將上篇user_service grpc 服務註冊、自動續約,以 gin 作爲客戶端調用服務。

服務註冊與發現概述

在分佈式架構中爲什麼需要服務註冊與發現?服務註冊與發現是誰在註冊,誰要發現?在分佈式系統中 (微服務是分佈式的其中一種) 一般都會部署多個實例,如果多個實例需要更新或故障移除時,客戶端就沒辦法及時知道可用的實例在哪裏,如當某個實例下線時會通知註冊中心,避免將請求再發到該實例,新實例上線並註冊到註冊中心;同時服務註冊發現爲負載均衡、容錯故障轉移提供了支持,根據一定的策略 (如少連接數、加權、負載情況) 將請求分發到不同的實例上;當某個實例發生故障時將請求轉移到健康的實例上。實際上服務註冊與發現要解決 "我要知道哪些機器部署了哪些服務",服務註冊的主體是服務提供者(服務端,只有服務端需要註冊自身地址),客戶端的任務是 "發現" 已註冊的業務,而非註冊自己。

註冊中心概述

服務註冊的組成部分

    在服務註冊中分成 3 個關鍵部分,分別是註冊中心 (存儲和管理所有服務端信息)、服務端(服務提供方,具體提供服務的實例) 和客戶端(調用其他服務的應用):註冊中心存儲服務端以下信息:

  1. 服務基本信息:包括服務名稱、版本、描述;

  2. 服務位置信息:包括服務地址 (IP 地址)、服務端口;

  3. 服務健康信息:包括心跳信號和健康檢查 URL;

  4. 服務元數據:包括負載均衡、路由、標籤等等。

    其中健康檢查時間間隔越短,對性能要求越高、但越快發現實例故障,但不好處理偶發性心跳失敗;反過來健康檢查時間間隔長越晚,優化方案可以是單次心跳失敗後,再次主動發起心跳檢測,避免將健康的實例錯當故障實例。

    然後服務端來說按照上面顯示註冊中心存儲的信息就知道要提供哪些信息給註冊中心了;客戶端的主要職責是服務發現,定期向註冊中心查詢可用的服務端列表、與註冊中心同步保持數據一致性,緩存服務端實例信息(地址、端口等),並通過負載均衡策略,動態調整請求流向。

服務註冊的流程

  1. 首先是服務端啓動起來,正常啓動後才能向註冊中心註冊;如果未啓動成功就不向註冊中心註冊了。

  2. 向註冊中心發起註冊請求,將服務名稱,網絡地址等信息。

  3. 註冊中心通過驗證後將服務端存儲在註冊表中。

  4. 服務端定期發送心跳消息 (或雙向心跳) 證明實例可用;若心跳失敗註冊中心將該實例移除,客戶端定期查詢或註冊中心同步實例列表。

  5. 若服務端某實例即將下線時,想註冊中心主動發送下線請求。該實例不再接受新的請求,讀取了一半的請求會在讀完後直接拒絕,不再交給後端業務代碼執行(可用中間件處理)。

  6. 客戶端啓動時向註冊中心請求目標服務的實例列表。

  7. 客戶端獲取列表後將其緩存避免頻繁訪問註冊中心。

  8. 客戶端向服務端發起調用。

![](https://mmbiz.qpic.cn/mmbiz_png/wVDe9U5Klicnic0DgQPT89Xiaj50m5XNF1AeSQhb7MictvgJASAsYgOhSzkoHiaJGDeB1Yga2DibXibIQYZPFlY8HsGeQ/640?wx_fmt=png&from=appmsg)

好的朋友們大概的理論部分說完了,下一部分我們將上一篇user_service例子改進使其接入基於etcd的註冊中心;常見的註冊中心有如 etcd、Consul、Nacos、EureKa、Zookeeper 等。

gRPC 接入註冊中心

我們要將上次user_service接入到註冊中心,那先將註冊中心 etcd 用 docker 啓動起來,再配置 gRPC 接入註冊中心

啓動 ETCD 配置中心

docker-compose.yaml 文件如下,使用docker-compose up -d啓動。

version: "3"
services:
etcd:
    image:bitnami/etcd:3.6.1
    ports:
      -"2379:2379"
    environment:
      -ALLOW_NONE_AUTHENTICATION=yes
      -ETCD_ADVERTISE_CLIENT_URLS=http://etcd:2379,http://localhost:2379
    volumes:
      -etcd_data:/bitnami/etcd

volumes:
  etcd_data:

實現簡易服務註冊功能

  1. 首先是修改一下 App struct, 要加入一個 etcdClient.
type App struct {
    grpcServer *grpc.Server
    gwServer   *http.Server
    userSvc    *UserService
    healthSvc  grpc_health_v1.HealthServer
    etcdClient *serviceReg.EtcdClient
}
  1. etcdClient 是怎麼組成的呢?實際上包含一個 etcd 客戶端用來連接操作 etcd,Config 保存 etcd 連接的配置,leaseID(記錄服務註冊時創建的租約 id, 該 id 由 etcd 服務器全局唯一生成,用 64 位整數表示在整個 etcd 中是唯一的),serviceKey(註冊到 etcd 使用的鍵名,這個也必須唯一,重複會導致覆蓋;若多實例部署可用 ip + 端口 + uuid 組合的方式)
type EtcdClient struct {
    Client     *clientv3.Client
    Config     *EtcdConfig
    leaseID    clientv3.LeaseID
    serviceKey string
}

func NewEtcdClient(config *EtcdConfig) (*EtcdClient, error) {
    cli, err := clientv3.New(clientv3.Config{
        Endpoints:   config.Endpoints,
        DialTimeout: config.DialTimeout,
    })
    if err != nil {
        returnnil, errors.New("create etcd client error:" + err.Error())
    }
    return &EtcdClient{Client: cli, Config: config}, nil
}
  1. 服務註冊、創建租約,心跳保持與自動續約
    通過 etcd 客戶端調用Grant方法創建租約,contxt.WithTimeout會創建可取消的上下文,在調用完畢後調用cancel()函數釋放資源。
    在函數 keepAlive() 中,通過調用 etcd 客戶端的KeepAlive函數,帶上 leaseID 實現自動續約。該函數會按照租約 1/3 的時間向 etcd 請求續約,租約會重置爲最初的 TTL(這裏設置 10s,可以看到差不多每 3s 就會請求一次續約)
func (e *EtcdClient) RegisterService(serviceName, serviceAddr string) error {
    // 創建租約
    ctx, cancel := context.WithTimeout(context.Background(), e.Config.DialTimeout)
    resp, err := e.Client.Grant(ctx, e.Config.LeaseTTL)
    cancel()

    if err != nil {
        return errors.New("create lease error" + err.Error())
    }

    // 註冊服務
    serviceKey := fmt.Sprintf("services/%s/%s", serviceName, serviceAddr)
    ctx, cancel = context.WithTimeout(context.Background(), e.Config.DialTimeout)
    _, err = e.Client.Put(ctx, serviceKey, serviceAddr, clientv3.WithLease(resp.ID))
    cancel()
    if err != nil {
        return errors.New("register service error:" + err.Error())
    }

    // 保存租約ID和 serviceKey
    e.leaseID = resp.ID
    e.serviceKey = serviceKey

    // 保持心跳
    go e.keepAlive()
    log.Println("serviceKey", serviceKey)
    log.Printf("register service %s %s success,TTL: %d", serviceName, serviceAddr, e.Config.LeaseTTL)
    returnnil
}

// 保持心跳 & 自動續約
func (e *EtcdClient) keepAlive() {
    ctx := context.Background()
    al, err := e.Client.KeepAlive(ctx, e.leaseID)
    if err != nil {
        log.Println("start keep alive error:" + err.Error())
        return
    }

    for {
        select {
        case k, ok := <-al:
            if !ok {
                log.Println("keep alive closed,重新註冊服務")
                return
            }
            log.Printf("keep alive TTL %d", k.TTL)
        }
    }
}

此時我們使用 etcd 的連接工具etcdctl,輸入serviceKey就能看到該實例已經註冊進去了。

實現簡易服務下線功能

  1. 撤銷租約的關鍵是Revoke函數,向 etcd 發送撤銷租約請求,etcd 收到請求後會立即停止租約的自動續約,並且與租約關聯的鍵會被刪除。

  2. 下面Delete函數是直接刪除 etcd 裏的鍵,屬於手動刪除鍵確保服務下線的方法。

func (e *EtcdClient) UnregisterService() error {
    // todo 1. 取消租約
    if e.leaseID != 0 {
        ctx, cancel := context.WithTimeout(context.Background(), e.Config.DialTimeout)
        _, err := e.Client.Revoke(ctx, e.leaseID)
        cancel()
        if err != nil {
            return errors.New("revoke lease error:" + err.Error())
        }
    }

    // todo 2. 刪除實例
    ctx, cancel := context.WithTimeout(context.Background(), e.Config.DialTimeout)
    _, err := e.Client.Delete(ctx, e.serviceKey)
    cancel()
    if err != nil {
        return errors.New("delete service error:" + err.Error())
    }

    log.Println("unregister service success")
    returnnil
}

測試簡易服務註冊與下線

Start函數這裏加一個定時器,到時間就將服務下線

    // 10s後服務下線
    t1 := time.NewTimer(10 * time.Second)
    go func() {
        <-t1.C
        log.Println("時間到嘍")
        if err := a.etcdClient.UnregisterService(); err != nil {
            log.Println("unregister service error:", err)
        } else {
            log.Println("unregister service success")
        }
        // 退出程序
        //os.Exit(0)
    }()

視頻詳情

實現簡易 gRPC 服務發現

如上面所說客戶端的職責是從註冊中心etcd中獲取服務列表,監聽變化更新地址,動態修改 gRPC 調用地址。

  1. 首先我們要有一個 gRPC 解析器 (resolver),裏面包括了下面這些內容:
```go
type etcdResolver struct {
    target     resolver.Target
    conn         resolver.ClientConn
    etcdClient *clientv3.Client
    ctx        context.Context
    cancel     context.CancelFunc
}

```
  1. 定義好解析器,就要一個方法來創建該解析器:
```go
func newEtcdResolver(target resolver.Target, conn resolver.ClientConn, etcdClient *clientv3.Client) resolver.Resolver {
    ctx, cancel := context.WithCancel(context.Background())
    r := &etcdResolver{
        target:     target,
        conn:       conn,
        etcdClient: etcdClient,
        ctx:        ctx,
        cancel:     cancel,
    }
    r.ResolveNow(resolver.ResolveNowOptions{})
    return r
}

// 解析etcd
func (r *etcdResolver) ResolveNow(options resolver.ResolveNowOptions) {
    gofunc() {
        if err := r.watchServices(); err != nil {
            log.Printf("Failed to watch services: %v", err)
        }
    }()
}

// Close 關閉解析器
func (r *etcdResolver) Close() {
    r.cancel()
}



```
  1. 實現 watchServices,負責從註冊中心獲取實例地址,監聽變化。首先是從 gRPC 目標中提取服務名,拼接前綴並查詢 etcd,查詢回來後把多出來的/刪掉;然後解析鍵值對,把地址寫到 addr 列表裏面;通過解析出來的列表用r.conn.UpdateState()更新 gRPC 調用的地址;通過r.etcdClient.Watch監聽 etcd 的服務變化。
```go
func (r *etcdResolver) watchServices() error {
    serviceName := r.target.URL.Path

    iflen(serviceName) > 0 && serviceName[0] == '/' {
        serviceName = serviceName[1:]
    }

    if serviceName == "" {
        return fmt.Errorf("empty service name in target: %+v", r.target)
    }

    serviceKeyPrefix := fmt.Sprintf("services/%s/", serviceName)

    // 首次獲取服務列表
    resp, err := r.etcdClient.Get(r.ctx, serviceKeyPrefix, clientv3.WithPrefix())
    if err != nil {
        return err
    }

    addrs := make([]resolver.Address, 0)
    for _, kv := range resp.Kvs {
        keyParts := strings.Split(string(kv.Key)"/")
        iflen(keyParts) < 3 {
            log.Printf("Invalid key format: %s", string(kv.Key))
            continue
        }

        // 最後一部分是服務地址
        addr := keyParts[len(keyParts)-1]
        if addr == "" {
            log.Printf("Empty address in key: %s", string(kv.Key))
            continue
        }

        addrs = append(addrs, resolver.Address{
            Addr:       addr,
            Attributes: attributes.New("service", serviceName),
            ServerName: serviceName,
        })
    }

    log.Printf("Resolved addresses service==> %s Addr ==> %v", serviceName, addrs)

    // 更新客戶端連接狀態
    iflen(addrs) == 0 {
        return fmt.Errorf("no addresses found for service %s", serviceName)
    }

    if err := r.conn.UpdateState(resolver.State{Addresses: addrs}); err != nil {
        return err
    }

    // 監聽服務變化
    rch := r.etcdClient.Watch(r.ctx, serviceKeyPrefix, clientv3.WithPrefix())
    for wresp := range rch {
        for _, ev := range wresp.Events {
            // 從鍵中提取服務地址
            keyParts := strings.Split(string(ev.Kv.Key)"/")
            iflen(keyParts) < 3 {
                log.Printf("Invalid key format in watch: %s", string(ev.Kv.Key))
                continue
            }

            addr := keyParts[len(keyParts)-1]
            if addr == "" {
                log.Printf("Empty address in watch key: %s", string(ev.Kv.Key))
                continue
            }

            switch ev.Type {
            case clientv3.EventTypePut:
                if !containsAddr(addrs, addr) {
                    addrs = append(addrs, resolver.Address{
                        Addr:       addr,
                        Attributes: attributes.New("service", serviceName),
                    })
                    log.Printf("Added new address: %s", addr)
                }
            case clientv3.EventTypeDelete:
                addrs = removeAddr(addrs, addr)
                log.Printf("Removed address: %s", addr)
            }
        }

        // 更新客戶端連接狀態
        if err := r.conn.UpdateState(resolver.State{Addresses: addrs}); err != nil {
            return err
        }
    }

    returnnil
}

```
  1. 修改主函數,初始化 etcdClient, 並使用 etcd 返回解析的地址調用服務
```go
func main() {
xxx
etcdEndpoints := []string{"http://localhost:2379"}
    etcdClient, err := et.InitEtcdClient(etcdEndpoints)
    if err != nil {
        log.Fatal("init etcd client error:", err)
    }
    defer etcdClient.Close()

    et.InitResolver(etcdClient)

    serviceName := "user_service"
    targetURI := fmt.Sprintf("%s:///%s", et.Scheme, serviceName)
    log.Printf("Target URI: %s", targetURI)
    conn, err := grpc.Dial(
        targetURI,
        grpc.WithTransportCredentials(insecure.NewCredentials()),
        grpc.WithResolvers(&et.EtcdResolverBuilder{EtcdClient: etcdClient}), // etcd解析回來的地址
        grpc.WithChainUnaryInterceptor(
            ji.JwtClientInterceptor(),
            li.LogUnaryServerInterceptor(),
        )) //
    if err != nil {
        log.Fatalf("did not connect: %v", err)
    }
    defer conn.Close()
}

```
  1. 最後調用一下:

本文由 Readfog 進行 AMP 轉碼,版權歸原作者所有。
來源https://mp.weixin.qq.com/s/d2LYlCihaqv5AtTZTzBRvQ