grpcclient

package
v0.3.5 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Aug 6, 2026 License: Apache-2.0 Imports: 33 Imported by: 0

Documentation

Index

Constants

View Source
const (
	BalancerName = "p2c_ewma"
)

Variables

This section is empty.

Functions

func Dial

func Dial(target string, opts ...DialOption) (*grpc.ClientConn, error)

Dial 不需要 context 的简化版本

func DialContext

func DialContext(ctx context.Context, target string, opts ...DialOption) (*grpc.ClientConn, error)

DialContext 统一拨号入口,支持两种模式:

直连:grpcclient.DialContext(ctx, "127.0.0.1:8080", grpcclient.WithDirect())
服务发现:grpcclient.DialContext(ctx, "etcd://127.0.0.1:2379/my-service")
          grpcclient.DialContext(ctx, "nacos://127.0.0.1:8848/my-service?env=prod")

func Now

func Now() time.Duration

Types

type ClientFactory

type ClientFactory struct {
	// contains filtered or unexported fields
}

ClientFactory gRPC客户端工厂

func NewClientFactory

func NewClientFactory(discovery discover.Discovery, defaultOpts ...ServiceDiscoveryOption) *ClientFactory

NewClientFactory 创建客户端工厂

func (*ClientFactory) Close

func (f *ClientFactory) Close() error

Close 停止后台 goroutine 并关闭所有连接

func (*ClientFactory) GetClient

func (f *ClientFactory) GetClient(serviceName string, opts ...ServiceDiscoveryOption) *ServiceDiscoveryClient

GetClient 获取指定服务的客户端

func (*ClientFactory) GetDiscovery

func (f *ClientFactory) GetDiscovery() discover.Discovery

GetDiscovery 获取服务发现实例

func (*ClientFactory) WatchAllServices

func (f *ClientFactory) WatchAllServices(ctx context.Context) error

WatchAllServices 监听所有服务的变化

type DialOption

type DialOption func(*dialConfig)

DialOption 拨号选项

func WithCacheInterceptor added in v0.3.0

func WithCacheInterceptor(store mwcache.Store, opts ...mwcache.Option) DialOption

WithCacheInterceptor 接入 gRPC 响应缓存(pkg/middleware/cache): 相同方法 + 相同请求命中缓存时直接返回,跳过网络调用。 拦截器注册在 chain 最外层:缓存命中不经过熔断/重试。

store := cache.NewMemoryStore(256)
conn, _ := grpcclient.DialContext(ctx, target,
    grpcclient.WithCacheInterceptor(store, cache.WithDefaultTTL(5*time.Minute)),
)

func WithCircuitBreakerInterceptor

func WithCircuitBreakerInterceptor(cb *mwcb.CircuitBreaker) DialOption

WithCircuitBreakerInterceptor 接入**请求级**熔断(pkg/middleware/circuitbreaker): 按整体请求成败驱动熔断状态,打开时直接拒绝、快速失败。直连与服务发现两种模式均生效 (经底层 grpc.WithChainUnaryInterceptor 接入)。

与服务发现版的 WithCircuitBreaker(节点级)互补、可叠加:

  • 节点级(WithCircuitBreaker):selectService 选实例时跳过已熔断的节点;
  • 请求级(本项):对逻辑调用整体熔断,直连模式也适用。

gRPC 的**重试**默认已通过 service config 开启(见 DefaultRetryPolicy),无需额外接线。 需要注入其它一元拦截器时用 WithGRPCDialOptions(grpc.WithChainUnaryInterceptor(...))。

func WithDirect

func WithDirect() DialOption

WithDirect 直连模式,target 直接作为 addr 传给 gRPC,不走服务发现。 用法:grpcclient.DialContext(ctx, "127.0.0.1:8080", grpcclient.WithDirect())

func WithEnvironment

func WithEnvironment(env string) DialOption

WithEnvironment 设置环境过滤

func WithGRPCDialOptions

func WithGRPCDialOptions(opts ...grpc.DialOption) DialOption

