一、前言
一些关于并发操作的组件。
二、学习代码
速率控制:
package main import ( "fmt" "time" ) func main() { request := make(chan int, 5) //模拟到来五个请求 for i := 0; i < 5; i++ { request <- i * 5 } close(request) //到来五个请求后停止接收请求 limiter := time.Tick(time.Second / 5) // 200ms每次接收 //Tick底层实现就是NewTicker,返回一个channel,每隔200ms向channel中放入一个时间,相当于无缓冲 for req := range request { <-limiter //每次先等待limiter的时间间隔 fmt.Println("request:", req, time.Now()) } //突发 brustLimiter := make(chan time.Time, 3) //模拟突发请求,实则为有缓冲的通道,突发计时器 for i := 0; i < 3; i++ { //注入三个时间到brust中 brustLimiter <- time.Now() } go func() { //理解为计时器协程,每200ms往brustLimiter中放入一个时间,保证每200ms可以处理一个请求 for t := range time.Tick(time.Second / 5) { //fmt.Println("brust time:", t, time.Now()) brustLimiter <- t //每200ms向brust中放入一个时间 } }() time.Sleep(50 * time.Millisecond) //主线程睡眠50ms,保证子线程有足够的时间执行 BrustRequest := make(chan int, 5) //模拟到来五个请求 for i := 0; i < 5; i++ { BrustRequest <- i * 3 } close(BrustRequest) //到来五个请求后停止接收请求 for req := range BrustRequest { <-brustLimiter //每次先等待brustLimiter的时间间隔 fmt.Println("brust request:", req, time.Now()) } }原子计数器:
package main import ( "fmt" "sync" "sync/atomic" ) func main() { var count uint64 = 0 var wg sync.WaitGroup for i := 0; i < 10; i++ { wg.Add(1) atomic.AddUint64(&count, 1) //原子操作加1 if i == 5 { fmt.Println(atomic.LoadUint64(&count)) //原子操作读取count的值 } wg.Done() } wg.Wait() fmt.Println("Final count:", count) atomic.StoreUint64(&count, 520) //原子操作存储count的值,即赋值 }互斥锁:
package main import ( "fmt" "sync" ) type Container struct { mu sync.Mutex //互斥锁,在一个时间段里只能有一个线程访问 counter map[string]int } func (c *Container) Increment(key string) { c.mu.Lock() defer c.mu.Unlock() c.counter[key]++ } func (c *Container) Get(key string) int { c.mu.Lock() defer c.mu.Unlock() return c.counter[key] } func main() { c := Container{ //互斥变量默认为可用 counter: map[string]int{"key1": 0, "key2": 0}, } var wg sync.WaitGroup doInc := func(key string, times int) { defer wg.Done() for i := 0; i < times; i++ { c.Increment(key) } } wg.Add(2) go doInc("key1", 1000) go doInc("key2", 1000) wg.Wait() fmt.Println("Final counts:", c.Get("key1"), c.Get("key2")) }状态协程:
package main import ( "fmt" "sync" ) // 状态协程,对信息的操作全部封装在一个协程里,有需求就向状态协程里发送消息,状态协程收到消息后进行处理,处理完后再返回结果给调用方 const ( Add = iota //默认为0 Sub //1 Get //2 ) type Msg struct { Op int //即前面const定义的操作类型 Value int Resp chan string //用于返回结果的通道 } func Operation(balance *int, msg <-chan Msg, done <-chan struct{}, wg *sync.WaitGroup) { defer fmt.Println("Operation goroutine exit") defer wg.Done() for { select { case <-done: return case cmd := <-msg: switch cmd.Op { case Add: *balance += cmd.Value cmd.Resp <- "Add Success" case Sub: *balance -= cmd.Value cmd.Resp <- "Sub Success" case Get: cmd.Resp <- fmt.Sprintf("Balance is %d", *balance) } } } } func main() { balance := 0 msg := make(chan Msg) done := make(chan struct{}) var wg sync.WaitGroup wg.Add(1) go Operation(&balance, msg, done, &wg) //发送加钱的消息 resp := make(chan string) msg <- Msg{Op: Add, Value: 100, Resp: resp} fmt.Println(<-resp) //发送减钱的消息 msg <- Msg{Op: Sub, Value: 50, Resp: resp} fmt.Println(<-resp) msg <- Msg{Op: Get, Resp: resp} fmt.Println(<-resp) // workGroup var workGroup sync.WaitGroup for i := 0; i < 10; i++ { workGroup.Add(1) go func(i int) { defer workGroup.Done() myResp := make(chan string) msg <- Msg{Op: Get, Resp: myResp} fmt.Println(<-myResp, i) }(i) } workGroup.Wait() //等待所有的需求操作的协程执行完毕 close(done) //这个Done是相对优雅的退出信号 wg.Wait() }