三、服务注册与发现

2021-02-18T14:21:02+08:00 | 26分钟阅读 | 更新于 2021-02-18T14:21:02+08:00

@

学习目标

学完本章你应该能够:

  1. 讲清服务注册与发现的演进四阶段(写死 IP → DNS → 注册中心 → 网关),以及每一阶段解决了什么、留下什么新问题。
  2. 描述注册中心基本模型的四个步骤:注册、心跳续约、订阅、变更通知,并画出推/拉两种服务发现模型的差异。
  3. 对照 gRPC、Dubbo-go、Kratos、EGO、go-micro,说出它们服务注册与发现的接口设计与注册维度(接口维度 vs 应用维度)。
  4. 基于 etcd 说出租约(Lease)机制如何解决"服务宕机自动下线",并讲清 Register / Unregister / ListServices / Subscribe 的实现要点。
  5. 在面试里把"客户端连不上注册中心怎么办"“注册中心崩溃怎么办"等容错场景讲成有取舍的判断,而不是背条文。

前置知识

  • Go 语言基础(interface、goroutine、channel、context
  • gRPC 基本用法(grpc.Dial、ClientConn 概念)
  • 分布式基础(心跳、CAP 取舍)

本章你会动手做的事

  1. 用 etcd 跑通"服务端注册 + 客户端通过服务名拨号"的最小 demo。
  2. 把 6.3 节的容错流程图抄一遍,并标注在你自己的系统里会用策略 A(一致性优先)还是策略 B(可用性优先)。
  3. RegisterleaseTTL 设不同值,观察服务宕机后注册中心多久才发现。

前言:从写死地址到服务发现

在前两章中,我们已经基本了解了 RPC 是怎么构建的。但在前面的代码演示中,我们都是通过直接写死的方式来连上服务器:

// 写死服务端地址
client, err := NewClientV2("127.0.0.1:9091", &JSONSerializer{})

实际中我们能这么直接写死吗?答案是:除非是开发或者联调,不然我们很少会通过写死的形式来指定下游

大多数时候,在微服务架构中,建议先在本地环境中和上下游联调好,直接使用对方的局域网地址,两边开着 DEBUG,很好联调测试。但在生产环境中,服务地址是动态变化的——实例会扩容、缩容、重启、迁移,写死地址根本无法应对。

那么,客户端怎么知道服务端在哪里?这就是服务注册与发现要解决的问题。

本教程涵盖以下主题:

  1. 服务注册与发现演进——从写死 IP 到注册中心的四个阶段
  2. 注册中心基本模型——注册、心跳、订阅、通知
  3. 主流框架对比——gRPC、Dubbo-go、Kratos、EGO、go-micro
  4. etcd 注册中心实现——Registry 接口、租约机制、服务注册
  5. gRPC 服务发现实现——Resolver、Builder、服务发现与监听
  6. 容错场景分析——各种故障情况下的应对策略
  7. 面试要点总结

一、服务注册与发现演进

服务注册与发现并非一开始就有注册中心,它经历了一个演化的过程。

1.1 第一阶段:直接 IP + 端口访问

flowchart LR
    C[客户端
配置: 172.10.0.1:1234] -->|直接连接| S[服务端
172.10.0.1:1234]
  • 客户端保存着 IP 端口信息,需要提前配置好
  • 缺点:IP 可能会变,如果实例很多,配置难以维护

这就像你把朋友的电话号码存在脑子里——如果朋友换了号码,你就联系不上了。而且如果你有 100 个朋友,全靠脑子记根本不现实。

1.2 第二阶段:域名解析(DNS)

flowchart LR
    C[客户端
service_a.yourcompany.com] -->|DNS解析| D[DNS服务器] D -->|返回IP| C C -->|连接| S[服务端
172.10.0.1:1234]
  • 客户端保存服务的 endpoint(核心是域名),客户端也可以缓存 DNS 解析的结果
  • 缺点:客户端不缓存则多了一次调用,缓存则存在不一致性
  • 基于 HTTP 的微服务架构还在大规模应用这种方式,gRPC 也可以直接使用这种形态

DNS 就像一个电话簿——你只需要记住朋友的名字,查一下就知道号码了。但电话簿更新有延迟,你可能拿到的是旧号码。

1.3 第三阶段:分布式协调——注册中心阶段

flowchart LR
    S[服务端] -->|set service_a| R[注册中心]
    C[客户端] -->|get service_a| R
    R -->|返回实例列表| C
    C -->|直接连接| S
    R -.->|通知变更| C
  • 服务端发起注册,客户端向注册中心询问
  • 注册中心维持住两端
  • 基本没有显著的缺点,不过在大规模集群下,注册中心容易成为瓶颈,网络中比较多探活流量

注册中心就像一个实时更新的通讯录——朋友换了号码会主动告诉通讯录,你每次打电话前查一下最新的通讯录就行。

1.4 变种:经过中间件(网关)的调用方式

flowchart LR
    C[客户端] -->|请求| G[网关/中间件]
    G -->|转发| S1[服务端1
172.10.0.1:1234] G -->|转发| S2[服务端2
172.10.0.2:1234] S1 -->|响应| G G -->|响应| C
  • 客户端连上中间件,中间件转发给服务端
  • 缺点是因为在客户端和服务端中间引入了额外的中间件,性能稍差
  • 请求和响应都要经过中间件,增加了网络开销

二、注册中心基本模型

2.1 核心流程

服务注册中心模式,核心依赖于一个第三方组件(注册中心)。基本模型如下:

flowchart TB
    subgraph 服务端
        S1[1.服务启动成功后主动注册] --> S2[2.与注册中心保持心跳]
    end
    subgraph 注册中心
        R1[接收注册] --> R2[维持心跳]
        R2 --> R3[通知变更]
    end
    subgraph 客户端
        C1[3.启动时主动订阅服务数据] --> C2[4.接收注册中心变更通知]
    end
    S1 --> R1
    S2 --> R2
    R3 --> C2
    C1 --> R1

四个基本步骤:

  1. 服务启动成功之后主动注册——把自己的地址信息写入注册中心
  2. 服务端和注册中心保持心跳——定期续约,证明自己还活着
  3. 客户端启动的时候要主动订阅对应服务的数据——拉取可用服务列表
  4. 注册中心要通知客户端变更——服务上下线时推送更新

2.2 理解难点

所有的面试都是围绕以下几个问题来进行的:

  • 运行过程中,客户端连不上注册中心怎么办?
  • 运行过程中,客户端拿到了注册数据,但连不上对应的服务端怎么办?
  • 运行过程中,注册中心没有收到服务端心跳怎么办?
  • 注册中心崩溃了,客户端和服务端怎么办?

这些问题我们会在第六章详细分析。

2.3 服务发现的两个基本模型

模型说明优缺点
推模型注册中心主动推送变更给客户端实时性好,但实现复杂
拉模型客户端轮询注册中心实现简单,但要考虑刷新频率——快了浪费资源,慢了数据不一致

gRPC 的 DNS Resolver 就是一个典型的拉模型——轮询 DNS 服务器来刷新本地可用的服务器列表。


三、主流框架的服务注册与发现

3.1 gRPC

gRPC 服务注册与发现的核心在 resolver 包,核心接口有:

  • Target:被解析的目标,是对服务的抽象
  • Builder:Builder 模式,用于构建一个 Resolver
  • Resolver:负责服务发现
  • Address:代表一个地址的抽象

为何没有服务注册? 因为 gRPC 本身并没有预设存在服务注册中心这种东西。gRPC 本身有一个 ServiceRegister,但它表达的是注册到本地(即 RegisterService 的用途)。

本质上来说,在 gRPC 里使用注册中心,只有 Resolver 这端(服务发现)设计了接口,服务注册需要框架使用者自己实现。

gRPC 服务注册与发现的重要实现就是基于 DNS 的实现:dnsResolverdnsBuilder。其中核心方法是 watcher,整体逻辑是一个拉模型。

3.2 Dubbo-go

Dubbo-go 的注册模型就是典型的服务注册与发现模型,多了一个 Monitor 组件。

核心接口是 Registry

  • 注册与取消注册
  • 订阅与取消订阅

比较有意思的是,Dubbo-go 的订阅采用的是 Listener + 回调的设计。一般来说,Go 里面因为有 channel,让调用者从 channel 里面捞数据会更加常见一点。

Dubbo-go 这种叠床架屋的设计过于复杂,而且还要保持和 Java 的 Dubbo 兼容,因此如果单纯设计一个 Go 微服务框架的话是用不上这些设计的。

Dubbo-go 有两种注册模式:

注册模式说明特点
接口维度注册一个接口注册一次,利用 ip+端口组成 key可控性强、粒度细,但开销大
应用维度注册只注册一个键值对,如 ip+port -> user-application粒度粗,但数据更小

3.3 Kratos

Kratos 本身是建立在 gRPC 上的,但它又有自己的服务注册与发现接口:

  • Registry:对应服务注册部分
  • Discovery:对应服务发现部分
  • Watcher:监听服务变更部分

Kratos 是应用维度注册,注册的就是 App 的信息,会把 App 里面的所有 Server 的信息写到注册中心。在保持心跳的过程中,核心就是不断续约。

而服务发现则是利用了 gRPC 的接口,传入了一个 gRPC 的 ClientOption。这符合之前的分析:使用 gRPC 作为底层通信的微服务框架,服务注册过程肯定是另起炉灶的,但服务发现可以直接利用 gRPC 的 Resolver 接口。

3.4 EGO

EGO 服务注册与发现的核心接口是 Registry,方法也分成三部分:服务注册、服务发现、监听服务。实现非常类似于 Kratos,也是利用 etcd 的 PUT 然后设置租约。

EGO 也是应用级别注册,注册的数据定义在 ServiceInfo 中。不过在 EGO 中,应用其实是指一个个 Server,不同 Server 使用不同的端口。

3.5 go-micro

核心接口是 Registry,方法也可以分成:服务注册、服务发现、监听服务。注册基本上和 Kratos、EGO 一样,都是利用 etcd 的 PUT,然后设置租约。

3.6 注册维度与注册数据总结

注册维度:

  • 接口维度:一个接口注册一次,服务发现时使用接口来查找。可控性更强、粒度更细,但写入的数据、心跳等开销更大
  • 应用维度:一个应用注册一次,一个应用有很多服务。粒度更粗,但数据更小

两种注册维度优劣并不是很明显,并且都有人用。

注册数据:

  • 定位信息:例如 IP + 端口
  • 其它:主要取决于微服务框架的具体功能,例如在 Dubbo-go 里面可以写入标签信息、分组信息

四、etcd 注册中心实现

4.1 为什么选择 etcd

我们将使用 gRPC 作为底层通信协议,然后为 gRPC 增加服务注册与发现功能,原因是:

  • gRPC 已经成了微服务选型的第一选择,第二选型是 HTTP 协议
  • 这年头 RPC 协议已经很难设计得比 gRPC 好了

而注册中心我们选择 etcd,因为它是 Go 生态中最流行的分布式 KV 存储,天然支持租约(Lease)和监听(Watch)机制,非常适合做服务注册与发现。

4.2 etcd 中存储什么数据

往 etcd 里注册服务,本质上就是写入数据。我们需要写入:

  • 服务实例的信息
    • 定位信息(IP + 端口)
    • 其它元数据

存储格式可以使用 JSON、XML 甚至 Protobuf,只要客户端知道如何解析就可以。我们使用 JSON 来存储。

引入两个 key 概念:

  • service key/micro/service-name——标识一个服务
  • instance key/micro/service-name/instance-name——标识服务下的一个实例
/micro/
  user-service/                    <- service key
    172.10.0.1:9091                <- instance key -> {"ip":"172.10.0.1","port":9091,...}
    172.10.0.2:9091                <- instance key -> {"ip":"172.10.0.2","port":9091,...}
  order-service/                   <- service key
    172.10.0.3:9092                <- instance key -> {"ip":"172.10.0.3","port":9092,...}

4.3 定义 Registry 接口

首先定义一个代表注册中心的接口:

package registry

import (
	"context"
	"time"
)

// ============================================================
// 第一部分:核心数据结构与接口定义
// ============================================================

// ServiceInstance 代表一个服务实例的信息
// 每个运行中的服务进程对应一个 ServiceInstance
type ServiceInstance struct {
	ServiceName string            `json:"serviceName"` // 服务名,例如 "user-service"
	InstanceID  string            `json:"instanceId"`  // 实例ID,唯一标识一个实例,例如 "172.10.0.1:9091"
	Address     string            `json:"address"`     // 服务地址(IP),例如 "172.10.0.1"
	Port        int               `json:"port"`        // 服务端口,例如 9091
	Metadata    map[string]string `json:"metadata"`    // 元数据,例如版本号、权重等
}

// Registry 是注册中心的抽象接口
// 不同的注册中心实现(etcd、consul、zookeeper)都实现这个接口
// 接口方法分成三部分:服务注册、服务发现、监听服务
type Registry interface {
	// --- 服务注册部分 ---

	// Register 将服务实例注册到注册中心
	// 参数 ctx 用于控制超时和取消
	// 参数 instance 是要注册的服务实例信息
	// 返回的 cancel 函数用于取消注册(服务下线时调用)
	Register(ctx context.Context, instance *ServiceInstance) (cancel context.CancelFunc, err error)

	// Unregister 从注册中心取消注册
	// 服务正常退出时调用,主动删除注册数据
	Unregister(ctx context.Context, instance *ServiceInstance) error

	// --- 服务发现部分 ---

	// ListServices 根据服务名列出所有可用的服务实例
	// 客户端启动时调用,拉取当前可用的服务列表
	ListServices(ctx context.Context, serviceName string) ([]*ServiceInstance, error)

	// --- 监听服务部分 ---

	// Subscribe 订阅服务变更
	// 当服务实例上线/下线时,通过 channel 通知调用方
	// 返回的 channel 会接收服务实例的变更事件
	Subscribe(ctx context.Context, serviceName string) (<-chan Event, error)
}

// Event 代表服务变更事件
type Event struct {
	Type     EventType         // 事件类型:新增或删除
	Instance *ServiceInstance  // 变更的服务实例
}

// EventType 事件类型
type EventType uint8

const (
	EventTypeAdd    EventType = iota // 服务实例新增
	EventTypeDelete                  // 服务实例删除
)

4.4 租约机制:解决服务宕机问题

服务注册的关键是写入一个**“不稳定"的数据**——万一服务节点崩溃了,写入的数据能够被注册中心自动删除。

正常来说,服务节点可以正常退出,退出之前可以让服务节点主动删除注册数据。但是实际中会出现服务节点突然宕机(例如停电),或者网络故障之类的问题,服务节点压根没机会删除注册数据。

flowchart TB
    subgraph 正常退出
        N1[服务端准备退出] --> N2[主动调用 Unregister]
        N2 --> N3[注册中心删除数据]
    end
    subgraph 异常宕机
        A1[服务端突然宕机] --> A2[无法调用 Unregister]
        A2 --> A3[租约过期]
        A3 --> A4[注册中心自动删除数据]
    end

etcd 的租约(Lease)API 就是用来解决这个问题的:

  • 每一个写入的数据都有一个存活时间,这个叫做租约
  • 写入方需要不断续约,否则数据过期之后就会被删除

那么,如果服务端宕机了,肯定没人续约了,过一段时间之后,这个数据就会被注册中心自动删除。

4.5 etcd 注册中心完整实现

package registry

import (
	"context"
	"encoding/json"
	"fmt"
	"sync"
	"time"

	clientv3 "go.etcd.io/etcd/client/v3"
)

// ============================================================
// 第二部分:etcd 注册中心实现
// ============================================================

// EtcdRegistry 基于 etcd 的注册中心实现
// 利用 etcd 的 Lease(租约)机制实现服务健康检查
// 利用 etcd 的 Watch(监听)机制实现服务变更通知
type EtcdRegistry struct {
	client    *clientv3.Client  // etcd 客户端连接
	keyPrefix string            // etcd 中 key 的前缀,例如 "/micro/"

	// 租约相关配置
	leaseTTL      int64         // 租约存活时间(秒),例如 15 秒
	keepaliveOnce sync.Once     // 确保心跳 goroutine 只启动一次

	// 每个服务实例对应的租约 ID
	// key: instance key, value: lease ID
	leases sync.Map             // 并发安全的 map,存储实例与租约的映射
}

// NewEtcdRegistry 创建一个新的 etcd 注册中心
// 参数 endpoints 是 etcd 集群的地址列表,例如 []string{"127.0.0.1:2379"}
func NewEtcdRegistry(endpoints []string) (*EtcdRegistry, error) {
	// 创建 etcd 客户端
	client, err := clientv3.New(clientv3.Config{
		Endpoints:   endpoints,              // etcd 地址列表
		DialTimeout: 5 * time.Second,       // 连接超时
	})
	if err != nil {
		return nil, fmt.Errorf("创建 etcd 客户端失败: %w", err)
	}

	return &EtcdRegistry{
		client:    client,
		keyPrefix: "/micro/",  // key 前缀,所有服务注册数据都存放在这个前缀下
		leaseTTL:  15,         // 租约 15 秒,即服务 15 秒不续约就认为宕机
	}, nil
}

// instanceKey 构造实例在 etcd 中的 key
// 格式:/micro/service-name/instance-id
func (r *EtcdRegistry) instanceKey(inst *ServiceInstance) string {
	return fmt.Sprintf("%s%s/%s", r.keyPrefix, inst.ServiceName, inst.InstanceID)
}

// servicePrefix 构造服务在 etcd 中的前缀
// 格式:/micro/service-name/
// 用于 ListServices 时按前缀查询所有实例
func (r *EtcdRegistry) servicePrefix(serviceName string) string {
	return fmt.Sprintf("%s%s/", r.keyPrefix, serviceName)
}

// ============================================================
// 第三部分:服务注册实现
// ============================================================

// Register 将服务实例注册到 etcd
// 核心步骤:
// 1. 创建租约(Lease)
// 2. 将实例信息以 JSON 格式写入 etcd,绑定租约
// 3. 启动 goroutine 持续续约(KeepAlive)
func (r *EtcdRegistry) Register(ctx context.Context, instance *ServiceInstance) (context.CancelFunc, error) {
	// 第一步:创建租约
	// 租约是一个存活时间,超过这个时间没有续约,数据就会被自动删除
	// 这就是服务宕机后能自动从注册中心移除的关键
	leaseResp, err := r.client.Grant(ctx, r.leaseTTL)
	if err != nil {
		return nil, fmt.Errorf("创建租约失败: %w", err)
	}
	leaseID := leaseResp.ID

	// 第二步:将实例信息序列化为 JSON
	data, err := json.Marshal(instance)
	if err != nil {
		return nil, fmt.Errorf("序列化实例信息失败: %w", err)
	}

	// 第三步:将数据写入 etcd,绑定租约
	// key = /micro/service-name/instance-id
	// value = JSON 格式的实例信息
	// lease = 刚创建的租约 ID
	key := r.instanceKey(instance)
	_, err = r.client.Put(ctx, key, string(data), clientv3.WithLease(leaseID))
	if err != nil {
		return nil, fmt.Errorf("写入 etcd 失败: %w", err)
	}

	// 记录租约 ID,用于后续取消注册
	r.leases.Store(key, leaseID)

	// 第四步:启动续约 goroutine
	// KeepAlive 会持续向 etcd 发送续约请求,保持租约不过期
	// 如果服务宕机,这个 goroutine 也会停止,租约就会过期
	keepAliveCtx, cancel := context.WithCancel(context.Background())
	go r.keepAlive(keepAliveCtx, leaseID)

	// 返回 cancel 函数,调用方可以在服务下线时调用
	// 调用 cancel 会停止续约 goroutine,租约随后过期,数据被删除
	return cancel, nil
}

// keepAlive 持续向 etcd 续约
// 只要这个 goroutine 在运行,租约就不会过期
// 当 ctx 被取消时(服务下线),续约停止,租约过期
func (r *EtcdRegistry) keepAlive(ctx context.Context, leaseID clientv3.LeaseID) {
	// client.KeepAlive 返回一个 channel
	// 每次续约成功,channel 就会收到一个响应
	// 如果 ctx 取消或 etcd 连接断开,channel 会被关闭
	ch, err := r.client.KeepAlive(ctx, leaseID)
	if err != nil {
		// 续约失败,租约会在 leaseTTL 秒后过期
		return
	}

	// 持续读取续约响应
	// 如果 ctx 被取消,KeepAlive 会关闭 channel,循环退出
	for range ch {
		// 收到续约响应,什么都不用做
		// 只要从 channel 能读到数据,就说明续约成功
	}
}

// Unregister 从 etcd 取消注册
// 服务正常退出时调用,主动删除注册数据
func (r *EtcdRegistry) Unregister(ctx context.Context, instance *ServiceInstance) error {
	key := r.instanceKey(instance)

	// 从 etcd 中删除数据
	_, err := r.client.Delete(ctx, key)
	if err != nil {
		return fmt.Errorf("删除注册数据失败: %w", err)
	}

	// 清理租约记录
	r.leases.Delete(key)
	return nil
}

4.6 服务注册时机

严格来说,服务注册的时机是很难做到完美的。大多数时候,我们是在服务启动完成之后才注册上去的。

而"服务启动完成"就是一个含糊的概念:

  • 监听端口成功了,算不算启动成功了?
  • 服务已经准备好了处理请求,算不算启动成功?
  • 服务已经预热好了(例如预加载了缓存),算不算启动成功?

我们的方案:将时机问题抛给用户。用户自己决定怎样才算是启动起来了。用户必须在所有的初始化工作都完成后,才能最终调用 Start 方法。而在 Start 方法里面,我们就是单纯启动端口之后就注册过去注册中心。

4.7 利用健康检查来决定注册时机

微服务框架里面有一种做法,就是微服务框架在启动服务之后,会给服务发一个健康检查或者心跳,如果收到了成功的响应,那么就认为服务启动成功了。只有在这个时候,微服务框架才会注册数据。

但实际上,这个东西没什么用——健康检查通过了,只能认为端口启动成功了,但你的微服务在业务层面上启动了没有,还是不知道。

// 服务注册的使用示例
func main() {
	// 创建 etcd 注册中心
	registry, err := NewEtcdRegistry([]string{"127.0.0.1:2379"})
	if err != nil {
		panic(err)
	}

	// 创建服务实例信息
	instance := &ServiceInstance{
		ServiceName: "user-service",
		InstanceID:  "172.10.0.1:9091", // 用 IP:Port 作为实例ID
		Address:     "172.10.0.1",
		Port:        9091,
		Metadata: map[string]string{
			"version": "v1.0.0",
			"weight":  "100",
		},
	}

	// 注意:用户必须确保所有初始化工作都完成后才调用注册
	// 例如数据库连接池已创建、缓存已预热等
	ctx := context.Background()
	cancel, err := registry.Register(ctx, instance)
	if err != nil {
		panic(err)
	}

	// 启动 gRPC 服务(省略具体代码)
	// servegrpc()

	// 服务退出时取消注册
	defer func() {
		cancel()                   // 停止续约
		registry.Unregister(ctx, instance) // 主动删除注册数据
	}()
}

五、gRPC 服务发现实现

5.1 gRPC 服务发现的核心接口

gRPC 服务发现的核心接口:

  • ClientConn:它是一个抽象,代表的是对一个服务的连接,而不是一个 TCP 连接。一个 ClientConn 内部可能维护了多个 TCP 连接
  • resolver.Builder:创建 Resolver。gRPC 会维护一个 scheme -> Builder 的映射
  • Resolver:它和服务进行绑定,一个服务一个 Resolver。它负责和注册中心交互,监听注册数据的变化
flowchart TB
    subgraph gRPC 内部
        D[grpc.Dial
scheme:///service-name] --> SB[scheme -> Builder 映射表] SB -->|找到对应Builder| B[resolver.Builder] B -->|Build| R[Resolver] R -->|更新地址列表| CC[ClientConn] end subgraph 注册中心 R -->|ListServices| E[etcd] R -->|Subscribe| E E -->|变更通知| R end

5.2 gRPC 注册中心基本流程

gRPC 本身没有提供服务注册接口,服务注册本身是我们自己管的。gRPC 只提供了服务解析的接口,也就是 Resolver 和对应的 Builder 两个接口。

一般步骤:

  1. 用户在初始化 gRPC 的时候指定 grpc.WithResolver 选项,传入自定义的 Resolver
  2. Dial 调用的时候传入服务标识符,一般形式是 scheme:///service-name。scheme 代表的是如何通信,大多数时候它就代表我们的注册中心
  3. gRPC 会根据 scheme 来找到我们注册的 Resolver,我们在 Resolver 里面更新可用的连接

5.3 Resolver 和 Builder 实现

package registry

import (
	"context"
	"encoding/json"
	"fmt"
	"sync"
	"time"

	"google.golang.org/grpc/resolver"
	clientv3 "go.etcd.io/etcd/client/v3"
)

// ============================================================
// 第四部分:gRPC Resolver 和 Builder 实现
// ============================================================

// etcdBuilder 实现 resolver.Builder 接口
// 负责根据 scheme 创建 Resolver
type etcdBuilder struct {
	client    *clientv3.Client // etcd 客户端
	keyPrefix string           // key 前缀
}

// NewEtcdBuilder 创建一个 etcd resolver builder
// 参数 scheme 是注册中心的协议标识,例如 "etcd"
// 这样 gRPC 拨号时使用 "etcd:///service-name" 就会使用这个 builder
func NewEtcdBuilder(client *clientv3.Client, scheme, keyPrefix string) *etcdBuilder {
	// 将 builder 注册到 gRPC 的 scheme 映射表中
	// gRPC 内部维护了一个 scheme -> builder 的 map
	resolver.Register(&etcdBuilder{
		client:    client,
		keyPrefix: keyPrefix,
	})

	return &etcdBuilder{
		client:    client,
		keyPrefix: keyPrefix,
	}
}

// Build 是 resolver.Builder 接口的核心方法
// 当 gRPC 拨号时遇到对应的 scheme,就会调用这个方法创建 Resolver
// 参数:
//   - target: 拨号目标,包含 scheme 和服务名
//   - cc: ClientConn,用于更新可用地址列表
//   - opts: resolver 选项
func (b *etcdBuilder) Build(target resolver.Target, cc resolver.ClientConn, opts resolver.BuildOptions) (resolver.Resolver, error) {
	// 创建 etcd Resolver
	r := &etcdResolver{
		client:    b.client,
		keyPrefix: b.keyPrefix,
		target:    target,           // 拨号目标(包含服务名)
		cc:        cc,               // ClientConn,用于更新地址
	}

	// 启动一个 goroutine 来处理服务发现
	go r.start()

	return r, nil
}

// Scheme 返回这个 builder 对应的 scheme
// gRPC 通过这个方法来匹配拨号地址中的 scheme
func (b *etcdBuilder) Scheme() string {
	return "etcd" // 对应 "etcd:///service-name" 中的 "etcd"
}

// ============================================================
// 第五部分:etcd Resolver 实现
// ============================================================

// etcdResolver 实现 resolver.Resolver 接口
// 负责从 etcd 拉取服务列表,并监听变更
type etcdResolver struct {
	client    *clientv3.Client  // etcd 客户端
	keyPrefix string            // key 前缀
	target    resolver.Target   // 拨号目标
	cc        resolver.ClientConn // gRPC 连接,用于更新地址列表
	mu        sync.Mutex        // 保护 addrs 的并发访问
}

// start 启动服务发现
// 1. 先从 etcd 拉取当前所有服务实例
// 2. 然后监听 etcd 中服务信息的变更
func (r *etcdResolver) start() {
	serviceName := r.target.Endpoint() // 从 target 中提取服务名
	prefix := fmt.Sprintf("%s%s/", r.keyPrefix, serviceName)

	// 第一步:从 etcd 拉取所有服务实例
	// 使用 Get + WithPrefix 查询指定前缀下的所有 key
	resp, err := r.client.Get(context.Background(), prefix, clientv3.WithPrefix())
	if err != nil {
		return
	}

	// 将拉取到的实例信息转换为 gRPC 的地址列表
	var addrs []resolver.Address
	for _, kv := range resp.Kvs {
		var instance ServiceInstance
		if err := json.Unmarshal(kv.Value, &instance); err != nil {
			continue // 跳过无法解析的数据
		}
		// 构造 gRPC 地址
		addrs = append(addrs, resolver.Address{
			Addr:       fmt.Sprintf("%s:%d", instance.Address, instance.Port),
			ServerName: instance.ServiceName,
		})
	}

	// 更新 gRPC 的连接地址列表
	// gRPC 内部会根据这个列表创建/维护 TCP 连接
	r.cc.UpdateState(resolver.State{
		Addresses: addrs,
	})

	// 第二步:监听 etcd 中服务信息的变更
	// 使用 Watch + WithPrefix 监听指定前缀下所有 key 的变更
	// 当有服务实例上线/下线时,etcd 会通过 channel 发送事件
	go r.watch(prefix)
}

// watch 监听 etcd 中服务信息的变更
// 当服务实例上线或下线时,更新 gRPC 的地址列表
func (r *etcdResolver) watch(prefix string) {
	// 创建 Watch channel
	// WithPrefix 表示监听指定前缀下所有 key 的变更
	// WithPrevKV 表示删除事件中会包含被删除 key 的旧值
	watchCh := r.client.Watch(context.Background(), prefix,
		clientv3.WithPrefix(), clientv3.WithPrevKV())

	for watchResp := range watchCh {
		for _, event := range watchResp.Events {
			switch event.Type {
			case clientv3.EventTypePut:
				// 有新的服务实例注册(或更新)
				// 解析新实例信息
				var instance ServiceInstance
				if err := json.Unmarshal(event.Kv.Value, &instance); err != nil {
					continue
				}
				// 这里采用简单策略:收到任何变更就重新拉取全部数据
				// 更精细的做法是只增加/删除变更的实例
				r.refreshAddresses(prefix)

			case clientv3.EventTypeDelete:
				// 有服务实例下线
				// event.PrevKv 包含被删除的旧值
				r.refreshAddresses(prefix)
			}
		}
	}
}

// refreshAddresses 重新从 etcd 拉取所有服务实例,更新地址列表
// 这对应 PDF 中提到的第一种实现思路:
// "监听到事件之后,不管事件的类型,直接拉取所有的服务信息,完整更新一下"
func (r *etcdResolver) refreshAddresses(prefix string) {
	resp, err := r.client.Get(context.Background(), prefix, clientv3.WithPrefix())
	if err != nil {
		return
	}

	var addrs []resolver.Address
	for _, kv := range resp.Kvs {
		var instance ServiceInstance
		if err := json.Unmarshal(kv.Value, &instance); err != nil {
			continue
		}
		addrs = append(addrs, resolver.Address{
			Addr:       fmt.Sprintf("%s:%d", instance.Address, instance.Port),
			ServerName: instance.ServiceName,
		})
	}

	r.cc.UpdateState(resolver.State{
		Addresses: addrs,
	})
}

// ResolveNow resolver.Resolver 接口要求实现的方法
// gRPC 在某些情况下会调用这个方法要求立即重新解析
// 我们使用 Watch 机制实时监听,所以这里不需要做任何事
func (r *etcdResolver) ResolveNow(o resolver.ResolveNowOptions) {}

// Close 关闭 Resolver,释放资源
func (r *etcdResolver) Close() {}

5.4 监听变更的两种实现思路

监听变更之后的核心就是更新可用服务节点列表,大体上可以分成两种实现思路:

思路说明优缺点
思路一监听到事件后,不管事件类型,直接拉取所有服务信息,完整更新实现简单,但每次变更都要拉取全部数据
思路二监听到事件后,根据事件类型决定动作:新增就添加,下线就删除性能好,但要维护本地状态,容易出错

思路一更简单可靠,推荐初学者使用。思路二性能更好,但需要仔细处理并发和各种边界情况。

5.5 etcd 服务发现实现

ListServices 就是根据服务名找到对应的所有的可用节点,本质上就是一个调用 API 的事情:

// ListServices 根据服务名列出所有可用服务实例
// 客户端启动时调用,拉取当前可用的服务列表
func (r *EtcdRegistry) ListServices(ctx context.Context, serviceName string) ([]*ServiceInstance, error) {
	// 构造查询前缀:/micro/service-name/
	prefix := r.servicePrefix(serviceName)

	// 从 etcd 中查询该前缀下的所有 key
	// WithPrefix 表示查询所有以 prefix 开头的 key
	resp, err := r.client.Get(ctx, prefix, clientv3.WithPrefix())
	if err != nil {
		return nil, fmt.Errorf("查询 etcd 失败: %w", err)
	}

	// 解析查询结果
	var instances []*ServiceInstance
	for _, kv := range resp.Kvs {
		var inst ServiceInstance
		if err := json.Unmarshal(kv.Value, &inst); err != nil {
			continue // 跳过无法解析的数据
		}
		instances = append(instances, &inst)
	}

	return instances, nil
}

// Subscribe 订阅服务变更
// 通过 etcd 的 Watch 机制监听服务实例的上线/下线
func (r *EtcdRegistry) Subscribe(ctx context.Context, serviceName string) (<-chan Event, error) {
	prefix := r.servicePrefix(serviceName)

	// 创建事件 channel
	eventCh := make(chan Event, 10) // 带缓冲,防止事件丢失

	// 启动 goroutine 监听 etcd 变更
	go func() {
		defer close(eventCh)

		// 创建 Watch
		watchCh := r.client.Watch(ctx, prefix,
			clientv3.WithPrefix(), clientv3.WithPrevKV())

		for watchResp := range watchCh {
			for _, event := range watchResp.Events {
				switch event.Type {
				case clientv3.EventTypePut:
					// 新服务实例注册
					var inst ServiceInstance
					if err := json.Unmarshal(event.Kv.Value, &inst); err != nil {
						continue
					}
					eventCh <- Event{
						Type:     EventTypeAdd,
						Instance: &inst,
					}

				case clientv3.EventTypeDelete:
					// 服务实例下线
					// event.PrevKv 包含被删除的旧值
					if event.PrevKv == nil {
						continue
					}
					var inst ServiceInstance
					if err := json.Unmarshal(event.PrevKv.Value, &inst); err != nil {
						continue
					}
					eventCh <- Event{
						Type:     EventTypeDelete,
						Instance: &inst,
					}
				}
			}
		}
	}()

	return eventCh, nil
}

5.6 完整使用示例

package main

import (
	"context"
	"fmt"
	"time"

	"google.golang.org/grpc"
	clientv3 "go.etcd.io/etcd/client/v3"
)

func main() {
	// ========================================
	// 服务端:注册服务到 etcd
	// ========================================

	// 创建 etcd 客户端
	etcdClient, err := clientv3.New(clientv3.Config{
		Endpoints:   []string{"127.0.0.1:2379"},
		DialTimeout: 5 * time.Second,
	})
	if err != nil {
		panic(err)
	}

	// 创建注册中心
	registry, err := NewEtcdRegistry([]string{"127.0.0.1:2379"})
	if err != nil {
		panic(err)
	}

	// 注册服务实例
	instance := &ServiceInstance{
		ServiceName: "user-service",
		InstanceID:  "172.10.0.1:9091",
		Address:     "172.10.0.1",
		Port:        9091,
		Metadata: map[string]string{
			"version": "v1.0.0",
		},
	}

	ctx := context.Background()
	cancel, err := registry.Register(ctx, instance)
	if err != nil {
		panic(err)
	}
	defer func() {
		cancel() // 停止续约
		registry.Unregister(ctx, instance) // 主动删除注册数据
	}()

	// 启动 gRPC 服务(省略具体代码)
	// grpcServer.Serve(lis)

	// ========================================
	// 客户端:通过 gRPC 服务发现连接服务端
	// ========================================

	// 注册 etcd resolver builder 到 gRPC
	// 这样 gRPC 拨号时遇到 "etcd" scheme 就会使用我们的 resolver
	NewEtcdBuilder(etcdClient, "etcd", "/micro/")

	// 使用服务名拨号,而不是写死 IP
	// "etcd:///user-service" 的含义:
	// - etcd: 使用 etcd resolver
	// - user-service: 要发现的服务名
	conn, err := grpc.Dial(
		"etcd:///user-service",           // 服务标识符,不写死 IP
		grpc.WithDefaultServiceConfig(`{"loadBalancingPolicy":"round_robin"}`), // 负载均衡策略
		grpc.WithInsecure(),               // 不使用 TLS(仅用于演示)
	)
	if err != nil {
		panic(err)
	}
	defer conn.Close()

	// 此时 conn 已经自动发现了 user-service 的所有实例
	// 后续的 RPC 调用会自动负载均衡到不同的实例
	fmt.Println("已连接到 user-service")

	// 订阅服务变更(可选)
	eventCh, _ := registry.Subscribe(ctx, "user-service")
	go func() {
		for event := range eventCh {
			switch event.Type {
			case EventTypeAdd:
				fmt.Printf("新实例上线: %s:%d\n", event.Instance.Address, event.Instance.Port)
			case EventTypeDelete:
				fmt.Printf("实例下线: %s:%d\n", event.Instance.Address, event.Instance.Port)
			}
		}
	}()
}

六、容错场景分析

服务注册与发现中最难的部分不是基本流程,而是各种异常情况下的容错处理。所有的面试都是围绕这些问题来进行的。

6.1 影响租约效果的参数

租约使用起来要比较小心,一般要考虑以下问题:

  • 间隔多久续约一次?
  • 续约失败了怎么办?
  • 续约能不能重试?如果可以怎么重试?重试几次以及每次重试的间隔怎么设置?

在大规模分布式系统里面,这些参数会影响:

  • 微服务节点崩溃了,隔多久注册中心才能发现?
  • 续约的请求会不会给注册中心庞大的压力?
  • 续约的请求会不会挤占正常业务请求的网络流量?

6.2 服务端崩溃后客户端多久才知道

取决于服务端和注册中心的交互方式:

交互方式发现速度取决于
注册中心主动发起心跳心跳间隔 + 多少次心跳失败才判定崩溃
服务端主动续约租约长短 + 续约的重试机制

心跳是一个很复杂的事情:

  • 心跳频繁:会挤占正常请求的资源,但容易发现服务端崩溃
  • 心跳不频繁:注册中心很难发现节点崩溃

心跳失败之后的判定——次数难以确定:

  • 一次心跳失败就判定节点失活:过于严苛,部分时候网络抖动就会引起误判
  • 连续多次心跳失败才判定节点失活:未能及时发现节点崩溃

6.3 客户端和注册中心连不上怎么办

可能原因:网络问题、注册中心崩溃

flowchart TD
    Q[客户端连不上注册中心] --> S1{策略选择}
    S1 -->|策略A: 停止服务| A1[直接停止服务
直到恢复连接] S1 -->|策略B: 继续服务| B1[使用本地缓存数据
尝试重连注册中心] B1 --> B2{重连成功?} B2 -->|是| B3[更新缓存,正常服务] B2 -->|否,超时| B4[停止服务] B2 -->|否,未超时| B5[继续用缓存数据服务
调用不通的实例移出缓存]

两种策略:

  • 策略A(一致性优先):客户端直接停止服务,直到恢复和注册中心的连接
  • 策略B(可用性优先):客户端继续服务,并且尝试重新连上注册中心,这段时间内使用的都是本地缓存数据。再过一段时间还连不上,再停止服务

策略B更实用。继续服务的时候,如果服务端调用不通了,要把服务端实例挪出本地缓存。

6.4 服务端和注册中心连不上怎么办

基本上就是服务端重试,直到连上,然后输出错误日志。一般来说服务端不会因为连不上注册中心就停止服务。

注册中心在连不上服务端之后,经过重试依旧失败就会判定服务端挂了,这时候就会通知客户端。

6.5 注册中心和两端都能连上,但客户端和服务端连不上

对于客户端来说就是服务调用必然失败。如果单纯站在服务发现这个角度,客户端如果发现服务端连不上,就要将服务端暂时性挪出可用列表,后续再考虑挪回来。

这个问题其实是和后面的节点筛选混在一起讨论的——可能是负载均衡、健康检查、熔断等机制要处理的问题。

6.6 注册中心崩溃怎么办

实际上还是问的客户端容错问题,等同于客户端和注册中心连不上。基本上就是两个点:

  • 客户端要利用缓存的数据尽最大努力调用
  • 客户端要将失活的服务端移出缓存

还有一些奇诡的技术,客户端会把注册中心的数据写到本地文件上,不过本质上也是一种缓存。

6.7 容器内部 IP 问题

在我们的 etcd 实现里面,使用的是 IP + 端口来识别。问题就在于如果微服务运行于容器内部(例如 Docker),那么拿到的 IP 永远都是 localhost 的。而客户端运行在另外一个容器内,使用 localhost 是无法连上服务端的。

在这种情况下,服务端只能选择:

  • 使用容器专属的注册与发现方式
  • 容器启动的时候使用宿主机的网络(--network host),而不是创建一个虚拟网络
  • 容器启动的时候想办法把宿主机的 IP 作为环境变量注入进去

6.8 中间件选型

体量小的时候,只要是主流的中间件都没问题:ZooKeeper、Nacos、etcd……

体量大之后要考虑:

考量维度选项说明
集群模型对等模式所有节点平等,一致性问题
集群模型主从模式脑裂问题,主节点瓶颈
CAPCP一致性优先,小集群偏向
CAPAP可用性优先,大集群偏向

七、面试要点总结

7.1 基本模型

  • 服务注册与发现的基本模型:按照注册中心的那种形态去回答,也可以讲讲利用 DNS 之类的来做服务注册与发现
  • 服务注册的步骤:要注意讨论服务注册的时机问题——理论上我们需要在业务层面上保证服务启动成功了才能注册,但实际上是做不到这一点的
  • 服务发现的步骤:要点在于讨论客户端缓存数据的必要性,以及这种缓存引起的一致性问题,以及利用缓存可以进行容错
  • 服务下线的步骤:要讨论是注册中心下线服务端(服务端可能还活着),还是服务端主动下线(关机)。服务端主动下线要记得把自己从注册中心主动删除
  • 心跳该怎么设置:考虑心跳的间隔和判定服务端失活的条件。心跳间隔过长和过短的问题,以及判定服务端失活的条件的影响

7.2 中间件选型

  • 可以用什么中间件做注册中心:Nacos、etcd、ZooKeeper 都可以,简而言之就是所有的配置中心都可以。Redis 之类的缓存中间件也可以,但不是十分契合
  • 选用中间件要考虑什么:核心就是集群模式和 CAP 的取舍。基本上对于小规模应用来说,使用什么中间件都没问题
  • 注册中心 CAP 用 CP 还是 AP:分情况讨论,强调可用性的大规模集群就是 AP,其它就是 CP
  • 对等模式和主从模式的影响:主要考察主从模式的脑裂问题和主节点瓶颈,对等模式的一致性问题和转发问题
  • 多注册中心:一般是出于资源隔离和安全的角度,超大规模的公司可能会考虑这种方案

7.3 容错

  • 客户端和注册中心连不上了怎么办:主要讨论一致性还是可用性的取舍。一致性就代表连不上就认为服务端不能用了。可用性就是继续使用缓存的数据,但要考虑服务节点不可用怎么办
  • 服务端和注册中心连不上了怎么办:服务端重试直到连上,输出错误日志。一般服务端不会因为连不上注册中心就停止服务
  • 客户端和服务端连不上怎么办:客户端如果发现服务端连不上,要将服务端暂时性挪出可用列表
  • 注册中心崩溃怎么办:客户端利用缓存数据尽最大努力调用,将失活的服务端移出缓存

总结

本教程从服务注册与发现的演进历史出发,讲解了注册中心的基本模型,对比了主流框架的实现方案,并基于 etcd 完整实现了一个服务注册中心。

关键知识点回顾:

  1. 服务发现演进经历了写死 IP -> DNS -> 注册中心 -> 网关四个阶段
  2. 注册中心基本模型包含四个步骤:服务注册、心跳续约、客户端订阅、变更通知
  3. 租约机制是解决服务宕机自动移除的关键——写入方持续续约,宕机后租约过期数据自动删除
  4. gRPC 只提供了服务发现接口(Resolver/Builder),服务注册需要自己实现
  5. 注册维度分为接口维度和应用维度,各有优劣
  6. 容错的核心是客户端缓存 + 失活实例移出,在一致性和可用性之间做取舍
  7. 中间件选型主要考虑集群模式(对等 vs 主从)和 CAP(CP vs AP)

下一章我们将讲解负载均衡、路由与集群,探讨客户端拿到服务列表后如何选择合适的实例进行调用。


自测题与动手练习

自测题(合上书能答出来,才算懂)

  1. 服务注册与发现演进的四个阶段分别是什么?DNS 方案相比写死 IP 解决了什么,又留下了什么新问题?
  2. 注册中心基本模型的四个步骤是哪四个?推模型和拉模型在"实时性"和"实现复杂度"上各有什么取舍?
  3. gRPC 为什么"只有服务发现、没有服务注册”?要自己实现注册得怎么做?它的 Resolver / Builder 各自负责什么?
  4. etcd 的租约(Lease)机制是怎样让"宕机服务自动下线"的?如果 keepAlive 的 goroutine 因为网络抖动短暂断连,会发生什么?
  5. 客户端连不上注册中心时,策略 A(一致性优先)和策略 B(可用性优先)分别怎么处理?为什么文中说策略 B 更实用?

动手练习(建议真做一遍)

  1. 照着文中 main 示例,用 etcd 跑通"服务端 Register + 客户端 etcd:///user-service 拨号”,再手动 kill 服务端,观察客户端多久从可用列表里摘掉它。
  2. 把 6.3 节的容错流程图抄下来,结合你们线上系统的实际情况,标注每个故障场景该用策略 A 还是策略 B,并写一句理由。
  3. NewEtcdRegistryleaseTTL 分别设成 5s 和 60s,各制造一次服务端崩溃,记录注册中心"发现宕机"的延迟差异,体会心跳/租约参数对故障发现速度的影响。

本章小结

  • 服务发现演进走过写死 IP → DNS → 注册中心 → 网关四个阶段,注册中心用"实时通讯录"解决了地址动态变化的问题。
  • 注册中心基本模型四步:服务注册、心跳续约、客户端订阅、变更通知;推模型实时但复杂,拉模型简单但要权衡刷新频率。
  • 主流框架里 gRPC 只给服务发现接口(Resolver/Builder),注册要自己实现;Kratos/EGO/go-micro 都是应用维度注册,多用 etcd 的 PUT + 租约。
  • 租约机制是"宕机自动下线"的关键:写入方持续续约,宕机后没人续约,租约过期注册中心自动删数据。
  • 容错的核心是客户端缓存 + 失活实例移出:在一致性和可用性之间做取舍,策略 B(可用性优先 + 本地缓存)更实用。
  • 中间件选型主要看集群模式(对等 vs 主从)和 CAP(CP vs AP),小规模用什么都行,大规模要权衡。

下一章将讲解负载均衡、路由与集群,探讨客户端拿到服务列表后如何选择合适的实例进行调用。

About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!