gRPC Resolver 实现 Nacos 服务发现

2026/8/10 16:07:01Henry85 阅读1 点赞4 评论

自行实现

Nacos SDK 提供了 Subscribe 接口,借助 gRPC manual.Resolver 可以实现健康实例的实时推送,在此 mark 一下。

基本流程是通过 gRPC 逻辑地址 nacos:///group/service 中的 <schema>:///,对 manual.NewBuilderWithScheme(scheme string) Resolver Schema 进行匹配。

因为 grpc.WithResolvers(rs ...resolver.Builder) 是可以注册多个 resolver 的,当然也就需要进行区分。

此外可以发现上述 gRPC 逻辑地址中,nacos:/// 后面的 endpoint 其实并没有什么用,这里只是为了可读性。

GO
/*
 * @Author: Henry csthenry@foxmail.com
 * @Date: 2026-08-10 14:14:29
 * @LastEditors: Henry csthenry@foxmail.com
 * @LastEditTime: 2026-08-10 15:53:07
 * @FilePath: /engine/internal/infra/nacos/subscriber.go
 * @Description:
 *
 * Copyright (c) 2026 by Henry email: csthenry@foxmail.com, All Rights Reserved.
 */
package nacos

import (
	"fmt"
	"net"
	"strconv"
	"engine/internal/global"

	"github.com/nacos-group/nacos-sdk-go/v2/model"
	"github.com/nacos-group/nacos-sdk-go/v2/vo"
	"go.uber.org/zap"
	"google.golang.org/grpc"
	"google.golang.org/grpc/credentials/insecure"
	"google.golang.org/grpc/resolver"
	"google.golang.org/grpc/resolver/manual"
)

// NacosRPCSubscriber 含 Nacos 订阅能力的 GRPC Client
type NacosRPCSubscriber struct {
	Conn     *grpc.ClientConn
	Resolver *manual.Resolver
}

// buildGRPCAddrs 将 Nacos 实例转换为 GRPC 地址
func buildGRPCAddrs(insts []*Instance) []resolver.Address {
	addrs := make([]resolver.Address, 0, len(insts))

	for _, inst := range insts {
		if !inst.Enable || !inst.Healthy || inst.Weight <= 0 {
			continue
		}
		addr := resolver.Address{
			Addr: net.JoinHostPort(inst.Ip, strconv.FormatUint(inst.Port, 10)),
		}
		addrs = append(addrs, addr)
	}
	return addrs
}

// NewNacosRPC 创建支持 Nacos 动态服务发现的 GRPC Client
func NewNacosRPC(client *NacosClient, serviceName, groupName string) (*NacosRPCSubscriber, error) {
	if client == nil {
		return nil, fmt.Errorf("unavailable nacos client")
	}
	// 获取所有健康实例
	insts, err := client.SelectInstances(serviceName, groupName, true)
	if len(insts) == 0 || err != nil {
		return nil, fmt.Errorf("no healthy nacos instance: service: %s group: %s",
			serviceName,
			groupName,
		)
	}

	// gRPC manual resolver
	r := manual.NewBuilderWithScheme("nacos")
	r.InitialState(resolver.State{
		Addresses: buildGRPCAddrs(insts),
	})

	// GRPC 逻辑地址 nacos:///group/service
	target := fmt.Sprintf(
		"nacos:///%s/%s",
		groupName,
		serviceName,
	)

	conn, err := grpc.NewClient(
		target,
		grpc.WithResolvers(r),
		grpc.WithTransportCredentials(
			insecure.NewCredentials(),
		),
		// Use round_robin LB policy.
		grpc.WithDefaultServiceConfig(
			`{"loadBalancingConfig":[{"round_robin":{}}]}`,
		),
	)
	if err != nil {
		return nil, fmt.Errorf("create grpc client: %w", err)
	}

	// Nacos 订阅
	logger := global.MustLogger()
	err = client.namingClient.Subscribe(
		&vo.SubscribeParam{
			ServiceName: serviceName,
			GroupName:   groupName,
			SubscribeCallback: func(instsRaw []model.Instance, subErr error) {
				if subErr != nil {
					logger.Error("[Nacos] subscribe error",
						zap.String("Service", serviceName),
						zap.String("Group", groupName),
						zap.Error(subErr))
				}

				// 获取新的订阅地址
				insts := make([]*Instance, 0, len(instsRaw))
				for _, inst := range instsRaw {
					insts = append(insts, (*Instance)(&inst))
				}
				addrs := buildGRPCAddrs(insts)
				logger.Info("[Nacos] service changed",
					zap.String("Service", serviceName),
					zap.String("Group", groupName),
					zap.Any("Addrs", addrs),
				)

				// 更新 ClientConn 地址
				r.UpdateState(resolver.State{
					Addresses: addrs,
				})

			},
		},
	)

	if err != nil {
		_ = conn.Close()
		return nil, fmt.Errorf("subscribe nacos failed: %w", err)
	}
	return &NacosRPCSubscriber{
		Conn:     conn,
		Resolver: r,
	}, nil
}

