resilience

package
v0.0.1 Latest Latest
Warning

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

Go to latest
Published: Jul 19, 2026 License: MIT Imports: 10 Imported by: 0

README

resilience 包 — 弹性与容错

所属层级: Infrastructure Layer
设计理念: 服务弹性,负载均衡,熔断容错
设计灵感: Spring Cloud + Netflix Hystrix + Resilience4j

概述

resilience 包提供服务弹性、负载均衡和容错能力,确保分布式系统的高可用性。

核心功能
功能 说明
服务注册与发现 抽象注册中心接口,支持服务实例注册、注销、发现和监听
负载均衡 10+ 种负载均衡策略(随机、轮询、加权、最少连接、一致性哈希等)
熔断器 熔断器模式实现,支持状态转换(Closed/Open/Half-Open)和故障恢复

核心接口

服务注册与发现
Registry 接口
type Registry interface {
    Register(ctx context.Context, info InstanceInfo) error
    Deregister(ctx context.Context, info InstanceInfo) error
    Discover(ctx context.Context, serviceName string) ([]InstanceInfo, error)
    Watch(ctx context.Context, serviceName string) (<-chan []InstanceInfo, error)
}
InstanceInfo 结构体
字段 类型 说明
ServiceName string 服务名称
ID string 实例唯一标识
Host string 实例主机地址
Port int 实例端口号
Weight int 权重,用于加权选择策略
Healthy bool 健康状态
Metadata map[string]string 附加元数据(版本号、区域等)
使用示例
var reg resilience.Registry // 具体实现由注册中心提供(如 etcd、consul、nacos)

// 注册实例
err := reg.Register(ctx, resilience.InstanceInfo{
    ServiceName: "user-service",
    ID:          "192.168.1.1:8080",
    Host:        "192.168.1.1",
    Port:        8080,
    Weight:      10,
    Healthy:     true,
    Metadata:    map[string]string{"version": "v1.0"},
})

// 发现服务
instances, err := reg.Discover(ctx, "user-service")

// 监听变化
ch, err := reg.Watch(ctx, "user-service")
for instances := range ch {
    // 处理实例变更
}
负载均衡
Selector 接口
type Selector interface {
    Select(instances []*InstanceInfo) (*InstanceInfo, error)
}
内置负载均衡策略
策略 说明 适用场景
RandomSelector 随机选择 简单无状态场景
RoundRobinSelector 轮询选择 后端性能相近
WeightedRandomSelector 加权随机 后端性能不均
WeightedRoundRobinSelector 加权轮询 后端性能不均且需均匀分布
LeastConnSelector 最少连接 长连接场景
ResponseTimeSelector 响应时间优先 追求最低延迟
AdaptiveSelector 自适应权重 动态调整负载
HashSelector 一致性哈希 缓存亲和性
StickySelector 粘性会话 会话保持
HealthSelector 健康状态优先 高可用场景
IPHashSelector IP 哈希 客户端亲和性
使用示例
// 轮询选择
selector := resilience.NewRoundRobinSelector()
selected, err := selector.Select(instances)

// 加权随机
selector := resilience.NewWeightedRandomSelector()
selected, err := selector.Select(instances)

// 最少连接
selector := resilience.NewLeastConnSelector()
selected, err := selector.Select(instances)
熔断器
CircuitBreaker 接口
type CircuitBreaker interface {
    Execute(fn func() (any, error)) (any, error)
    IsOpen() bool
    RecordSuccess()
    RecordFailure()
}
熔断器状态
状态 常量 说明
Closed StateClosed 关闭状态,正常处理请求
Open StateOpen 打开状态,快速失败
HalfOpen StateHalfOpen 半开状态,尝试恢复
状态转换
Closed ──(失败率超过阈值)──> Open
Open ──(等待时间过后)──> HalfOpen
HalfOpen ──(请求成功)──> Closed
HalfOpen ──(请求失败)──> Open
使用示例
// 基本使用
breaker := resilience.NewCircuitBreaker(
    resilience.WithFailureThreshold(5),
    resilience.WithSuccessThreshold(3),
    resilience.WithTimeout(30 * time.Second),
)

result, err := breaker.Execute(func() (any, error) {
    return callExternalService()
})

// 带降级逻辑
breaker := resilience.NewCircuitBreaker(
    resilience.WithFailureThreshold(5),
    resilience.WithFallback(func() (any, error) {
        // 降级逻辑
        return cachedResult, nil
    }),
)

快速开始

服务注册与发现
package main

import (
    "context"
    "fmt"
    "github.com/xudefa/enhance/resilience"
)

func main() {
    // 创建注册中心(以内存实现为例)
    reg := resilience.NewInMemoryRegistry()

    // 注册服务实例
    reg.Register(context.Background(), resilience.InstanceInfo{
        ServiceName: "user-service",
        ID:          "instance-1",
        Host:        "192.168.1.1",
        Port:        8080,
        Weight:      10,
        Healthy:     true,
    })

    // 发现服务
    instances, _ := reg.Discover(context.Background(), "user-service")
    for _, instance := range instances {
        fmt.Printf("Found instance: %s at %s:%d\n", instance.ID, instance.Host, instance.Port)
    }
}

API 参考

负载均衡
// 创建服务实例列表
instances := []*resilience.InstanceInfo{
    {ID: "instance-1", Host: "192.168.1.1", Port: 8080, Weight: 10, Healthy: true},
    {ID: "instance-2", Host: "192.168.1.2", Port: 8080, Weight: 20, Healthy: true},
    {ID: "instance-3", Host: "192.168.1.3", Port: 8080, Weight: 15, Healthy: true},
}

// 使用加权随机选择器
selector := resilience.NewWeightedRandomSelector()
selected, err := selector.Select(instances)
if err != nil {
    // 处理错误
}
fmt.Printf("Selected instance: %s\n", selected.ID)
熔断器
// 创建熔断器
breaker := resilience.NewCircuitBreaker(
    resilience.WithFailureThreshold(5),
    resilience.WithSuccessThreshold(3),
    resilience.WithTimeout(30 * time.Second),
    resilience.WithFallback(func() (any, error) {
        return "fallback response", nil
    }),
)

// 执行受保护的调用
result, err := breaker.Execute(func() (any, error) {
    return http.Get("https://api.example.com/data")
})

if err != nil {
    // 处理错误(可能是熔断器打开)
}

使用示例

场景 1: 微服务调用链
type OrderService struct {
    registry   resilience.Registry
    selector   resilience.Selector
    breaker    resilience.CircuitBreaker
}

func (s *OrderService) CreateOrder(ctx context.Context, order *Order) error {
    // 发现库存服务
    instances, err := s.registry.Discover(ctx, "inventory-service")
    if err != nil {
        return err
    }

    // 选择实例
    instance, err := s.selector.Select(instances)
    if err != nil {
        return err
    }

    // 使用熔断器调用库存服务
    _, err = s.breaker.Execute(func() (any, error) {
        return http.Post(fmt.Sprintf("http://%s:%d/reserve", instance.Host, instance.Port), "application/json", order)
    })

    return err
}
场景 2: 动态服务发现与负载均衡
func (s *GatewayService) StartLoadBalancer() {
    // 监听服务变化
    ch, _ := s.registry.Watch(context.Background(), "user-service")
    
    go func() {
        for instances := range ch {
            // 更新本地实例列表
            s.updateInstances(instances)
        }
    }()
}

func (s *GatewayService) HandleRequest(w http.ResponseWriter, r *http.Request) {
    // 选择实例
    instance, err := s.selector.Select(s.instances)
    if err != nil {
        http.Error(w, "No available instances", 503)
        return
    }

    // 转发请求
    proxy := httputil.NewSingleHostReverseProxy(&url.URL{
        Scheme: "http",
        Host:   fmt.Sprintf("%s:%d", instance.Host, instance.Port),
    })
    proxy.ServeHTTP(w, r)
}

最佳实践

1. 选择合适的负载均衡策略
// ✅ 推荐:根据场景选择策略
// 后端性能相近:使用轮询
selector := resilience.NewRoundRobinSelector()

// 后端性能不均:使用加权随机
selector := resilience.NewWeightedRandomSelector()

// 长连接场景:使用最少连接
selector := resilience.NewLeastConnSelector()

// 缓存亲和性:使用一致性哈希
selector := resilience.NewHashSelector()

// ⚠️ 不推荐:不考虑场景随意选择
selector := resilience.NewRandomSelector() // 可能导致负载不均
2. 熔断器配置建议
// ✅ 推荐:合理配置熔断器
breaker := resilience.NewCircuitBreaker(
    resilience.WithFailureThreshold(5),      // 失败 5 次后打开
    resilience.WithSuccessThreshold(3),      // 成功 3 次后关闭
    resilience.WithTimeout(30*time.Second),  // 打开 30 秒后尝试半开
    resilience.WithFallback(fallbackFunc),   // 提供降级逻辑
)

// ⚠️ 不推荐:不配置降级逻辑
breaker := resilience.NewCircuitBreaker(
    resilience.WithFailureThreshold(5),
    // 没有 Fallback,熔断后直接返回错误
)
3. 服务健康检查
// ✅ 推荐:定期检查服务健康状态
func (s *ServiceMonitor) StartHealthCheck() {
    ticker := time.NewTicker(10 * time.Second)
    go func() {
        for range ticker.C {
            for _, instance := range s.instances {
                if !s.checkHealth(instance) {
                    instance.Healthy = false
                    s.registry.Deregister(context.Background(), instance)
                }
            }
        }
    }()
}

// ⚠️ 不推荐:不检查健康状态
// 可能导致请求发送到不健康的实例
4. 与依赖注入集成
// ✅ 推荐:将弹性组件注册为 Bean
container.Register(
    reflect.TypeOf(&resilience.CircuitBreaker{}),
    core.Bean(createCircuitBreaker()),
    core.Singleton(),
)

container.Register(
    reflect.TypeOf(&resilience.Selector{}),
    core.Bean(createSelector()),
    core.Singleton(),
)

// 注入使用
type OrderService struct {
    Breaker  resilience.CircuitBreaker `inject:"circuitBreaker"`
    Selector resilience.Selector       `inject:"selector"`
}
5. 设计要点
  • 参考 Spring Cloud 设计理念
  • 接口抽象,易于扩展
  • 零外部依赖(仅使用 Go 标准库)
  • 并发安全设计

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

View Source
const (
	// DefaultMaxRequests 半开状态下允许的最大试探请求数。
	DefaultMaxRequests = 10
	// DefaultErrorThreshold 错误率阈值(50%)。
	DefaultErrorThreshold = 0.5
	// DefaultWaitDuration 打开状态等待时间。
	DefaultWaitDuration = 30 * time.Second
)

默认熔断器配置常量。

Variables

View Source
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 NewAdaptiveWeight

func NewAdaptiveWeight() *AdaptiveWeight

NewAdaptiveWeight 创建自适应权重负载均衡器

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()

func NewBreaker

func NewBreaker(opts ...BreakerOption) Breaker

NewBreaker 创建熔断器。

type BreakerOption

type BreakerOption func(breaker Breaker)

BreakerOption 熔断器配置选项函数。

func WithErrorThreshold

func WithErrorThreshold(threshold float64) BreakerOption

WithErrorThreshold 设置错误率阈值。

func WithMaxRequests

func WithMaxRequests(max int) BreakerOption

WithMaxRequests 设置半开状态最大请求数。

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 NewHealthAware

func NewHealthAware(inner Balancer) *HealthAware

NewHealthAware 创建健康感知负载均衡器

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 NewIPHash

func NewIPHash() *IPHash

NewIPHash 创建 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) Build

func (b *InstanceBuilder) Build() InstanceInfo

Build 构建实例信息

func (*InstanceBuilder) Host

func (b *InstanceBuilder) Host(host string) *InstanceBuilder

Host 设置主机地址

func (*InstanceBuilder) ID

ID 设置实例ID

func (*InstanceBuilder) Metadata

func (b *InstanceBuilder) Metadata(key, value string) *InstanceBuilder

Metadata 添加元数据

func (*InstanceBuilder) Port

func (b *InstanceBuilder) Port(port int) *InstanceBuilder

Port 设置端口

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 NewMemoryRegistry

func NewMemoryRegistry() *MemoryRegistry

NewMemoryRegistry 创建内存注册中心

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 NewRandom

func NewRandom() *Random

NewRandom 创建随机负载均衡器

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 NewRegistryBuilder

func NewRegistryBuilder() *RegistryBuilder

NewRegistryBuilder 创建注册中心构建器

func (*RegistryBuilder) Address

func (b *RegistryBuilder) Address(addr string) *RegistryBuilder

Address 设置注册中心地址

func (*RegistryBuilder) Build

func (b *RegistryBuilder) Build() (Registry, error)

Build 构建注册中心

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 NewRoundRobin

func NewRoundRobin() *RoundRobin

NewRoundRobin 创建轮询负载均衡器

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 NewRandomSelector

func NewRandomSelector() Selector

NewRandomSelector 创建随机选择器。

func NewRoundRobinSelector

func NewRoundRobinSelector() Selector

NewRoundRobinSelector 创建轮询选择器。

func NewWeightedRandomSelector

func NewWeightedRandomSelector() Selector

NewWeightedRandomSelector 创建加权随机选择器。

type SelectorBuilder

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

SelectorBuilder 选择器构建器

func NewSelectorBuilder

func NewSelectorBuilder() *SelectorBuilder

NewSelectorBuilder 创建选择器构建器

func (*SelectorBuilder) Build

func (b *SelectorBuilder) Build() (Selector, error)

Build 构建选择器

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(半开):允许少量请求试探,评估是否恢复
const (
	// StateClosed 关闭状态,正常处理请求。
	StateClosed State = iota
	// StateOpen 打开状态,快速失败。
	StateOpen
	// StateHalfOpen 半开状态,尝试恢复。
	StateHalfOpen
)

熔断器状态常量。

func (State) String

func (s State) String() string

String 返回状态的字符串表示。

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 根据权重选择后端

Jump to

Keyboard shortcuts

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