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