、状态协程)
一、前言一些关于并发操作的组件。二、学习代码速率控制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() }