gRPC Resolver 实现 Nacos 服务发现
自行实现
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 其实并没有什么用,这里只是为了可读性。
/*
* @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:
/*
* @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 的生命周期即可,某服务示例:
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] 服务实例订阅成功")
封装库
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,
})
}
评论区
很有技术力,感谢分享
湖南省macOSChrome
@llx 鉴定为广告!已隐藏你的链接!
广东省macOSChrome
@Henry 哎哟不是广告,混个脸熟
福建省macOSChrome
@llx 信你一回🫠
北京macOSChrome