WithGRPCDialOptions 设置底层 gRPC 连接选项

func WithInsecure

func WithInsecure() DialOption

WithInsecure 明文连接(不加密),通常用于开发或内网可信环境。

func WithLabelFilter

func WithLabelFilter(filter *ServiceLabelFilter) DialOption

WithLabelFilter 设置标签过滤器

func WithLoadBalancer

func WithLoadBalancer(strategy string) DialOption

WithLoadBalancer 设置负载均衡策略(服务发现模式有效)

func WithRegion

func WithRegion(region string) DialOption

WithRegion 设置地域过滤

func WithRegionFilter

func WithRegionFilter(regions, zones, campuses, environments []string) DialOption

WithRegionFilter 设置地域过滤

func WithRegistry

func WithRegistry(registry discover.Discovery) DialOption

WithRegistry 设置服务注册中心

func WithTimeout

func WithTimeout(timeout time.Duration) DialOption

WithTimeout 设置连接超时时间

func WithVersion

func WithVersion(version string) DialOption

WithVersion 只路由到指定版本的实例,用于灰度发布或版本级隔离。 服务端需通过 grpcserver.WithVersion("v2") 将版本写入注册信息。

conn, _ := grpcclient.Dial("etcd://host/svc", grpcclient.WithVersion("v2"))

func WithVersionIn

func WithVersionIn(versions ...string) DialOption

WithVersionIn 路由到版本在给定集合中的实例,支持同时灰度多个版本。

conn, _ := grpcclient.Dial("etcd://host/svc", grpcclient.WithVersionIn("v2", "v3"))

type LoadBalanceStrategy

type LoadBalanceStrategy int

LoadBalanceStrategy 负载均衡策略

const (
	RoundRobin LoadBalanceStrategy = iota
	Random
	WeightedRoundRobin
	LeastConnections
)

type Resolver

type Resolver struct {
	// contains filtered or unexported fields
}

func NewResolver

func NewResolver(cc resolver.ClientConn, serviceName string, discovery discover.Discovery) *Resolver

func (*Resolver) Close

func (r *Resolver) Close()

func (*Resolver) ResolveNow

func (r *Resolver) ResolveNow(opt resolver.ResolveNowOptions)

func (*Resolver) Start

func (r *Resolver) Start()

type RetryPolicy

type RetryPolicy struct {
	// MaxAttempts 最大尝试次数(含首次),范围 [2, 5],gRPC 规范限制上限为 5。
	MaxAttempts int
	// InitialBackoff 首次重试等待时间,格式为 Go duration string,如 "0.1s"。
	InitialBackoff string
	// MaxBackoff 最大退避时间。
	MaxBackoff string
	// BackoffMultiplier 退避倍数,每次重试等待时间乘以该值。
	BackoffMultiplier float64
	// RetryableStatusCodes 触发重试的 gRPC 状态码。
	// 常用:UNAVAILABLE(服务不可用)、RESOURCE_EXHAUSTED(过载)。
	// 注意:DEADLINE_EXCEEDED 和 CANCELLED 不能加入此列表(gRPC 规范禁止)。
	RetryableStatusCodes []string
}

RetryPolicy 对应 gRPC ServiceConfig 中的 retryPolicy 字段。 注入到 grpc.WithDefaultServiceConfig 后,由 gRPC 传输层自动处理,业务代码无感知。

适用场景:

  • 服务端滚动发布期间短暂返回 UNAVAILABLE
  • 网络抖动导致的瞬态连接失败
  • 服务端触发背压返回 RESOURCE_EXHAUSTED

与 Call()/failover 的区别:

  • RetryPolicy 在单个 conn.Invoke() 内部由 gRPC 自动重试,不换节点
  • Call()/failover 每次重试重新 GetClient() 选节点,处理节点彻底不可用的情况
  • 两者互补,建议同时使用

func DefaultRetryPolicy

func DefaultRetryPolicy() RetryPolicy