上述逻辑使用到了我重新封装的 Nacos Client:

GO
/*
 * @Author: Henry csthenry@foxmail.com
 * @Date: 2026-08-05 21:11:51
 * @LastEditors: Henry csthenry@foxmail.com
 * @LastEditTime: 2026-08-10 14:24:19
 * @FilePath: /engine/internal/infra/nacos/client.go
 * @Description:
 *
 * Copyright (c) 2026 by Henry email: csthenry@foxmail.com, All Rights Reserved.
 */
package nacos

import (
	"fmt"

	"github.com/nacos-group/nacos-sdk-go/v2/clients"
	"github.com/nacos-group/nacos-sdk-go/v2/clients/config_client"
	"github.com/nacos-group/nacos-sdk-go/v2/clients/naming_client"
	"github.com/nacos-group/nacos-sdk-go/v2/common/constant"
	"github.com/nacos-group/nacos-sdk-go/v2/vo"
)

type NacosClient struct {
	namingClient naming_client.INamingClient
	configClient config_client.IConfigClient
}

func NewNacosClient(clientCfg NacosClientConfig, serverCfgs []NacosServerConfig) (*NacosClient, error) {
	clientConfig := constant.NewClientConfig(
		constant.WithNamespaceId(clientCfg.NamespaceId),
		constant.WithUsername(clientCfg.Username),
		constant.WithPassword(clientCfg.Password),
		constant.WithLogDir(clientCfg.LogDir),
		constant.WithLogLevel(clientCfg.LogLevel),
	)
	serverConfigs := make([]constant.ServerConfig, 0, len(serverCfgs))
	for _, serverCfg := range serverCfgs {
		serverConfigs = append(serverConfigs, *constant.NewServerConfig(
			serverCfg.IpAddr,
			serverCfg.Port,
			constant.WithScheme(serverCfg.Scheme),
		))
	}

	naming, err := clients.NewNamingClient(
		vo.NacosClientParam{
			ClientConfig:  clientConfig,
			ServerConfigs: serverConfigs,
		},
	)

	if err != nil {
		return nil, fmt.Errorf("create nacos client failed, err: %w", err)
	}

	config, err := clients.NewConfigClient(
		vo.NacosClientParam{
			ClientConfig:  clientConfig,
			ServerConfigs: serverConfigs,
		},
	)

	if err != nil {
		return nil, fmt.Errorf("create nacos client failed, err: %w", err)
	}

	return &NacosClient{
		namingClient: naming,
		configClient: config,
	}, nil
}

func (c *NacosClient) NamingClient() naming_client.INamingClient {
	return c.namingClient
}

func (c *NacosClient) ConfigClient() config_client.IConfigClient {
	return c.configClient
}

// Register 向 Nacos 注册实例
func (c *NacosClient) Register(ip string, port uint64, serviceName, groupName string, metadata map[string]string) error {
	success, err := c.namingClient.RegisterInstance(vo.RegisterInstanceParam{
		Ip:          ip,
		Port:        port,
		ServiceName: serviceName,
		GroupName:   groupName,
		Weight:      1,
		Enable:      true,
		Healthy:     true,
		Ephemeral:   true,
		Metadata:    metadata,
	})

	if !success || err != nil {
		return fmt.Errorf("register instance failed: %w", err)
	}
	return nil

}

// Deregister 向 Nacos 注销实例
func (c *NacosClient) Deregister(ip string, port uint64, serviceName, groupName string) error {
	success, err := c.namingClient.DeregisterInstance(vo.DeregisterInstanceParam{
		Ip:          ip,
		Port:        port,
		ServiceName: serviceName,
		GroupName:   groupName,
		Ephemeral:   true,
	})

	if !success || err != nil {
		return fmt.Errorf("deregister instance failed: %w", err)
	}
	return nil
}

