gRPC Resolver 实现 Nacos 订阅

2026/8/10 16:07:01Henry3 阅读0 点赞0 评论

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] 服务实例订阅成功")

评论区