DefaultRetryPolicy 返回适合大多数微服务场景的默认 retry policy:

  • 最多 3 次尝试(1 次首次 + 2 次重试)
  • 首次重试等 100ms,最长等 1s,指数退避 2 倍
  • 仅对 UNAVAILABLE 重试(最保守,幂等安全)

func (RetryPolicy) WithResourceExhausted

func (p RetryPolicy) WithResourceExhausted() RetryPolicy

WithResourceExhausted 在默认策略基础上追加 RESOURCE_EXHAUSTED 重试, 适用于服务端限流会短暂返回该状态码的场景。

type ServiceClient

type ServiceClient struct {
	*ServiceDiscoveryClient
	// contains filtered or unexported fields
}

ServiceClient 特定服务的客户端包装器

func NewServiceClient

func NewServiceClient(discovery discover.Discovery, serviceName string, opts ...ServiceDiscoveryOption) *ServiceClient

NewServiceClient 创建特定服务的客户端

func (*ServiceClient) Call

func (c *ServiceClient) Call(ctx context.Context, method string, req, resp any, opts ...grpc.CallOption) error

Call 调用服务方法

func (*ServiceClient) Close

func (c *ServiceClient) Close() error

Close 停止后台 goroutine 并关闭所有连接

func (*ServiceClient) NewStream

func (c *ServiceClient) NewStream(ctx context.Context, desc *grpc.StreamDesc, method string, opts ...grpc.CallOption) (grpc.ClientStream, error)

NewStream 创建流

type ServiceDiscoveryClient

type ServiceDiscoveryClient struct {
	// contains filtered or unexported fields
}

ServiceDiscoveryClient 基于服务发现的gRPC客户端

func NewServiceDiscoveryClient

func NewServiceDiscoveryClient(discovery discover.Discovery, serviceName string, opts ...ServiceDiscoveryOption) *ServiceDiscoveryClient

NewServiceDiscoveryClient 创建基于服务发现的客户端

func (*ServiceDiscoveryClient) Call

func (c *ServiceDiscoveryClient) Call(ctx context.Context, method string, req, resp any, opts ...grpc.CallOption) error

Call 调用服务方法,支持指数退避重试(带 ±25% jitter)。 maxRetries 为额外重试次数,0 表示不重试,总调用次数为 maxRetries+1。 context.Canceled / context.DeadlineExceeded 不重试,直接返回。 Call 调用服务方法,支持指数退避重试(带 ±25% jitter)。 maxRetries 为额外重试次数,0 表示不重试,总调用次数为 maxRetries+1。 context.Canceled / context.DeadlineExceeded 不重试,直接返回。 重试链内:失败节点自动加入 bannednodes(本次请求不再重复选)+ 反馈给熔断器(跨请求熔断)。

func (*ServiceDiscoveryClient) Close

func (c *ServiceDiscoveryClient) Close() error

Close 停止后台 goroutine 并关闭所有连接

func (*ServiceDiscoveryClient) GetClient

GetClient 获取一个连接,按负载均衡策略选择实例

func (*ServiceDiscoveryClient) GetServiceInfo

func (c *ServiceDiscoveryClient) GetServiceInfo() []discover.ServiceInfo

GetServiceInfo 获取当前缓存的服务列表

func (*ServiceDiscoveryClient) Start

func (c *ServiceDiscoveryClient) Start(ctx context.Context) (err error)

Start 启动客户端:拉取初始服务列表,启动 watch 和健康检查。 调用 Stop() 或取消传入的 ctx 均可停止后台 goroutine。 Start 是幂等的,多次调用只有第一次生效。

func (*ServiceDiscoveryClient) Stop

func (c *ServiceDiscoveryClient) Stop()

Stop 停止后台 goroutine 并等待它们退出

func (*ServiceDiscoveryClient) WatchServices

func (c *ServiceDiscoveryClient) WatchServices(ctx context.Context) error

WatchServices 监听服务变化,自动更新缓存并关闭失效连接。 回调处理在独立 goroutine 中执行,避免阻塞底层 watcher 事件循环。

