Documentation
¶
Overview ¶
Package resilience 提供弹性容错功能,用于 enhance 框架。
Package resilience 提供弹性容错功能,用于 enhance 框架。
Package resilience 提供弹性容错功能,用于 enhance 框架。
该模块提供熔断器、负载均衡、限流等弹性容错机制,提升系统的稳定性和可用性。 参考 Resilience4j 的设计理念。
架构设计 ¶
- Breaker: 熔断器接口,防止级联故障
- State: 熔断器状态枚举
- BreakerOption: 熔断器配置选项函数
- Selector: 负载均衡选择器接口
- Registry: 服务注册中心接口
- InstanceInfo: 服务实例信息
- HealthStatus: 健康状态枚举
- ServiceInstance: 服务实例结构
- Balancer: 负载均衡接口
核心功能 ¶
- 熔断器: 支持 OPEN、HALF_OPEN、CLOSED 三种状态
- 负载均衡: 支持轮询、随机、最少连接等策略
- 限流: 支持自适应限流,动态调整速率
- 健康检查: 支持主动和被动健康检查
- 服务注册: 支持服务注册和发现
使用方式 ¶
使用熔断器:
cb := resilience.NewBreaker(
resilience.WithErrorThreshold(0.5),
resilience.WithWaitDuration(30 * time.Second),
)
if err := cb.Allow(); err != nil {
return err // 熔断器打开,快速失败
}
cb.RecordSuccess()
使用负载均衡:
sel := resilience.NewRoundRobinSelector() instance, err := sel.Select(instances)
熔断器状态 ¶
- CLOSED: 正常状态,请求正常通过
- OPEN: 熔断状态,请求直接失败
- HALF_OPEN: 半开状态,允许部分请求通过进行测试
Package resilience 提供弹性容错功能,用于 enhance 框架。
Package resilience 提供弹性容错功能,用于 enhance 框架。
Package resilience 提供弹性容错功能,用于 enhance 框架。
Index ¶
- Constants
- Variables
- type AdaptiveWeight
- func (aw *AdaptiveWeight) GetAllStats() map[string]*BackendStats
- func (aw *AdaptiveWeight) GetStats(backendURL string) (*BackendStats, bool)
- func (aw *AdaptiveWeight) Next(backends []*ServiceInstance) (*ServiceInstance, error)
- func (aw *AdaptiveWeight) RecordConnection(backendURL string, delta int64)
- func (aw *AdaptiveWeight) RecordRequest(backendURL string, responseTimeMs float64, failed bool)
- type BackendStats
- type Balancer
- type Breaker
- type BreakerOption
- type ConsistentHash
- type DefaultBreaker
- type HealthAware
- type HealthStatus
- type IPHash
- type InMemoryRegistry
- func (r *InMemoryRegistry) Deregister(ctx context.Context, info InstanceInfo) error
- func (r *InMemoryRegistry) Discover(ctx context.Context, serviceName string) ([]InstanceInfo, error)
- func (r *InMemoryRegistry) Register(ctx context.Context, info InstanceInfo) error
- func (r *InMemoryRegistry) Watch(ctx context.Context, serviceName string) (<-chan []InstanceInfo, error)
- type InstanceBuilder
- func (b *InstanceBuilder) Build() InstanceInfo
- func (b *InstanceBuilder) Host(host string) *InstanceBuilder
- func (b *InstanceBuilder) ID(id string) *InstanceBuilder
- func (b *InstanceBuilder) Metadata(key, value string) *InstanceBuilder
- func (b *InstanceBuilder) Port(port int) *InstanceBuilder
- func (b *InstanceBuilder) Weight(weight int) *InstanceBuilder
- type InstanceInfo
- type LeastConnections
- type MemoryRegistry
- func (r *MemoryRegistry) Deregister(ctx context.Context, info InstanceInfo) error
- func (r *MemoryRegistry) Discover(ctx context.Context, serviceName string) ([]InstanceInfo, error)
- func (r *MemoryRegistry) Register(ctx context.Context, info InstanceInfo) error
- func (r *MemoryRegistry) Watch(ctx context.Context, serviceName string) (<-chan []InstanceInfo, error)
- type Random
- type Registry
- type RegistryBuilder
- func (b *RegistryBuilder) Address(addr string) *RegistryBuilder
- func (b *RegistryBuilder) Build() (Registry, error)
- func (b *RegistryBuilder) Heartbeat(seconds int) *RegistryBuilder
- func (b *RegistryBuilder) Metadata(key, value string) *RegistryBuilder
- func (b *RegistryBuilder) MustBuild() Registry
- func (b *RegistryBuilder) Registry(reg Registry) *RegistryBuilder
- func (b *RegistryBuilder) Token(token string) *RegistryBuilder
- func (b *RegistryBuilder) Type(typ string) *RegistryBuilder
- type ResponseTimeWeighted
- type RoundRobin
- type Selector
- type SelectorBuilder
- type ServiceInstance
- type State
- type StickySession
- func (ss *StickySession) GetSessionBackend(sessionID string) (*ServiceInstance, bool)
- func (ss *StickySession) GetSessionCount() int
- func (ss *StickySession) Next(backends []*ServiceInstance) (*ServiceInstance, error)
- func (ss *StickySession) NextWithSession(backends []*ServiceInstance, sessionID string) (*ServiceInstance, error)
- func (ss *StickySession) RemoveSession(sessionID string)
- type WeightedRoundRobin
Constants ¶
const ( // DefaultMaxRequests 半开状态下允许的最大试探请求数。 DefaultMaxRequests = 10 // DefaultErrorThreshold 错误率阈值(50%)。 DefaultErrorThreshold = 0.5 // DefaultWaitDuration 打开状态等待时间。 DefaultWaitDuration = 30 * time.Second )
默认熔断器配置常量。
Variables ¶
var ( // ErrCircuitOpen 熔断器处于打开状态。 ErrCircuitOpen error // ErrCircuitHalfOpen 熔断器处于半开状态且已达到最大请求数。 ErrCircuitHalfOpen error // ErrNoInstances 没有可用实例。 ErrNoInstances error // ErrNoBackends 没有可用后端服务。 ErrNoBackends error )
Breaker errors.
Functions ¶
This section is empty.
Types ¶
type AdaptiveWeight ¶
type AdaptiveWeight struct {
// contains filtered or unexported fields
}
AdaptiveWeight 自适应权重负载均衡器 使用 sync.Map 优化读多写少场景的并发性能
func (*AdaptiveWeight) GetAllStats ¶
func (aw *AdaptiveWeight) GetAllStats() map[string]*BackendStats
GetAllStats 获取所有后端统计信息
func (*AdaptiveWeight) GetStats ¶
func (aw *AdaptiveWeight) GetStats(backendURL string) (*BackendStats, bool)
GetStats 获取后端统计信息
func (*AdaptiveWeight) Next ¶
func (aw *AdaptiveWeight) Next(backends []*ServiceInstance) (*ServiceInstance, error)
Next 选择权重最高的后端
func (*AdaptiveWeight) RecordConnection ¶
func (aw *AdaptiveWeight) RecordConnection(backendURL string, delta int64)
RecordConnection 记录连接数变化
func (*AdaptiveWeight) RecordRequest ¶
func (aw *AdaptiveWeight) RecordRequest(backendURL string, responseTimeMs float64, failed bool)
RecordRequest 记录请求统计
type BackendStats ¶
type BackendStats struct {
TotalRequests atomic.Int64
FailedRequests atomic.Int64
ActiveConnections atomic.Int64
LastUpdated atomic.Int64 // Unix 时间戳
// contains filtered or unexported fields
}
BackendStats 后端统计信息 使用 atomic 操作优化计数器,float64 使用 mu 保护
func (*BackendStats) GetAvgResponseTime ¶
func (s *BackendStats) GetAvgResponseTime() float64
GetAvgResponseTime 获取平均响应时间
func (*BackendStats) UpdateAvgResponseTime ¶
func (s *BackendStats) UpdateAvgResponseTime(responseTimeMs float64)
UpdateAvgResponseTime 更新平均响应时间
type Balancer ¶
type Balancer interface {
// Next 选择下一个服务实例。
Next(backends []*ServiceInstance) (*ServiceInstance, error)
}
Balancer 负载均衡接口。
type Breaker ¶
type Breaker interface {
// Allow 检查是否允许请求通过。
Allow() error
// RecordSuccess 记录成功请求。
RecordSuccess()
// RecordFailure 记录失败请求。
RecordFailure()
// State 获取当前状态。
State() State
}
Breaker 熔断器接口。
熔断器模式用于防止级联故障。当依赖服务不可用时, 快速失败而不是等待超时,保护系统资源。
使用示例:
breaker := resilience.NewBreaker(
resilience.WithErrorThreshold(0.5),
resilience.WithWaitDuration(30 * time.Second),
)
if err := breaker.Allow(); err != nil {
return err // 熔断器打开,快速失败
}
// 执行业务逻辑
breaker.RecordSuccess() // 或 RecordFailure()
type BreakerOption ¶
type BreakerOption func(breaker Breaker)
BreakerOption 熔断器配置选项函数。
func WithErrorThreshold ¶
func WithErrorThreshold(threshold float64) BreakerOption
WithErrorThreshold 设置错误率阈值。
func WithWaitDuration ¶
func WithWaitDuration(duration time.Duration) BreakerOption
WithWaitDuration 设置打开状态等待时间。
type ConsistentHash ¶
type ConsistentHash struct {
// contains filtered or unexported fields
}
ConsistentHash 一致性哈希负载均衡器 无状态设计,无需锁保护
func NewConsistentHash ¶
func NewConsistentHash(replicas ...int) *ConsistentHash
NewConsistentHash 创建一致性哈希负载均衡器
func (*ConsistentHash) Next ¶
func (ch *ConsistentHash) Next(backends []*ServiceInstance) (*ServiceInstance, error)
Next 使用一致性哈希选择后端
func (*ConsistentHash) NextByKey ¶
func (ch *ConsistentHash) NextByKey(backends []*ServiceInstance, key string) (*ServiceInstance, error)
NextByKey 根据键选择后端
type DefaultBreaker ¶
type DefaultBreaker struct {
// contains filtered or unexported fields
}
DefaultBreaker 默认熔断器实现(已弃用,请使用 Breaker 接口)。
使用 atomic 操作优化计数器的并发性能, 适用于高并发场景下的请求保护。
配置参数(只读,无需锁保护):
- maxRequests: 半开状态下允许的最大试探请求数
- errorThreshold: 错误率阈值(0.0-1.0),超过此值触发熔断
- waitDuration: 打开状态等待时间,之后进入半开状态
状态管理(使用 atomic 操作):
- state: 当前熔断器状态
- errors/successes/requests: 请求统计
- lastStateChange: 上次状态变更时间
- halfOpenRequests: 半开状态已处理请求数
type HealthAware ¶
type HealthAware struct {
// contains filtered or unexported fields
}
HealthAware 健康感知负载均衡器
func (*HealthAware) GetFailureCount ¶
func (ha *HealthAware) GetFailureCount(backendURL string) int
GetFailureCount 获取失败计数
func (*HealthAware) Next ¶
func (ha *HealthAware) Next(backends []*ServiceInstance) (*ServiceInstance, error)
Next 选择健康的后端
func (*HealthAware) RecordFailure ¶
func (ha *HealthAware) RecordFailure(backendURL string)
RecordFailure 记录后端失败
func (*HealthAware) RecordSuccess ¶
func (ha *HealthAware) RecordSuccess(backendURL string)
RecordSuccess 记录后端成功
type HealthStatus ¶
type HealthStatus int
HealthStatus 健康状态。
const ( // HealthUp 健康。 HealthUp HealthStatus = iota // HealthDown 不健康。 HealthDown // HealthUnknown 未知。 HealthUnknown )
健康状态常量。
type IPHash ¶
type IPHash struct{}
IPHash IP 哈希负载均衡器 无状态设计,无需锁保护
func (*IPHash) Next ¶
func (ih *IPHash) Next(backends []*ServiceInstance) (*ServiceInstance, error)
Next 选择后端(不带 IP 信息)
func (*IPHash) NextByIP ¶
func (ih *IPHash) NextByIP(backends []*ServiceInstance, clientIP string) (*ServiceInstance, error)
NextByIP 根据客户端 IP 选择后端
type InMemoryRegistry ¶
type InMemoryRegistry struct {
// contains filtered or unexported fields
}
InMemoryRegistry 内存中的服务注册中心实现。 适用于测试和开发环境,不支持持久化。
func NewInMemoryRegistry ¶
func NewInMemoryRegistry() *InMemoryRegistry
NewInMemoryRegistry 创建内存注册中心。
func (*InMemoryRegistry) Deregister ¶
func (r *InMemoryRegistry) Deregister(ctx context.Context, info InstanceInfo) error
Deregister 注销服务实例。
func (*InMemoryRegistry) Discover ¶
func (r *InMemoryRegistry) Discover(ctx context.Context, serviceName string) ([]InstanceInfo, error)
Discover 发现服务实例。
func (*InMemoryRegistry) Register ¶
func (r *InMemoryRegistry) Register(ctx context.Context, info InstanceInfo) error
Register 注册服务实例。
func (*InMemoryRegistry) Watch ¶
func (r *InMemoryRegistry) Watch(ctx context.Context, serviceName string) (<-chan []InstanceInfo, error)
Watch 监听服务实例变更。
type InstanceBuilder ¶
type InstanceBuilder struct {
// contains filtered or unexported fields
}
InstanceBuilder 实例构建器,简化服务实例创建
func NewInstanceBuilder ¶
func NewInstanceBuilder(serviceName string) *InstanceBuilder
NewInstanceBuilder 创建实例构建器
func (*InstanceBuilder) Host ¶
func (b *InstanceBuilder) Host(host string) *InstanceBuilder
Host 设置主机地址
func (*InstanceBuilder) Metadata ¶
func (b *InstanceBuilder) Metadata(key, value string) *InstanceBuilder
Metadata 添加元数据
func (*InstanceBuilder) Weight ¶
func (b *InstanceBuilder) Weight(weight int) *InstanceBuilder
Weight 设置权重
type InstanceInfo ¶
type InstanceInfo struct {
// ServiceName 服务名称。
ServiceName string
// ID 实例唯一标识。
ID string
// Host 主机地址。
Host string
// Port 端口号。
Port int
// Weight 负载均衡权重。
Weight int
// Healthy 是否健康。
Healthy bool
// Metadata 扩展元数据。
Metadata map[string]string
}
InstanceInfo 注册中心中的服务实例信息。
包含服务名、实例 ID、网络地址、权重、健康状态和元数据。
type LeastConnections ¶
type LeastConnections struct{}
LeastConnections 最少连接负载均衡器
func NewLeastConnections ¶
func NewLeastConnections() *LeastConnections
NewLeastConnections 创建最少连接负载均衡器
func (*LeastConnections) Next ¶
func (lc *LeastConnections) Next(backends []*ServiceInstance) (*ServiceInstance, error)
Next 选择连接数最少的后端
type MemoryRegistry ¶
type MemoryRegistry struct {
// contains filtered or unexported fields
}
MemoryRegistry 内存注册中心(简单实现)
func (*MemoryRegistry) Deregister ¶
func (r *MemoryRegistry) Deregister(ctx context.Context, info InstanceInfo) error
func (*MemoryRegistry) Discover ¶
func (r *MemoryRegistry) Discover(ctx context.Context, serviceName string) ([]InstanceInfo, error)
func (*MemoryRegistry) Register ¶
func (r *MemoryRegistry) Register(ctx context.Context, info InstanceInfo) error
func (*MemoryRegistry) Watch ¶
func (r *MemoryRegistry) Watch(ctx context.Context, serviceName string) (<-chan []InstanceInfo, error)
type Random ¶
type Random struct{}
Random 随机负载均衡器
func (*Random) Next ¶
func (r *Random) Next(backends []*ServiceInstance) (*ServiceInstance, error)
Next 随机选择一个后端
type Registry ¶
type Registry interface {
// Register 注册服务实例到注册中心。
Register(ctx context.Context, info InstanceInfo) error
// Deregister 从注册中心注销服务实例。
Deregister(ctx context.Context, info InstanceInfo) error
// Discover 发现指定服务的所有实例。
Discover(ctx context.Context, serviceName string) ([]InstanceInfo, error)
// Watch 监听服务实例变更,返回实例列表变更通道。
Watch(ctx context.Context, serviceName string) (<-chan []InstanceInfo, error)
}
Registry 服务注册中心接口。
提供服务实例的注册、注销、发现和监听功能。
使用示例:
var reg resilience.Registry
err := reg.Register(ctx, resilience.InstanceInfo{
ServiceName: "user-service",
ID: "192.168.1.1:8080",
Host: "192.168.1.1",
Port: 8080,
})
instances, err := reg.Discover(ctx, "user-service")
type RegistryBuilder ¶
type RegistryBuilder struct {
// contains filtered or unexported fields
}
RegistryBuilder 注册中心构建器,支持链式配置
func (*RegistryBuilder) Address ¶
func (b *RegistryBuilder) Address(addr string) *RegistryBuilder
Address 设置注册中心地址
func (*RegistryBuilder) Heartbeat ¶
func (b *RegistryBuilder) Heartbeat(seconds int) *RegistryBuilder
Heartbeat 设置心跳间隔(秒)
func (*RegistryBuilder) Metadata ¶
func (b *RegistryBuilder) Metadata(key, value string) *RegistryBuilder
Metadata 添加元数据
func (*RegistryBuilder) MustBuild ¶
func (b *RegistryBuilder) MustBuild() Registry
MustBuild 构建注册中心,失败则panic
func (*RegistryBuilder) Registry ¶
func (b *RegistryBuilder) Registry(reg Registry) *RegistryBuilder
Registry 使用自定义注册中心
func (*RegistryBuilder) Token ¶
func (b *RegistryBuilder) Token(token string) *RegistryBuilder
Token 设置认证令牌
func (*RegistryBuilder) Type ¶
func (b *RegistryBuilder) Type(typ string) *RegistryBuilder
Type 设置注册中心类型
type ResponseTimeWeighted ¶
type ResponseTimeWeighted struct {
// contains filtered or unexported fields
}
ResponseTimeWeighted 响应时间加权负载均衡器 使用 sync.Map 优化读多写少场景的并发性能
func NewResponseTimeWeighted ¶
func NewResponseTimeWeighted(decay ...float64) *ResponseTimeWeighted
NewResponseTimeWeighted 创建响应时间加权负载均衡器
func (*ResponseTimeWeighted) GetAvgResponseTime ¶
func (rtw *ResponseTimeWeighted) GetAvgResponseTime(backendURL string) (float64, bool)
GetAvgResponseTime 获取后端平均响应时间
func (*ResponseTimeWeighted) Next ¶
func (rtw *ResponseTimeWeighted) Next(backends []*ServiceInstance) (*ServiceInstance, error)
Next 选择响应时间最短的后端
func (*ResponseTimeWeighted) RecordResponseTime ¶
func (rtw *ResponseTimeWeighted) RecordResponseTime(backendURL string, responseTimeMs float64)
RecordResponseTime 记录后端响应时间
type RoundRobin ¶
type RoundRobin struct {
// contains filtered or unexported fields
}
RoundRobin 轮询负载均衡器
func (*RoundRobin) Next ¶
func (rr *RoundRobin) Next(backends []*ServiceInstance) (*ServiceInstance, error)
Next 选择下一个后端
type Selector ¶
type Selector interface {
// Select 选择服务实例。
Select(instances []InstanceInfo) (InstanceInfo, error)
}
Selector 负载均衡选择器接口。
从一组服务实例中按策略选出一个目标实例。
使用示例:
sel := resilience.NewRoundRobinSelector() inst, err := sel.Select(instances)
var ( // RandomSelect 随机选择器(已弃用,请使用 NewRandomSelector)。 // 保留用于向后兼容。 RandomSelect Selector = NewRandomSelector() // RoundRobinSelect 轮询选择器(已弃用,请使用 NewRoundRobinSelector)。 // 保留用于向后兼容。 RoundRobinSelect Selector = NewRoundRobinSelector() )
func NewWeightedRandomSelector ¶
func NewWeightedRandomSelector() Selector
NewWeightedRandomSelector 创建加权随机选择器。
type SelectorBuilder ¶
type SelectorBuilder struct {
// contains filtered or unexported fields
}
SelectorBuilder 选择器构建器
func (*SelectorBuilder) MustBuild ¶
func (b *SelectorBuilder) MustBuild() Selector
MustBuild 构建选择器,失败则panic
func (*SelectorBuilder) Strategy ¶
func (b *SelectorBuilder) Strategy(strategy string) *SelectorBuilder
Strategy 设置选择策略(random、round_robin、weighted)
func (*SelectorBuilder) Weight ¶
func (b *SelectorBuilder) Weight(instanceID string, weight int) *SelectorBuilder
Weight 添加权重配置
type ServiceInstance ¶
type ServiceInstance struct {
// ID 服务实例 ID。
ID string
// URL 服务地址。
URL string
// Weight 权重。
Weight int
// Metadata 元数据。
Metadata map[string]string
// Health 健康状态。
Health HealthStatus
// Active 活跃连接数。
Active int64
}
ServiceInstance 服务实例。
func SortByResponseTime ¶
func SortByResponseTime(backends []*ServiceInstance, stats map[string]*BackendStats) []*ServiceInstance
SortByResponseTime 按响应时间排序后端
type State ¶
type State int32
State 熔断器状态。
熔断器有三种状态:
- Closed(关闭):正常处理请求,统计错误率
- Open(打开):快速失败,拒绝所有请求
- HalfOpen(半开):允许少量请求试探,评估是否恢复
type StickySession ¶
type StickySession struct {
// contains filtered or unexported fields
}
StickySession 会话保持负载均衡器
func NewStickySession ¶
func NewStickySession(sessionCookieName string) *StickySession
NewStickySession 创建会话保持负载均衡器
func (*StickySession) GetSessionBackend ¶
func (ss *StickySession) GetSessionBackend(sessionID string) (*ServiceInstance, bool)
GetSessionBackend 获取会话绑定的后端
func (*StickySession) GetSessionCount ¶
func (ss *StickySession) GetSessionCount() int
GetSessionCount 获取会话数量
func (*StickySession) Next ¶
func (ss *StickySession) Next(backends []*ServiceInstance) (*ServiceInstance, error)
Next 选择后端(不带会话信息)
func (*StickySession) NextWithSession ¶
func (ss *StickySession) NextWithSession(backends []*ServiceInstance, sessionID string) (*ServiceInstance, error)
NextWithSession 根据会话 ID 选择后端
func (*StickySession) RemoveSession ¶
func (ss *StickySession) RemoveSession(sessionID string)
RemoveSession 移除会话绑定
type WeightedRoundRobin ¶
type WeightedRoundRobin struct {
// contains filtered or unexported fields
}
WeightedRoundRobin 加权轮询负载均衡器
func NewWeightedRoundRobin ¶
func NewWeightedRoundRobin() *WeightedRoundRobin
NewWeightedRoundRobin 创建加权轮询负载均衡器
func (*WeightedRoundRobin) Next ¶
func (wrr *WeightedRoundRobin) Next(backends []*ServiceInstance) (*ServiceInstance, error)
Next 根据权重选择后端