一. hystrix-go

  1. hystrix-go 是基于Java的Hystrix类似的执行语义实现的一个集成了流量统计, 限流、熔断、降级等功能的三方类库
  2. 使用时需要执行"go get github.com/afex/hystrix-go/hystrix" 命令将hystrix-go拉取到本地
  3. hystrix的设计原则
  1. 防止任何单个依赖服务耗尽所有用户线程
  2. 直接响应失败,而不是一直等待
  3. 提供错误返回接口,而不是让用户线程直接处理依赖服务抛出的异常
  4. 使用隔离或熔断技术来降低并限制单个依赖对整个系统造成的影响
  1. 熔断器三种状态
  1. 关闭状态: 服务正常,维护失败率统计,达到失败率阈值时,转到开启状态。
  2. 开启状态: 服务异常,调用fallback函数,一段时间后,进入半开启状态。
  3. 半开启状态: 尝试恢复服务,失败率高于阈值,进入开启状态,低于阈值,进入关闭状态。
  1. hystrix-go 基本使用示例
  1. 创建hystrix流量统计服务器
  2. 设置流控规则
  3. 创建熔断器
  4. 熔断器与业务接口绑定,监控指定业务接口流量
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)
		}
	}
}

  1. Do内部调用的Doc方法,Go内部调用的是Goc方法,在Doc方法内部最终调用的还是Goc方法,只是在Doc方法内做了同步逻辑
  2. 实际可以理解为启动了一个专门的hystrix服务用来管理熔断器,通过熔断器监控指定的业务接口的流量情况,当启动此处的hystrix流控服务后可以访问这个地址,例如此处访问:127.0.0.1:8074 就可以拿到监控的访问情况

二. hystrix-go 原理相关

  1. 参考转载博客1 ,参考转载博客2
  2. hystrix-go中Do和Go方法内部都是调用了hystrix.GoC方法,只是Do方法处理了异步的过程
  3. 可以把hystrix-go分为: 初始化,执行,执行过程中的状态上报与开关判断几个位置去理解

初始化

  1. 在我们使用hystrix-go时,需要先创建hystrix.CommandConfig封装限流规则,如不不创建则使用默认的限流规则,查看hystrix.CommandConfig结构体属性与默认限流规则
  1. Timeout:定义执行command的超时时间,时间单位是ms,默认时间是1000ms;
  2. MaxConcurrnetRequests:定义command的最大并发量,默认值是10并发量;
  3. SleepWindow:熔断器被打开后使用,在熔断器被打开后,根据SleepWindow设置的时间控制多久后尝试服务是否可用,默认时间为5000ms;
  4. RequestVolumeThreshold:判断熔断开关的条件之一,统计10s(代码中写死了)内请求数量,达到这个请求数量后再根据错误率判断是否要开启熔断;
  5. ErrorPercentThreshold:判断熔断开关的条件之一,统计错误百分比,请求数量大于等于RequestVolumeThreshold并且错误率到达这个百分比后就会启动熔断 默认值是50;
  1. 然后执行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,
	}
}
  1. circuitSettings 这个map是什么时候初始化的?, 该文件中提供了init(),通过该方法进行初始化
func init() {
	circuitSettings = make(map[string]*Settings)
	settingsMutex = &sync.RWMutex{}
	log = DefaultLogger
}

执行

  1. 查看hystrix-go中对业务函数监控的方法, hystrix.Do内部调用的Doc方法,Go内部调用的是Goc方法,在Doc方法内部最终调用的还是Goc方法,只是在Doc方法内做了同步逻辑, Goc 是核心方法,查看该方法,方法内部的大致流程
  1. 针对每次请求会封装一个command
  2. 获取到熔断器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. 在看这段逻辑会有这样几个问题:

1.熔断器何时开启,如何开启?
2.熔断器何时关闭?
3.如何去上报这些事件,包括成功的失败的?
4.如何对这些上报事件进行处理
5.令牌的逻辑是什么?

1. 针对每次请求封装command,对command解释

  1. 以GoC()为例查看业务执行时内部的逻辑,首先会针对每次请求封装一个command,查看command内部关键的两个字段是events和circuit
  1. events主要是存储事件类型信息,比如执行成功的success,或者失败的timeout、context_canceled等
  2. 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
}
  1. command主要是记录单个执行的状态以及和熔断器进行一些运行交互,例如主要向CircuitBreaker上报执行状态事件

