文章目录
- 1. 惰性生成器
- 1.1 基于通道创建
- 1.2 Iter.Seq 实现
- 1.3 自己实现一个 Seql
- 1.4 生成更通用的惰性生成器
- 2. Future 模式
- 2.1 基础 Future 模式实现
- 2.2 基础使用
- 2.2.1 http 异步调用
- 2.2.2 带超时控制的 Future
- 2.3 多任务并发执行
- 3. 小结
本系列文章:
- Java 转 go 学习 - 项目管理
- Java 转 go 学习 - 基本语法
- Java 转 go 学习 - 类型转换
- Java 转 go 学习 - 流程控制结构
- Java 转 go 学习 - 数组和切片
- Java 转 go 学习 - map
- Java 转 go 学习 - 函数(1)
- Java 转 go 学习 - 函数(2)
- Java 转 go 学习 - 结构体
- Java 转 go 学习 - 接口
- Java 转 go 学习 - 并发编程(1)
- Java 转 go 学习 - 并发编程(2)
1. 惰性生成器
惰性生成器就是按需计算并返回值的机制,通过channel和 goroutine 结合实现,仅仅在需要的时候才生成下一个值,而不是一次性返回整个序列,通过这种方式可以节省内存,提高处理大数据流的效率。
具体的表现就是:
- 创建一个通道,调用者通过这个通道来获取生成的值。
- 后台通过 goroutine 来按需生成值然后发送到通道。
- 每次从通道读取的时候才会生成下一个值。
1.1 基于通道创建
首先来看下最基础的创建方式,基于通道和 goroutine 创建。
// 返回一个只读通道funcintegers()<-chanint{yield:=make(chanint)count:=0gofunc(){for{yield<-count count++}}()returnyield}functest61(){gen:=integers()fmt.Println(<-gen)// 0fmt.Println(<-gen)// 1fmt.Println(<-gen)// 2}这种生成器的核心逻辑就是:通道阻塞,我们知道 channel 的特点就是如果没有缓冲区,那么往里面写入数据必须要等到从通道获取数据之后才能继续往里面写入。
虽然integers方法中通过 go 启动了一个协程不断往通道写入数据,但是yield <- count写入第一个数据之后会直接阻塞,所以integers方法将这个只读通道返回了,注意返回值是<- chan,也就是说虽然我们创建出来的 channel 是双向的,但是函数会转成只读通道返回。
拿到这个通道之后,只有调用<- gen才会从通道获取数据,然后协程才会继续往里面写入,同时由于协程阻塞会休眠,不会消耗 CPU,所以不会导致 CPU 空转。
1.2 Iter.Seq 实现
Go 1.23+ 之后官方提供了iter包,可以用来实现更简洁的惰性生成器。
// integersSeq 返回一个基于 iter.Seq 的惰性生成器,生成 1 到 5。funcintegersSeq()iter.Seq[int]{returnfunc(yieldfunc(int)bool){fori:=1;;i++{if!yield(i){return}}}}functest62(){// 直接遍历值//for v := range integersSeq() {// fmt.Println(v)// // 1// // 2// // 3// // 4// // ...//}// 也可以手动获取下一个值的函数next,stop:=iter.Pull(integersSeq())fmt.Println(next())// 1 truefmt.Println(next())// 2 truefmt.Println(next())// 3 truefmt.Println(next())// 4 truefmt.Println(next())// 5 truestop()fmt.Println(next())// 0 false}下面来看下这个函数:
typeSeq[V any]func(yieldfunc(V)bool)bool这个函数接收一个 func(V) bool 类型的参数,返回 bool 的函数,当我们在函数内调用yield,就会把值传递给消费者,比如for-range遍历的时候就会不断接收从yield方法传过来的值,next()也是。
什么时候yield会返回 false,比如调用了stop()函数,for-range遍历完成,或者在for-range中进行break了,上面这种写法不能用for-range遍历,因为生成器从 1 开始不断往里面通过yield传数据,所以我们可以通过Pull来拉取数据,每一次就 + 1,返回的第二个参数代表生成器有没有停止。
1.3 自己实现一个 Seql
我们也可以来模拟一个yield的写法。
首先定义迭代器类型:
// 定义迭代器类型typeMySeq[V any]func(yieldfunc(V)bool)然后我们定义MyPull方法,这个方法参数是MySeq类型,返回值是(next func() (V, bool), stop func()),通过 next 可以获取下一个值,通过 stop 可以关闭生成器的通道。
// MyPull 把 push 风格的迭代器转成 pull 风格的 next/stop。// stop 返回后,后续 next 一定只会得到零值和 false。funcMyPull[V any](seq MySeq[V])(nextfunc()(V,bool),stopfunc()){ch:=make(chanV)stopCh:=make(chanstruct{})done:=make(chanstruct{})gofunc(){deferclose(done)deferclose(ch)seq(func(v V)bool{select{case<-stopCh:returnfalsecasech<-v:returntrue}})}()next=func()(V,bool){v,ok:=<-chreturnv,ok}stop=func(){select{case<-done:returncase<-stopCh:<-donereturndefault:close(stopCh)<-done}}returnnext,stop}可以看到这个函数中ch就是生成数的通道,stopCh和done是暂停生成的通道,在函数里面启动一个 go 协程,执行seq函数,由于这个函数接收参数是yield func(V) bool,所以往通道写入的函数就是yield,做的事情很简单,就是把yield的参数往ch通道写入,同时监听stopCh通道,如果发现stopCh关闭了就会监听到这个分支然后返回 false。
为什么要两个通道stopCh和done呢,这两个通道代表的含义不同:
stopCh: 你对生成器说“别再产出了,准备停”done:生成器说“我已经退出了,后续不会再生成数据”
假设没有 done,我们可能会写成下面的写法。
stop=func(){close(stopCh)}这种情况下有可能会发生下面的情况(我们考虑多个程序同时调用 next()),假设先前调用了 next,此时 goroutine 是阻塞在 select 等待进入某个 case 的:
- 程序1 调用 stop(),此时
<-stopCh分支就绪 - 程序2 同时调用 next() 获取数据
- 由于此时
close(stopCh)和ch <- v都准备继续,那么 goroutine 的两个 case 都有可能进入 - 这时候外层看来就是,明明先调用了
stop,但是还能通过next获取数据
但是加上 done 之后,流程就是如下了:
- 程序1 调用 stop(),进入
default分支,close(stopCh) 发送停止信号,然后阻塞在<-done上面。 - 程序2 同时调用 next() 获取数据。
- 由于此时
close(stopCh)和ch <- v都准备继续,那么 goroutine 的两个 case 都有可能进入,此时没有语义上面的问题,因为这种情况下stop()卡在<-done上面,说明 stop 还没有调用完成,通道不关闭也是对的。 - goroutine 假设先处理完 next,再进入
case <-stopCh分支,此时 return false 退出seq。 - 假设此时还有 next() 请求到来,就会卡在
v, ok := <-ch,因为 goroutine 已经不生产数据了。 - goroutine 执行
close(ch)、close(done),执行完成之后其他 next() 请求会返回0, false,因为通道关闭了。 - 到这里程序1
stop()才收到信号,结束。
总之记住一句话:调用close之后从无缓冲通道获取到的就是通道类型的零值了,但是如果是有缓冲的通道就会先获取出缓冲区的值,最后再返回零值。
输出如下:
functest63(){next,stop:=MyPull(generator())fmt.Println(next())fmt.Println(next())fmt.Println(next())fmt.Println(next())stop()fmt.Println(next())fmt.Println(next())// 1 true// 2 true// 3 true// 4 true// 0 false// 0 false}上面解释了一些问题,下面来看下这种写法又有什么问题。
funcMyPull[V any](seq MySeq[V])(nextfunc()(V,bool),stopfunc()){ch:=make(chanV)stopCh:=make(chanstruct{})gofunc(){deferclose(ch)seq(func(v V)bool{select{casech<-v:returntruedefault:returnfalse}})}()next=func()(V,bool){v,ok:=<-chreturnv,ok}stop=func(){stopCh<-struct{}{}}returnnext,stop}这种写法是非阻塞的,也就是seq函数中,如果 ch 中暂时没有值默认走的是 default 分支,我们想的是用 default 代替<- stopCh,但是这是不对的,比如下面的语序:
- 主程序执行
next(),获取到 goroutine 设置到通道里面的值。 - 此时还没等主程序继续执行
next(),goroutine 继续遍历,走到 default 分支返回 false,退出,然后关闭 ch。 - 主程序此时执行了几次
next()返回的都是0, false。 - 主程序最后执行
stop()报错fatal error: all goroutines are asleep - deadlock!,因为这时候往stopCh写入的数据,阻塞住了,但是所有 goroutine 都退出了,没有协程再去处理这个通道。
输出如下:
1true0false0false0falsefatalerror:all goroutines are asleep-deadlock!如果 stop 改成close(stopCh),就会输出这个:
1true0false0false0false0false0false不管怎么样都不对。
1.4 生成更通用的惰性生成器
我们可以通过工厂函数创建惰性生成器。
typeAnyinterface{}typeEvalFuncfunc(Any)(Any,Any)funcBuildLazyIntEvaluator(evalFunc EvalFunc,initState Any)func()Any{// 数据写入通道retValChan:=make(chanAny)loopFunc:=func(){// 初始值varactState Any=initState// retVal 是经过我们的函数 evalFunc 计算出来的值varretVal Any// 循环写入通道for{// 将初始值传进去, 可以进行 +-*/ 等操作retVal,actState=evalFunc(actState)retValChan<-retVal}}// 函数, 从通道获取数据retFunc:=func()Any{return<-retValChan}// 启动协程, 往 retValChan 写入数据goloopFunc()returnretFunc}函数BuildLazyIntEvaluator接收evalFunc类型的参数,evalFunc是一个函数,可以将任意类型传进去,然后计算出下一个要返回的值retVal,最后将这个值写入通道,actState则是实时更新。
下面来看下用法:
functest65(){evenFunc:=func(state Any)(Any,Any){os:=state.(int)ns:=os+2returnos,ns}even:=BuildLazyIntEvaluator(evenFunc,0)fori:=0;i<10;i++{fmt.Printf("%vth even: %v\n",i,even())}}我们初始化函数evenFunc,这个函数会将state + 2返回给 go 协程,让协程写入通道,同时记录state + 2的值,下一次再根据这个值再传入evalFunc来获取下一个生成值。
2. Future 模式
Future 模式就是:用户向系统提交一个任务,系统会通过协程去执行这个任务,然后将结果写到 Future,用户执行完其他任务之后可以通过 Future 获取到之前提交任务的执行结果,简单来说就是将一些操作解耦出来异步执行。
Go 语言中实现 Future 是使用 goroutine 配合带缓冲的 channel,下面来看下具体实现。
2.1 基础 Future 模式实现
最常见的就是创建一个返回<- chan Result的函数,通过goroutine执行任务,然后通过带缓冲的 channel 返回结果。
packagemainimport("fmt""go-learn-1/learn-4/util""time")typeResultstruct{Data any Errerror}funcasyncCompute()<-chanResult{// 创建带缓冲的 channelch:=make(chanResult,1)// 开启协程调用函数gofunc(){deferclose(ch)data,err:=work()ch<-Result{Data:data,Err:err}}()returnch}funcwork()(any,error){fmt.Printf("[%s]work............\n",util.Now())time.Sleep(1*time.Second)return"模拟异常返回结果",fmt.Errorf("模拟抛出异常")}funcmain(){ch:=asyncCompute()// 阻塞res:=<-chifres.Err!=nil{fmt.Printf("[%s]%v:%v\n",util.Now(),res.Data,res.Err)}}这里的关键点就是必须使用带缓冲的channel(make(chan Result, 1)),否则 goroutine 可能在发送的时候阻塞。然后就是结果用一个结构体将 error 和返回结果包装起来。
最后接收方可以通过val, ok := <-ch判断结果是否已经送达,不需要外部调用手动 close(channel),在 goroutine 内部使用defer close(ch)来关闭,因为 ch 只是用来存结果,也就是只会往里面写入一次,所以直接关闭之后结果会写入缓冲区。
2.2 基础使用
2.2.1 http 异步调用
首先就是 http 异步获取结果写入 channel,下面是简单写法,可以看到我们通过 go 协程去 Get 数据,然后处理数据到 body,将结果写入channel。
funcfetchURL(urlstring)<-chanstring{ch:=make(chanstring,1)gofunc(){resp,err:=http.Get(url)iferr!=nil{ch<-""return}deferresp.Body.Close()body,_:=ioutil.ReadAll(resp.Body)ch<-string(body)}()returnch}// 使用html:=<-fetchURL("https://example.com")2.2.2 带超时控制的 Future
我们可以在 Future 中加上超时控制,如果处理结果超时,就往里面写入一个 Error。
funcasyncCompute(timeout time.Duration)<-chanResult{ch:=make(chanResult,1)gofunc(){resultCh:=make(chanResult,1)gofunc(){data,err:=work()resultCh<-Result{Data:data,Err:err}}()timer:=time.NewTimer(timeout)defertimer.Stop()select{caseres:=<-resultCh:ch<-res// 超时控制case<-timer.C:ch<-Result{Err:fmt.Errorf("async compute timeout after %s: %w",timeout,context.DeadlineExceeded)}}}()returnch}可以看到这里是开了两个协程,一个协程用来控制超时,一个用来处理业务逻辑。
2.3 多任务并发执行
funcasyncTask(namestring,duration time.Duration)<-chanResult{ch:=make(chanResult,1)gofunc(){deferclose(ch)fmt.Printf("[%s]%s start\n",util.Now(),name)// 模拟任务请求耗时time.Sleep(duration)fmt.Printf("[%s]%s done\n",util.Now(),name)ch<-Result{Data:fmt.Sprintf("%s finished in %s",name,duration),}}()returnch}funcwaitAll(tasks...<-chanResult)[]Result{results:=make([]Result,len(tasks))// 阻塞等待所有任务的返回结果fori,task:=rangetasks{results[i]=<-task}returnresults}多任务就是创建多个 goroutine 来执行任务,然后通过for-range来阻塞等待结果,下面来看下调用。
funcmain(){fmt.Printf("[%s]wait all tasks\n",util.Now())results:=waitAll(asyncTask("task-1",1*time.Second),asyncTask("task-2",500*time.Millisecond),asyncTask("task-3",2*time.Second),)fori,result:=rangeresults{fmt.Printf("[%s]result-%d data=%v err=%v\n",util.Now(),i+1,result.Data,result.Err)}}输出如下:
[2026-03-2915:26:07]wait all tasks[2026-03-2915:26:07]task-3start[2026-03-2915:26:07]task-1start[2026-03-2915:26:07]task-2start[2026-03-2915:26:08]task-2done[2026-03-2915:26:08]task-1done[2026-03-2915:26:09]task-3done[2026-03-2915:26:09]result-1data=task-1finished in 1s err=<nil>[2026-03-2915:26:09]result-2data=task-2finished in 500ms err=<nil>[2026-03-2915:26:09]result-3data=task-3finished in 2s err=<nil>3. 小结
好了,go 的协程先学到这里,先把后续的数据库和 web 操作学了再说。
如有错误,欢迎指出!!!