type ServiceDiscoveryOption

type ServiceDiscoveryOption func(*ServiceDiscoveryClient)

ServiceDiscoveryOption 服务发现客户端选项

func WithCircuitBreaker

func WithCircuitBreaker(cb governancecb.CircuitBreaker) ServiceDiscoveryOption

WithCircuitBreaker 设置节点级熔断器。selectService 选实例前先过 Available 检查, 跳过已熔断节点;Call 调用结束自动 Report 结果。默认 NoopBreaker(不熔断)。

func WithDiscoveryDialOptions

func WithDiscoveryDialOptions(opts ...grpc.DialOption) ServiceDiscoveryOption

WithDiscoveryDialOptions 设置连接选项

func WithDiscoveryDrainTimeout

func WithDiscoveryDrainTimeout(d time.Duration) ServiceDiscoveryOption

WithDiscoveryDrainTimeout 设置服务实例从发现列表移除后,连接的排空等待时间。 在此期间连接不会被关闭,正在进行的请求有机会完成。默认 5s,设为 0 立即关闭。

func WithDiscoveryFailover

func WithDiscoveryFailover(maxRetries int, retryDelay time.Duration) ServiceDiscoveryOption

WithDiscoveryFailover 设置故障重试

func WithDiscoveryHealthCheck

func WithDiscoveryHealthCheck(enabled bool, interval time.Duration) ServiceDiscoveryOption

WithDiscoveryHealthCheck 设置健康检查

func WithDiscoveryInsecure

func WithDiscoveryInsecure() ServiceDiscoveryOption

WithDiscoveryInsecure 明确声明使用明文连接(不加密)。 生产环境应通过 WithDiscoveryDialOptions(grpc.WithTransportCredentials(...)) 提供 TLS 凭证; 此选项仅用于开发或内网可信环境。

func WithDiscoveryLabelFilter

func WithDiscoveryLabelFilter(filter *ServiceLabelFilter) ServiceDiscoveryOption

WithDiscoveryLabelFilter 设置标签过滤器,替换由 WithDiscoveryRegionFilter 设置的过滤条件。 若需同时使用两种过滤,请在同一个 ServiceLabelFilter 上链式调用后再传入。

func WithDiscoveryRegionFilter

func WithDiscoveryRegionFilter(regions, zones, campuses, environments []string) ServiceDiscoveryOption

WithDiscoveryRegionFilter 追加地域过滤条件(可与 WithDiscoveryLabelFilter 叠加)

func WithDiscoveryRetryPolicy

func WithDiscoveryRetryPolicy(p RetryPolicy) ServiceDiscoveryOption

WithDiscoveryRetryPolicy 设置 gRPC 原生 retry policy,覆盖默认策略。 传入零值 RetryPolicy{} 或空 RetryableStatusCodes 可禁用重试。

示例——对 UNAVAILABLE 和 RESOURCE_EXHAUSTED 均重试,最多 5 次:

WithDiscoveryRetryPolicy(grpcclient.DefaultRetryPolicy().WithResourceExhausted())

func WithDiscoveryStrategy

func WithDiscoveryStrategy(strategy LoadBalanceStrategy) ServiceDiscoveryOption

WithDiscoveryStrategy 设置负载均衡策略

func WithDiscoveryVersionFilter

func WithDiscoveryVersionFilter(versions ...string) ServiceDiscoveryOption

WithDiscoveryVersionFilter 只路由到 version 在给定集合中的实例。 服务端通过 grpcserver.WithVersion("v2") 注册版本信息,客户端用此 Option 过滤。

灰度示例:同时保留 v1(稳定)和 v2(灰度),按流量比例路由见 WithDiscoveryStrategy。

// 只调 v2 实例
client := grpcclient.NewServiceDiscoveryClient(reg, "order-svc",
    grpcclient.WithDiscoveryVersionFilter("v2"),
)

func WithServiceRouter

