Java 转 go 学习 - 并发编程(3)
发布时间:2026/10/6 2:14:37 作者:尧图编辑部 阅读量:1,286
)
文章目录1. 惰性生成器1.1 基于通道创建1.2 Iter.Seq 实现1.3 自己实现一个 Seql1.4 生成更通用的惰性生成器2. Future 模式2.1 基础 Future 模式实现2.2 基础使用2.2.1 http 异步调用2.2.2 带超时控制的 Future2.3 多任务并发执行3. 小结本系列文章Java 转 go 学习 - 项目管理Java 转 go 学习 - 基本语法Java 转 go 学习 - 类型转换Java 转 go 学习 - 流程控制结构Java 转 go 学习 - 数组和切片Java 转 go 学习 - mapJava 转 go 学习 - 函数1Java 转 go 学习 - 函数2Java 转 go 学习 - 结构体Java 转 go 学习 - 接口Java 转 go 学习 - 并发编程1Java 转 go 学习 - 并发编程21. 惰性生成器惰性生成器就是按需计算并返回值的机制通过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}})}()nextfunc()(V,bool){v,ok:-chreturnv,ok}stopfunc(){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我们可能会写成下面的写法。stopfunc(){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因为通道关闭了。到这里程序1stop()才收到信号结束。总之记住一句话调用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}})}()nextfunc()(V,bool){v,ok:-chreturnv,ok}stopfunc(){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 AnyinitState// retVal 是经过我们的函数 evalFunc 计算出来的值varretVal Any// 循环写入通道for{// 将初始值传进去, 可以进行 -*/ 等操作retVal,actStateevalFunc(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:os2returnos,ns}even:BuildLazyIntEvaluator(evenFunc,0)fori:0;i10;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(fmtgo-learn-1/learn-4/utiltime)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(),i1,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-1datatask-1finished in 1s errnil[2026-03-2915:26:09]result-2datatask-2finished in 500ms errnil[2026-03-2915:26:09]result-3datatask-3finished in 2s errnil3. 小结好了go 的协程先学到这里先把后续的数据库和 web 操作学了再说。如有错误欢迎指出