2. 真实熔断器CircuitBreaker的创建

  1. 先看一下CircuitBreaker熔断器结构中的字段解释
  1. open表示当前熔断器是否开启
  2. executorPool是流量控制中心,所有的请求都需要先获取到令牌
  3. metrics的类型是*metricExchange, 可以看成是上报执行状态事件的载体。通过它把执行状态信息存储到实际熔断器执行各个维度状态 (成功次数,失败次数,超时……) 的数据集合中
type CircuitBreaker struct {
	Name                   string
	open                   bool
	forceOpen              bool
	mutex                  *sync.RWMutex
	openedOrLastTestedTime int64

	executorPool *executorPool
	metrics      *metricExchange
}
  1. 针对每次请求执行时会调用"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
}
  1. 创建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
}
  1. 在创建熔断器的同时也初始化了executorPool 和 metrics,metrics是metricExchange 类型
创建熔断器时封装的 metricExchange
  1. 在创建熔断器的同时也初始化了executorPool 和 metrics,metrics是metricExchange 类型
  1. Updates是一个channel类型,通过Updates上报执行事件集合。
  2. metricCollectors存储的是metricCollector.MetricCollector切片,而metricCollector.MetricCollector是一个接口类型,默认使用DefaultMetricCollector,内部存储了熔断器执行状态的所有信息
  1. 该方法内重点是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
}
  1. 可以看到,初始化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
  1. 在执行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
}
  1. DefaultMetricCollector实现了MetricCollector所有方法,内部存储了熔断器执行状态所有信息
2. 执行状态事件信息上报
  1. 在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()
	}
}
  1. 在接收到事件信息后,调用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()
}
  1. 在接收到事件信息后,调用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
  1. 创建熔断器的newCircuitBreaker(name string)方法中会调用newExecutorPool()获取executorPool ,
  2. 在executorPool 中关注两个字段
  1. Tickets表示的就是访问令牌带缓冲通道的channel,初始化channel容量取决于一开始设置的MaxConcurrentRequests当有请求到来时,从channel中拿出一个令牌,调用后重新归还
  2. poolMetrics就是流量控制的具体指标
  3. 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
}
  1. 此外在 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()
}
  1. 自此获取真实熔断器的内部流程执行完毕

3. GOC()中的上报执行事件

  1. 在GoC()方法中会封装一个名为"reportAllEvent "函数,该函数就是用来上报执行事件的
	//上方GoC()方法中第4步骤  
    reportAllEvent := func() {
        err := cmd.circuit.ReportEvent(cmd.events, cmd.start, cmd.runDuration)
        if err != nil {
            log.Printf(err.Error())
        }
    }
  1. 在
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()中第一个协程

  1. 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")
        })
    }()
如何判断是否能请求
  1. circuit.AllowRequest(),首先判断熔断器是否被强制开启或者已经开启,如果是,直接返回true。否则说明是当前熔断器处于关闭。接着判断过去十秒内各个桶值的和是否小于设置的RequestVolumeThreshold值,如果小于,说明熔断器还应该是关闭状态,返回false。如果大于等于,那么应该进一步去判断错误百分比 是否超出自己的设置的ErrorPercentThreshold。如果超出了,那么说明错误率过高,此时需要开启熔断器
  1. 如果 circuit.IsOpen() 不符合,那么再看 circuit.allowSingleTest()。虽然熔断器是开启的,但是如果当前的时间已经大于 (上次开启熔断器的时间 +SleepWindow 的时间),这时候熔断器属于半开的状态,可以执行下一步。那么就会返回 true。
  2. 如果 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()
}
获取令牌
  1. 在上面执行执行circuit.IsOpen() 与 circuit.allowSingleTest()后拿到可以向下执行的结果后,执行第协程中的第3步骤开始向下执行获取令牌
  2. 如果拿不到访问令牌,那么和刚才一样,上报当前请求已超过并发数事件 (max concurrency),运行保底操作,响应
  3. 如果拿到访问令牌,那么真正执行自己的业务代码 run(ctx)。套路和上面相似

5. GoC()中第二个协程

  1. 查看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()的特殊处理
  1. 至于为什么说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 中其它流控工具

参考博客

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