// PublishConfig 向 Nacos 推送配置
func (c *NacosClient) PublishConfig(config []byte, dataId string, group string) error {
	success, err := c.configClient.PublishConfig(
		vo.ConfigParam{
			DataId:  dataId,
			Group:   group,
			Content: string(config),
			Type:    "json",
		},
	)
	if !success || err != nil {
		return fmt.Errorf("publish instance config failed: %w", err)
	}
	return nil
}

// DeleteConfig 删除 Nacos 上的配置
func (c *NacosClient) DeleteConfig(dataId string, group string) error {
	success, err := c.configClient.DeleteConfig(
		vo.ConfigParam{
			DataId: dataId,
			Group:  group,
		},
	)
	if !success || err != nil {
		return fmt.Errorf("delete instance config failed: %w", err)
	}
	return nil
}

// SelectOneHealthyInstance 从 Nacos 获取一个健康实例
func (c *NacosClient) SelectOneHealthyInstance(serviceName, groupName string) (*Instance, error) {
	inst, err := c.namingClient.SelectOneHealthyInstance(
		vo.SelectOneHealthInstanceParam{
			ServiceName: serviceName,
			GroupName:   groupName,
		},
	)
	if err != nil {
		return nil, err
	}
	return (*Instance)(inst), nil
}

// SelectInstances 获取 Nacos 服务中的所有实例
func (c *NacosClient) SelectInstances(serviceName, groupName string, healthyOnly bool) ([]*Instance, error) {
	rawInsts, err := c.namingClient.SelectInstances(
		vo.SelectInstancesParam{
			ServiceName: serviceName,
			GroupName:   groupName,
			HealthyOnly: healthyOnly,
		},
	)
	if err != nil {
		return nil, err
	}

	insts := make([]*Instance, 0, len(rawInsts))
	for _, inst := range rawInsts {
		insts = append(insts, (*Instance)(&inst))
	}

	return insts, nil
}

使用时十分简单,后续只需要管理一下 gRPC Conn 和 Nacos Client 的生命周期即可,某服务示例:

GO
grpcConn, err := nacos.NewNacosRPC(res.nacosClient, relationCfg.ServiceName, relationCfg.GroupName)
if err != nil {
	global.MustLogger().Error("[Nacos] 服务实例订阅失败", zap.Error(err))
}
relationRes, _ := initRelationResourceFromConn(grpcConn.Conn)
res.relationConn = relationRes.conn
res.relationClient = relationRes.client
global.MustLogger().Info("[Nacos] 服务实例订阅成功")

封装库

GO
package resolver

// Package resolver 实现基于 Nacos 的 gRPC 服务发现 Resolver。
// 通过 import 即可自动注册 "nacos" scheme,在 gRPC Dial 时使用 "nacos:///服务名" 的 target。

import (
	"fmt"
	"net"
	"strconv"

	"github.com/nacos-group/nacos-sdk-go/v2/clients/naming_client"
	"github.com/nacos-group/nacos-sdk-go/v2/model"
	"github.com/nacos-group/nacos-sdk-go/v2/vo"
	"google.golang.org/grpc/resolver"
)

// Option 可选的订阅参数配置
type Option func(*nacosResolverBuilder)

// WithGroupName 指定分组名,默认 "DEFAULT_GROUP"
func WithGroupName(group string) Option {
	return func(b *nacosResolverBuilder) {
		b.groupName = group
	}
}

// WithClusters 指定集群名,默认 "DEFAULT"
func WithClusters(clusters []string) Option {
	return func(b *nacosResolverBuilder) {
		b.clusters = clusters
	}
}

// Register 使用给定的 Nacos namingClient 注册 gRPC Resolver。
// 一般在 main 中调用一次即可,因为 resolver 是全局注册。
//
// 使用示例:
//
//	resolver.Register(namingClient, resolver.WithGroupName("MY_GROUP"))
func Register(client naming_client.INamingClient, opts ...Option) {
	b := &nacosResolverBuilder{
		client:    client,
		groupName: "DEFAULT_GROUP",
		clusters:  []string{"DEFAULT"},
	}
	for _, o := range opts {
		o(b)
	}
	resolver.Register(b)
}