WithServiceRouter 设置路由过滤层。selectService 选实例前先过 router.Filter, 用于灰度/地域亲和等。默认 NoopRouter(不过滤)。与 WithDiscoveryLabelFilter 区别: labelFilter 在缓存层过滤(refreshServices 时),router 在选实例时过滤(每次 selectService)。

func WithStreamInterceptors

func WithStreamInterceptors(interceptors ...grpc.StreamClientInterceptor) ServiceDiscoveryOption

WithStreamInterceptors 设置流拦截器

func WithUnaryInterceptors

func WithUnaryInterceptors(interceptors ...grpc.UnaryClientInterceptor) ServiceDiscoveryOption

WithUnaryInterceptors 设置一元拦截器

type ServiceLabelFilter

type ServiceLabelFilter struct {
	*selector.LabelFilter
}

ServiceLabelFilter gRPC服务标签过滤器,基于通用的LabelFilter

func NewLabelFilter

func NewLabelFilter() *ServiceLabelFilter

NewLabelFilter 为了向后兼容,保留原来的函数名

func NewServiceLabelFilter

func NewServiceLabelFilter() *ServiceLabelFilter

NewServiceLabelFilter 创建服务标签过滤器

func (*ServiceLabelFilter) Filter

Filter 过滤服务实例

func (*ServiceLabelFilter) WithCampusIn

func (f *ServiceLabelFilter) WithCampusIn(campuses ...string) *ServiceLabelFilter

WithCampusIn 添加园区 in 表达式(便捷方法)

func (*ServiceLabelFilter) WithEnvironmentIn

func (f *ServiceLabelFilter) WithEnvironmentIn(environments ...string) *ServiceLabelFilter

WithEnvironmentIn 添加环境 in 表达式(便捷方法)

func (*ServiceLabelFilter) WithExpression

func (f *ServiceLabelFilter) WithExpression(key string, operator selector.FilterOperator, values ...string) *ServiceLabelFilter

WithExpression 添加表达式匹配

func (*ServiceLabelFilter) WithMatchLabel

func (f *ServiceLabelFilter) WithMatchLabel(key, value string) *ServiceLabelFilter

WithMatchLabel 添加单个精确匹配的标签

func (*ServiceLabelFilter) WithMatchLabels

func (f *ServiceLabelFilter) WithMatchLabels(labels map[string]string) *ServiceLabelFilter

WithMatchLabels 添加精确匹配的标签

func (*ServiceLabelFilter) WithRegionIn

func (f *ServiceLabelFilter) WithRegionIn(regions ...string) *ServiceLabelFilter

WithRegionIn 添加地域 in 表达式(便捷方法)

func (*ServiceLabelFilter) WithVersionIn

func (f *ServiceLabelFilter) WithVersionIn(versions ...string) *ServiceLabelFilter

WithVersionIn 只路由到 version 标签在给定集合中的实例,用于灰度发布。 服务端通过 grpcserver.WithVersion("v2") / webserver.WithVersion("v2") 写入 metadata, 客户端通过 WithVersionIn("v2") 过滤,实现版本级流量隔离。

func (*ServiceLabelFilter) WithZoneIn

func (f *ServiceLabelFilter) WithZoneIn(zones ...string) *ServiceLabelFilter

WithZoneIn 添加可用区 in 表达式(便捷方法)

Directories

Path Synopsis
resolver
consul
Package consul registers a Consul-backed gRPC name resolver.
Package consul registers a Consul-backed gRPC name resolver.
etcd
Package etcd registers an etcd-backed gRPC name resolver.
Package etcd registers an etcd-backed gRPC name resolver.
k8s
Package k8s registers a Kubernetes-backed gRPC name resolver.
Package k8s registers a Kubernetes-backed gRPC name resolver.
nacos
Package nacos registers a Nacos-backed gRPC name resolver.
Package nacos registers a Nacos-backed gRPC name resolver.
polaris
Package polaris registers a Polaris-backed gRPC name resolver.
Package polaris registers a Polaris-backed gRPC name resolver.
Package xds 启用客户端 xDS 支持。
Package xds 启用客户端 xDS 支持。

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL