go 进阶 限流相关: 二. hystrix-go 类库
·
目录
一. hystrix-go
- hystrix-go 是基于Java的Hystrix类似的执行语义实现的一个集成了流量统计, 限流、熔断、降级等功能的三方类库
- 使用时需要执行"go get github.com/afex/hystrix-go/hystrix" 命令将hystrix-go拉取到本地
- hystrix的设计原则
- 防止任何单个依赖服务耗尽所有用户线程
- 直接响应失败,而不是一直等待
- 提供错误返回接口,而不是让用户线程直接处理依赖服务抛出的异常
- 使用隔离或熔断技术来降低并限制单个依赖对整个系统造成的影响
- 熔断器三种状态
- 关闭状态: 服务正常,维护失败率统计,达到失败率阈值时,转到开启状态。
- 开启状态: 服务异常,调用fallback函数,一段时间后,进入半开启状态。
- 半开启状态: 尝试恢复服务,失败率高于阈值,进入开启状态,低于阈值,进入关闭状态。
- hystrix-go 基本使用示例
- 创建hystrix流量统计服务器
- 设置流控规则
- 创建熔断器
- 熔断器与业务接口绑定,监控指定业务接口流量
package main
import (
"errors"
"github.com/afex/hystrix-go/hystrix"
"log"
"net/http"
"testing"
"time"
)
func Test_main(t *testing.T) {
//1.创建hystrix流量统计服务器
//(实际可以理解为启动一个专门的hystrix服务用来管理熔断器,
//通过熔断器监控指定的业务接口的流量情况,当启动此处的hystrix流控服务后可以访问这个地址,
//例如此处访问:127.0.0.1:8074 就可以拿到监控的访问情况)
hystrixStreamHandler := hystrix.NewStreamHandler()
hystrixStreamHandler.Start()
go http.ListenAndServe(":8074", hystrixStreamHandler)
//2.配置限流规则hystrix.CommandConfig
commandConfig := hystrix.CommandConfig{
Timeout: 1000, //单次请求超时时间,默认时间是1000毫秒
MaxConcurrentRequests: 1, // 最大并发量,默认值是10(注意此处并不是设置为1就是1)
SleepWindow: 5000, // 熔断后多久去尝试服务是否可用,默认值是5000毫秒(熔断器打开到半打开的时间)
RequestVolumeThreshold: 1, //一个统计窗口10秒内请求数量。达到这个请求数量后才去判断是否要开启熔断,默认值是20(比如10秒内接到了11个请求只超过了1个,当前为1没有超过熔断限制,则不熔断)
ErrorPercentThreshold: 1, //错误百分比,默认值是50(当错误百分比超过这个限制时则进行熔断)
}
//3.设置熔断器
//第一个参数:当前创建的熔断器名称
//第二个参数: hystrix.CommandConfig配置的限流规则
hystrix.ConfigureCommand("aaa", commandConfig)
//4.hystrix.Do()同步: hystrix-go与业务方法的同步整合
//假设当前有100个请求依次进来
for i := 0; i < 100; i++ {
//5.通过hystrix.Do()对业务方法实现流控
//hystrix.Do需要三个参数
//参数一: 指定使用哪个熔断器,设置对应的熔断器名称,例如上面创建了"aaa"熔断器,当前则可以设置"aaa"
//参数二: 需要被流控的业务方法
//参数三: 降级方法,当代码返回一个错误时,或者当它基于各种健康检查无法完成时,就会触发此事件
//对应Do()方法hystrix还提供了一个Doc(),Do()内部实际调用的就是Doc()
//Doc()中多了一个Context,思考一下能否使用Context传递参数与接收响应
err := hystrix.Do("aaa", func() error {
//test case 1 并发测试
if i == 0 {
return errors.New("service error")
}
//test case 2 超时测试
//time.Sleep(2 * time.Second)
log.Println("do services")
return nil
}, nil)
//6.hystrix.Do()执行完毕后返回的error,该error会影响熔断
//返回的error会记录到熔断器的错误百分比中,当超过设置阈值则触发熔断,
//熔断后当到达熔断器设置的SleepWindow时间后,开始尝试恢复
if err != nil {
log.Println("hystrix err:" + err.Error())
time.Sleep(1 * time.Second)
log.Println("sleep 1 second")
}
}
time.Sleep(100 * time.Second)
}
// 异步调用使用 hystrix.Go
func test() {
//7.hystrix.Go()异步: hystrix-go与业务方法的异步整合
//假设当前有100个请求依次进来
for i := 0; i < 100; i++ {
//接收响应通道
output := make(chan bool, 1)
//与hystrix.DO()方法大体相同,不同的是Go()时,内部业务方法是异步调用的
//并且返回值是一个chan error,调用hystrix.Go就像启动一个goroutine
//会收到一个可以select监控的error channel,
//对应Go()方法hystrix还提供了一个Goc(),Go()内部实际调用的就是Goc()
//Goc()中多了一个Context,思考一下能否使用Context传递参数与接收响应
chanErr := hystrix.Go("aaa", func() error {
//test case 1 并发测试
if i == 0 {
return errors.New("service error")
}
//test case 2 超时测试
//time.Sleep(2 * time.Second)
log.Println("do services")
return nil
}, nil)
//接收响应
select {
case out := <-output:
// success
fmt.Println(out)
case err := <-chanErr:
// failure
fmt.Println(err)
}
}
}
- Do内部调用的Doc方法,Go内部调用的是Goc方法,在Doc方法内部最终调用的还是Goc方法,只是在Doc方法内做了同步逻辑
- 实际可以理解为启动了一个专门的hystrix服务用来管理熔断器,通过熔断器监控指定的业务接口的流量情况,当启动此处的hystrix流控服务后可以访问这个地址,例如此处访问:127.0.0.1:8074 就可以拿到监控的访问情况
二. hystrix-go 原理相关
- 参考转载博客1 ,参考转载博客2
- hystrix-go中Do和Go方法内部都是调用了hystrix.GoC方法,只是Do方法处理了异步的过程
- 可以把hystrix-go分为: 初始化,执行,执行过程中的状态上报与开关判断几个位置去理解
初始化
- 在我们使用hystrix-go时,需要先创建hystrix.CommandConfig封装限流规则,如不不创建则使用默认的限流规则,查看hystrix.CommandConfig结构体属性与默认限流规则
- Timeout:定义执行command的超时时间,时间单位是ms,默认时间是1000ms;
- MaxConcurrnetRequests:定义command的最大并发量,默认值是10并发量;
- SleepWindow:熔断器被打开后使用,在熔断器被打开后,根据SleepWindow设置的时间控制多久后尝试服务是否可用,默认时间为5000ms;
- RequestVolumeThreshold:判断熔断开关的条件之一,统计10s(代码中写死了)内请求数量,达到这个请求数量后再根据错误率判断是否要开启熔断;
- ErrorPercentThreshold:判断熔断开关的条件之一,统计错误百分比,请求数量大于等于RequestVolumeThreshold并且错误率到达这个百分比后就会启动熔断 默认值是50;
- 然后执行hystrix.ConfigureCommand()方法,查看该方法源码,主要是创建熔断器,为熔断器命名并绑定限流规则,如果没有创建hystrix.CommandConfig,则使用默认的限流规则,查看该方法源码,会把限流规则封装成Settings对象作为value,以传进来的熔断器名称为key存储到circuitSettings一个map中
func ConfigureCommand(name string, config CommandConfig) {
//1.并发安全加锁处理
settingsMutex.Lock()
defer settingsMutex.Unlock()
//2.判断是否设置了指定的限流规则,如果没有使用默认的
timeout := DefaultTimeout
if config.Timeout != 0 {
timeout = config.Timeout
}
max := DefaultMaxConcurrent
if config.MaxConcurrentRequests != 0 {
max = config.MaxConcurrentRequests
}
volume := DefaultVolumeThreshold
if config.RequestVolumeThreshold != 0 {
volume = config.RequestVolumeThreshold
}
sleep := DefaultSleepWindow
if config.SleepWindow != 0 {
sleep = config.SleepWindow
}
errorPercent := DefaultErrorPercentThreshold
if config.ErrorPercentThreshold != 0 {
errorPercent = config.ErrorPercentThreshold
}
//3.将限流配置封装为Settings对象指针,以传进来的熔断器名称为key,
//Settings为value存储到circuitSettings一个map中
circuitSettings[name] = &Settings{
Timeout: time.Duration(timeout) * time.Millisecond,
MaxConcurrentRequests: max,
RequestVolumeThreshold: uint64(volume),
SleepWindow: time.Duration(sleep) * time.Millisecond,
ErrorPercentThreshold: errorPercent,
}
}
- circuitSettings 这个map是什么时候初始化的?, 该文件中提供了init(),通过该方法进行初始化
func init() {
circuitSettings = make(map[string]*Settings)
settingsMutex = &sync.RWMutex{}
log = DefaultLogger
}
执行
- 查看hystrix-go中对业务函数监控的方法, hystrix.Do内部调用的Doc方法,Go内部调用的是Goc方法,在Doc方法内部最终调用的还是Goc方法,只是在Doc方法内做了同步逻辑, Goc 是核心方法,查看该方法,方法内部的大致流程
- 针对每次请求会封装一个command
- 获取到熔断器CircuitBreaker
func GoC(ctx context.Context, name string, run runFuncC, fallback fallbackFuncC) chan error {
//1.可以简单理解为这次请求
cmd := &command{
run: run,
fallback: fallback,
start: time.Now(),
errChan: make(chan error, 1),
finished: make(chan bool, 1),
}
//2.根据熔断器名称获取真实熔断器CircuitBreaker
circuit, _, err := GetCircuit(name)
if err != nil {
cmd.errChan <- err
return cmd.errChan
}
cmd.circuit = circuit
ticketCond := sync.NewCond(cmd)
ticketChecked := false
//3.令牌返回逻辑
returnTicket := func() {
cmd.Lock()
for !ticketChecked {
ticketCond.Wait()
}
cmd.circuit.executorPool.Return(cmd.ticket)
cmd.Unlock()
}
//它的作用是确保由最快那个Goroutine运行errWithFallback()和reportAllEvent(),而且保证只会执行一次
returnOnce := &sync.Once{}
//4.上报此次请求处理结果逻辑
reportAllEvent := func() {
err := cmd.circuit.ReportEvent(cmd.events, cmd.start, cmd.runDuration)
if err != nil {
log.Printf(err.Error())
}
}
go func() {
//5.此次请求结束的标志
defer func() { cmd.finished <- true }()
//6.不允许请求,熔断器开启,直接走兜底处理逻辑
if !cmd.circuit.AllowRequest() {
cmd.Lock()
// It's safe for another goroutine to go ahead releasing a nil ticket.
ticketChecked = true
ticketCond.Signal()
cmd.Unlock()
returnOnce.Do(func() {
returnTicket()
cmd.errorWithFallback(ctx, ErrCircuitOpen)
reportAllEvent()
})
return
}
//可以请求,
cmd.Lock()
select {
//7.可以获取到令牌,走正常处理逻辑
case cmd.ticket = <-circuit.executorPool.Tickets:
ticketChecked = true
ticketCond.Signal()
cmd.Unlock()
default:
//8.获取不到令牌,走兜底逻辑
ticketChecked = true
ticketCond.Signal()
cmd.Unlock()
returnOnce.Do(func() {
returnTicket()
cmd.errorWithFallback(ctx, ErrMaxConcurrency)
reportAllEvent()
})
return
}
runStart := time.Now()
runErr := run(ctx)
returnOnce.Do(func() {
defer reportAllEvent()
cmd.runDuration = time.Since(runStart)
returnTicket()
if runErr != nil {
cmd.errorWithFallback(ctx, runErr)
return
}
cmd.reportEvent("success")
})
}()
//9.专门开的协程对超时情况进行兜底处理
go func() {
timer := time.NewTimer(getSettings(name).Timeout)
defer timer.Stop()
select {
case <-cmd.finished:
// returnOnce has been executed in another goroutine
case <-ctx.Done():
returnOnce.Do(func() {
returnTicket()
cmd.errorWithFallback(ctx, ctx.Err())
reportAllEvent()
})
return
case <-timer.C:
returnOnce.Do(func() {
returnTicket()
cmd.errorWithFallback(ctx, ErrTimeout)
reportAllEvent()
})
return
}
}()
return cmd.errChan
}
- 在看这段逻辑会有这样几个问题:
1.熔断器何时开启,如何开启?
2.熔断器何时关闭?
3.如何去上报这些事件,包括成功的失败的?
4.如何对这些上报事件进行处理
5.令牌的逻辑是什么?
1. 针对每次请求封装command,对command解释
- 以GoC()为例查看业务执行时内部的逻辑,首先会针对每次请求封装一个command,查看command内部关键的两个字段是events和circuit
- events主要是存储事件类型信息,比如执行成功的success,或者失败的timeout、context_canceled等
- circuit 是指针类型的CircuitBreaker,也就是真正的熔断器
type command struct {
sync.Mutex
ticket *struct{}
start time.Time
errChan chan error
finished chan bool
circuit *CircuitBreaker
run runFuncC
fallback fallbackFuncC
runDuration time.Duration
events []string
}
- command主要是记录单个执行的状态以及和熔断器进行一些运行交互,例如主要向CircuitBreaker上报执行状态事件
2. 真实熔断器CircuitBreaker的创建
- 先看一下CircuitBreaker熔断器结构中的字段解释
- open表示当前熔断器是否开启
- executorPool是流量控制中心,所有的请求都需要先获取到令牌
- metrics的类型是*metricExchange, 可以看成是上报执行状态事件的载体。通过它把执行状态信息存储到实际熔断器执行各个维度状态 (成功次数,失败次数,超时……) 的数据集合中
type CircuitBreaker struct {
Name string
open bool
forceOpen bool
mutex *sync.RWMutex
openedOrLastTestedTime int64
executorPool *executorPool
metrics *metricExchange
}
- 针对每次请求执行时会调用"circuit, _, err := GetCircuit(name)"根据熔断器名称在获取真实熔断器CircuitBreaker,
func GetCircuit(name string) (*CircuitBreaker, bool, error) {
circuitBreakersMutex.RLock()
//1.判断circuitBreakers中是否存在指定名称的熔断器
//如果存在则获取,不存在则创建
_, ok := circuitBreakers[name]
if !ok {
circuitBreakersMutex.RUnlock()
circuitBreakersMutex.Lock()
defer circuitBreakersMutex.Unlock()
if cb, ok := circuitBreakers[name]; ok {
return cb, false, nil
}
//创建熔断器
circuitBreakers[name] = newCircuitBreaker(name)
} else {
defer circuitBreakersMutex.RUnlock()
}
return circuitBreakers[name], !ok, nil
}
- 创建CircuitBreaker熔断器
func newCircuitBreaker(name string) *CircuitBreaker {
c := &CircuitBreaker{}
c.Name = name
//是metricExchange结构
c.metrics = newMetricExchange(name)
c.executorPool = newExecutorPool(name)
c.mutex = &sync.RWMutex{}
return c
}
- 在创建熔断器的同时也初始化了executorPool 和 metrics,metrics是metricExchange 类型
创建熔断器时封装的 metricExchange
- 在创建熔断器的同时也初始化了executorPool 和 metrics,metrics是metricExchange 类型
- Updates是一个channel类型,通过Updates上报执行事件集合。
- metricCollectors存储的是metricCollector.MetricCollector切片,而metricCollector.MetricCollector是一个接口类型,默认使用DefaultMetricCollector,内部存储了熔断器执行状态的所有信息
- 该方法内重点是1,3步骤,初始化Updates通道,获取到存储了熔断器执行状态所有信息的DefaultMetricCollector与执行"go m.Monitor()"通过协程去接收channel类型的Updates信息,即执行事件状态信息,最终把整合后的执行状态事件信息上报
func newMetricExchange(name string) *metricExchange {
m := &metricExchange{}
m.Name = name
//1.初始化Updates通道
m.Updates = make(chan *commandExecution, 2000)
m.Mutex = &sync.RWMutex{}
//2获取到存储了熔断器执行状态所有信息的DefaultMetricCollector
m.metricCollectors = metricCollector.Registry.InitializeMetricCollectors(name)
m.Reset()
//3.通过协程将执行状态事件信息上报
go m.Monitor()
return m
}
type metricExchange struct {
Name string
Updates chan *commandExecution
Mutex *sync.RWMutex
metricCollectors []metricCollector.MetricCollector
}
- 可以看到,初始化Updates通道的的容量是2000。初始化metricCollectors主要逻辑在InitializeMetricCollectors
func (m *metricCollectorRegistry) InitializeMetricCollectors(name string) []MetricCollector {
m.lock.RLock()
defer m.lock.RUnlock()
metrics := make([]MetricCollector, len(m.registry))
for i, metricCollectorInitializer := range m.registry {
metrics[i] = metricCollectorInitializer(name)
}
return metrics
}
1. DefaultMetricCollector
- 在执行InitializeMetricCollectors()函数时,发现创建了一个metricCollectorRegistry类型变量,在创建时会执行newDefaultMetricCollector()创建一个DefaultMetricCollector
var Registry = metricCollectorRegistry{
lock: &sync.RWMutex{},
registry: []func(name string) MetricCollector{
//查看该方法会返回DefaultMetricCollector
newDefaultMetricCollector,
},
}
func newDefaultMetricCollector(name string) MetricCollector {
m := &DefaultMetricCollector{}
m.mutex = &sync.RWMutex{}
m.Reset()
return m
}
- DefaultMetricCollector实现了MetricCollector所有方法,内部存储了熔断器执行状态所有信息
2. 执行状态事件信息上报
- 在newMetricExchange()函数中执行了"go m.Monitor()"开启协程去接收channel类型的Updates信息,即执行事件状态信息
func (m *metricExchange) Monitor() {
for update := range m.Updates {
// we only grab a read lock to make sure Reset() isn't changing the numbers.
m.Mutex.RLock()
totalDuration := time.Since(update.Start)
wg := &sync.WaitGroup{}
for _, collector := range m.metricCollectors {
wg.Add(1)
//在接收到事件信息后,
//调用IncrementMetrics先做状态信息的整合
//最终把整合后的执行状态事件信息上报collector.Update(r)
go m.IncrementMetrics(wg, collector, update, totalDuration)
}
wg.Wait()
m.Mutex.RUnlock()
}
}
- 在接收到事件信息后,调用IncrementMetrics先做状态信息的整合,最终把整合后的执行状态事件信息上报collector.Update®
func (m *metricExchange) IncrementMetrics(wg *sync.WaitGroup, collector metricCollector.MetricCollector, update *commandExecution, totalDuration time.Duration) {
// granular metrics
r := metricCollector.MetricResult{
Attempts: 1,
TotalDuration: totalDuration,
RunDuration: update.RunDuration,
ConcurrencyInUse: update.ConcurrencyInUse,
}
switch update.Types[0] {
case "success":
r.Successes = 1
case "failure":
r.Failures = 1
r.Errors = 1
case "rejected":
r.Rejects = 1
r.Errors = 1
case "short-circuit":
r.ShortCircuits = 1
r.Errors = 1
case "timeout":
r.Timeouts = 1
r.Errors = 1
case "context_canceled":
r.ContextCanceled = 1
case "context_deadline_exceeded":
r.ContextDeadlineExceeded = 1
}
if len(update.Types) > 1 {
// fallback metrics
if update.Types[1] == "fallback-success" {
r.FallbackSuccesses = 1
}
if update.Types[1] == "fallback-failure" {
r.FallbackFailures = 1
}
}
collector.Update(r)
wg.Done()
}
- 在接收到事件信息后,调用IncrementMetrics先做状态信息的整合,最终把整合后的执行状态事件信息上报collector.Update®
func (d *DefaultMetricCollector) Update(r MetricResult) {
d.mutex.RLock()
defer d.mutex.RUnlock()
d.numRequests.Increment(r.Attempts)
d.errors.Increment(r.Errors)
d.successes.Increment(r.Successes)
d.failures.Increment(r.Failures)
d.rejects.Increment(r.Rejects)
d.shortCircuits.Increment(r.ShortCircuits)
d.timeouts.Increment(r.Timeouts)
d.fallbackSuccesses.Increment(r.FallbackSuccesses)
d.fallbackFailures.Increment(r.FallbackFailures)
d.contextCanceled.Increment(r.ContextCanceled)
d.contextDeadlineExceeded.Increment(r.ContextDeadlineExceeded)
d.totalDuration.Add(r.TotalDuration)
d.runDuration.Add(r.RunDuration)
}
创建熔断器时封装的executorPool
- 创建熔断器的newCircuitBreaker(name string)方法中会调用newExecutorPool()获取executorPool ,
- 在executorPool 中关注两个字段
- Tickets表示的就是访问令牌带缓冲通道的channel,初始化channel容量取决于一开始设置的MaxConcurrentRequests当有请求到来时,从channel中拿出一个令牌,调用后重新归还
- poolMetrics就是流量控制的具体指标
- Executed 表示当前桶已经处理的请求数量
type executorPool struct {
Name string
Metrics *poolMetrics
Max int
Tickets chan *struct{}
}
//1创建熔断器
func newExecutorPool(name string) *executorPool {
p := &executorPool{}
p.Name = name
//获取PoolMetrics
p.Metrics = newPoolMetrics(name)
p.Max = getSettings(name).MaxConcurrentRequests
p.Tickets = make(chan *struct{}, p.Max)
for i := 0; i < p.Max; i++ {
p.Tickets <- &struct{}{}
}
return p
}
- 此外在 newExecutorPool(name) 函数中,执行newPoolMetrics()获取poolMetrics时,也启动一个 go m.Monitor() 专门去更新当前桶的最大值
func newPoolMetrics(name string) *poolMetrics {
m := &poolMetrics{}
m.Name = name
m.Updates = make(chan poolMetricsUpdate)
m.Mutex = &sync.RWMutex{}
m.Reset()
//开启协程,更新当前桶的最大值
go m.Monitor()
return m
}
func (m *poolMetrics) Monitor() {
for u := range m.Updates {
m.Mutex.RLock()
m.Executed.Increment(1)
//更新当前桶的最大值
m.MaxActiveRequests.UpdateMax(float64(u.activeCount))
m.Mutex.RUnlock()
}
}
func (r *Number) UpdateMax(n float64) {
r.Mutex.Lock()
defer r.Mutex.Unlock()
b := r.getCurrentBucket()
if n > b.Value {
b.Value = n
}
r.removeOldBuckets()
}
- 自此获取真实熔断器的内部流程执行完毕
3. GOC()中的上报执行事件
- 在GoC()方法中会封装一个名为"reportAllEvent "函数,该函数就是用来上报执行事件的
//上方GoC()方法中第4步骤
reportAllEvent := func() {
err := cmd.circuit.ReportEvent(cmd.events, cmd.start, cmd.runDuration)
if err != nil {
log.Printf(err.Error())
}
}
- 在
func (circuit *CircuitBreaker) ReportEvent(eventTypes []string, start time.Time, runDuration time.Duration) error {
if len(eventTypes) == 0 {
return fmt.Errorf("no event types sent for metrics")
}
circuit.mutex.RLock()
o := circuit.open
circuit.mutex.RUnlock()
//1.当前熔断器是开启的,并且已经过了SleepWindow时间,
//此时请求就属于半开的状态,允许尝试执行,如果执行成功,
//那么就说明服务恢复了,可以关闭熔断器了
if eventTypes[0] == "success" && o {
circuit.setClose()
}
var concurrencyInUse float64
if circuit.executorPool.Max > 0 {
concurrencyInUse = float64(circuit.executorPool.ActiveCount()) / float64(circuit.executorPool.Max)
}
//2.组装执行状态状态事件,然后塞进Updates通道中
//正好被初始化metricExchange另开的Goroutine接收
//这样这个上报流程就对应上了
select {
case circuit.metrics.Updates <- &commandExecution{
Types: eventTypes,
Start: start,
RunDuration: runDuration,
ConcurrencyInUse: concurrencyInUse,
}:
default:
return CircuitError{Message: fmt.Sprintf("metrics channel (%v) is at capacity", circuit.Name)}
}
return nil
}
4. GoC()中第一个协程
- GoC()方法中上报事件函数封装完毕后,接下来有一个函数
go func() {
//1.作为正常运行结束的通知
defer func() { cmd.finished <- true }()
//2.判断是否运行请求,不允许请求,熔断器开启,直接走兜底处理逻辑
//circuit.AllowRequest()就是判断是否能请求的核心逻辑
if !cmd.circuit.AllowRequest() {
cmd.Lock()
// It's safe for another goroutine to go ahead releasing a nil ticket.
ticketChecked = true
ticketCond.Signal()
cmd.Unlock()
returnOnce.Do(func() {
returnTicket()
cmd.errorWithFallback(ctx, ErrCircuitOpen)
reportAllEvent()
})
return
}
//3.可以请求
cmd.Lock()
select {
//4.可以获取到令牌,走正常处理逻辑
case cmd.ticket = <-circuit.executorPool.Tickets:
ticketChecked = true
ticketCond.Signal()
cmd.Unlock()
default:
//8.获取不到令牌,走兜底逻辑
ticketChecked = true
ticketCond.Signal()
cmd.Unlock()
returnOnce.Do(func() {
returnTicket()
cmd.errorWithFallback(ctx, ErrMaxConcurrency)
reportAllEvent()
})
return
}
runStart := time.Now()
runErr := run(ctx)
returnOnce.Do(func() {
defer reportAllEvent()
cmd.runDuration = time.Since(runStart)
returnTicket()
if runErr != nil {
cmd.errorWithFallback(ctx, runErr)
return
}
cmd.reportEvent("success")
})
}()
如何判断是否能请求
- circuit.AllowRequest(),首先判断熔断器是否被强制开启或者已经开启,如果是,直接返回true。否则说明是当前熔断器处于关闭。接着判断过去十秒内各个桶值的和是否小于设置的RequestVolumeThreshold值,如果小于,说明熔断器还应该是关闭状态,返回false。如果大于等于,那么应该进一步去判断错误百分比 是否超出自己的设置的ErrorPercentThreshold。如果超出了,那么说明错误率过高,此时需要开启熔断器
- 如果 circuit.IsOpen() 不符合,那么再看 circuit.allowSingleTest()。虽然熔断器是开启的,但是如果当前的时间已经大于 (上次开启熔断器的时间 +SleepWindow 的时间),这时候熔断器属于半开的状态,可以执行下一步。那么就会返回 true。
- 如果 cmd.circuit.AllowRequest 返回 false,那么就是执行 returnTicket 归还令牌 (尽管这时候还没有令牌可言)。这段代码很有趣,通过变量 ticketChecked 加 sync.NewCond 实现的逻辑。cmd.errorWithFallback(),上报熔断器已开启事件 (circuit open) 以及运行 fallBakck 保底函数 (如果存在的话),执行结束,响应
//1.第一种熔断器是关闭的对应circuit.IsOpen()
func (circuit *CircuitBreaker) AllowRequest() bool {
return !circuit.IsOpen() || circuit.allowSingleTest()
}
获取令牌
- 在上面执行执行circuit.IsOpen() 与 circuit.allowSingleTest()后拿到可以向下执行的结果后,执行第协程中的第3步骤开始向下执行获取令牌
- 如果拿不到访问令牌,那么和刚才一样,上报当前请求已超过并发数事件 (max concurrency),运行保底操作,响应
- 如果拿到访问令牌,那么真正执行自己的业务代码 run(ctx)。套路和上面相似
5. GoC()中第二个协程
- 查看GoC()中的第二个协程,三个case分别表示:正常执行结束、业务执行被取消以及超时
go func() {
timer := time.NewTimer(getSettings(name).Timeout)
defer timer.Stop()
select {
case <-cmd.finished:
// returnOnce has been executed in another goroutine
case <-ctx.Done():
returnOnce.Do(func() {
returnTicket()
cmd.errorWithFallback(ctx, ctx.Err())
reportAllEvent()
})
return
case <-timer.C:
returnOnce.Do(func() {
returnTicket()
cmd.errorWithFallback(ctx, ErrTimeout)
reportAllEvent()
})
return
}
}()
DoC()的特殊处理
- 至于为什么说Do是同步操作,因为在Doc最后,在GoC执行后有通过select的阻塞等待处
func DoC(ctx context.Context, name string, run runFuncC, fallback fallbackFuncC) error {
done := make(chan struct{}, 1)
r := func(ctx context.Context) error {
err := run(ctx)
if err != nil {
return err
}
done <- struct{}{}
return nil
}
f := func(ctx context.Context, e error) error {
err := fallback(ctx, e)
if err != nil {
return err
}
done <- struct{}{}
return nil
}
var errChan chan error
if fallback == nil {
errChan = GoC(ctx, name, r, nil)
} else {
errChan = GoC(ctx, name, r, f)
}
//在GoC执行后有通过select的阻塞等待处
select {
case <-done:
return nil
case err := <-errChan:
return err
}
}
三. go 中其它流控工具
更多推荐
所有评论(0)