func NewNacosResolverBuilder(client naming_client.INamingClient, opts ...Option) resolver.Builder {
	b := &nacosResolverBuilder{
		client:    client,
		groupName: "DEFAULT_GROUP",
		clusters:  []string{"DEFAULT"},
	}
	for _, o := range opts {
		o(b)
	}
	return b
}

// nacosResolverBuilder 实现 resolver.Builder 接口
type nacosResolverBuilder struct {
	client    naming_client.INamingClient
	groupName string
	clusters  []string
}

// Scheme 返回 "nacos" 作为 gRPC target scheme
func (b *nacosResolverBuilder) Scheme() string {
	return "nacos"
}

// Build 根据 target 创建 resolver.Resolver
func (b *nacosResolverBuilder) Build(target resolver.Target, cc resolver.ClientConn,
	_ resolver.BuildOptions) (resolver.Resolver, error) {

	serviceName := target.Endpoint()
	if serviceName == "" {
		return nil, fmt.Errorf("nacos resolver: target.Endpoint() is empty, use nacos:///serviceName")
	}

	r := &nacosResolver{
		cc:          cc,
		client:      b.client,
		serviceName: serviceName,
		groupName:   b.groupName,
		clusters:    b.clusters,
		stopCh:      make(chan struct{}),
	}
	r.start()
	return r, nil
}

// nacosResolver 实现 resolver.Resolver 接口
type nacosResolver struct {
	cc          resolver.ClientConn
	client      naming_client.INamingClient
	serviceName string
	groupName   string
	clusters    []string
	stopCh      chan struct{}
}

func (r *nacosResolver) start() {
	// 1. 初次获取:同步拉取健康实例
	r.fetch()

	// 2. 订阅变更:Nacos 推送实例上下线通知
	param := &vo.SubscribeParam{
		ServiceName: r.serviceName,
		GroupName:   r.groupName,
		Clusters:    r.clusters,
		SubscribeCallback: func(instances []model.Instance, err error) {
			if err != nil {
				return
			}
			r.updateAddrs(instances)
		},
	}
	if err := r.client.Subscribe(param); err != nil {
		// 订阅失败不 panic,仅日志记录(初次已拉取到实例)
		return
	}
}

func (r *nacosResolver) fetch() {
	instances, err := r.client.SelectInstances(vo.SelectInstancesParam{
		ServiceName: r.serviceName,
		GroupName:   r.groupName,
		Clusters:    r.clusters,
		HealthyOnly: true,
	})
	if err != nil {
		r.cc.ReportError(err)
		return
	}
	r.updateAddrs(instances)
}

func (r *nacosResolver) updateAddrs(instances []model.Instance) {
	addrs := make([]resolver.Address, 0, len(instances))
	for _, ins := range instances {
		if !ins.Healthy || !ins.Enable {
			continue
		}
		addr := net.JoinHostPort(ins.Ip, strconv.FormatUint(ins.Port, 10))
		addrs = append(addrs, resolver.Address{Addr: addr})
	}
	// if len(addrs) == 0 {
	// 	return
	// }
	r.cc.UpdateState(resolver.State{Addresses: addrs})
}

// ResolveNow 实现 resolver.Resolver 接口,主动触发解析
func (r *nacosResolver) ResolveNow(_ resolver.ResolveNowOptions) {
	r.fetch()
}

// Close 实现 resolver.Resolver 接口,取消订阅
func (r *nacosResolver) Close() {
	select {
	case <-r.stopCh:
	default:
		close(r.stopCh)
	}
	_ = r.client.Unsubscribe(&vo.SubscribeParam{
		ServiceName: r.serviceName,
		GroupName:   r.groupName,
		Clusters:    r.clusters,
	})
}

评论区

  • llx
    #1
    llx
    2026/8/25 16:57:03

    很有技术力,感谢分享

    • Henry
      #1
      Henry
      2026/8/25 17:12:54
      个人认证YuelaiGroup, Software Engineer

      @llx 鉴定为广告!已隐藏你的链接!

    • llx
      #2
      2026/8/25 17:22:02

      @Henry 哎哟不是广告,混个脸熟

    • Henry
      #3
      Henry
      2026/8/25 17:37:59
      个人认证YuelaiGroup, Software Engineer

      @llx 信你一回🫠