Go并发模式实战Pipeline与Fan-in-Fan-out生产案例
Go并发模式实战Pipeline与Fan-in-Fan-out生产案例文章导语Go的并发模式不仅仅是goroutinechannelPipeline流水线、Fan-out扇出和Fan-in扇入是Go并发编程中最经典的设计模式。它们将复杂任务分解为可组合的并发阶段在高并发数据处理、消息处理和实时流计算中发挥着关键作用。一、Pipeline流水线模式1.1 基础流水线// 一个完整的Pipeline生成→处理→消费funcgen(nums...int)-chanint{out:make(chanint)gofunc(){for_,n:rangenums{out-n}close(out)}()returnout}funcsquare(in-chanint)-chanint{out:make(chanint)gofunc(){forn:rangein{out-n*n}close(out)}()returnout}funcmain(){// 管道连接c:gen(2,3)out:square(c)forv:rangeout{fmt.Println(v)// 4, 9}}1.2 Pipeline的取消机制funcgenWithCancel(ctx context.Context,nums...int)-chanint{out:make(chanint)gofunc(){deferclose(out)for_,n:rangenums{select{caseout-n:case-ctx.Done():return}}}()returnout}funcsquareWithCancel(ctx context.Context,in-chanint)-chanint{out:make(chanint)gofunc(){deferclose(out)forn:rangein{select{caseout-n*n:case-ctx.Done():return}}}()returnout}二、扇出Fan-out模式扇出将同一个输入channel分发给多个worker并发处理。funcfanOut(in-chanWork,workerCountint)[]-chanResult{workers:make([]-chanResult,workerCount)fori:0;iworkerCount;i{workers[i]worker(i,in)}returnworkers}funcworker(idint,in-chanWork)-chanResult{out:make(chanResult)gofunc(){deferclose(out)forw:rangein{out-process(id,w)}}()returnout}三、扇入Fan-in模式扇入将多个channel合并为一个。funcfanIn(channels...-chanResult)-chanResult{varwg sync.WaitGroup out:make(chanResult)multiplex:func(c-chanResult){deferwg.Done()forr:rangec{out-r}}wg.Add(len(channels))for_,c:rangechannels{gomultiplex(c)}gofunc(){wg.Wait()close(out)}()returnout}四、实战并发爬虫系统funcCrawl(urls[]string,concurrencyint)[]Page{ctx,cancel:context.WithTimeout(context.Background(),30*time.Second)defercancel()// Stage 1: URL生成器urlCh:make(chanstring,len(urls))gofunc(){deferclose(urlCh)for_,u:rangeurls{select{caseurlCh-u:case-ctx.Done():return}}}()// Stage 2: Fan-out——并发抓取fetchChs:make([]-chanFetchedPage,concurrency)fori:0;iconcurrency;i{fetchChs[i]fetchWorker(ctx,urlCh)}// Stage 3: Fan-in——合并结果mergedCh:fanIn(fetchChs...)// Stage 4: 结果处理varpages[]Pageforpage:rangemergedCh{ifpage.Err!nil{log.Printf(抓取失败 %s: %v,page.URL,page.Err)continue}pagesappend(pages,page.Page)}returnpages}funcfetchWorker(ctx context.Context,urls-chanstring)-chanFetchedPage{out:make(chanFetchedPage)gofunc(){deferclose(out)forurl:rangeurls{resp,err:http.Get(url)fp:FetchedPage{URL:url}iferr!nil{fp.Errerr}else{body,_:io.ReadAll(resp.Body)resp.Body.Close()fp.Bodybody}select{caseout-fp:case-ctx.Done():return}}}()returnout}五、Pipeline最佳实践5.1 显式关闭channel// 原则发送方负责关闭channelfuncproducer()-chanint{ch:make(chanint)gofunc(){deferclose(ch)// 发送方关闭fori:0;i10;i{ch-i}}()returnch}5.2 缓冲channel减少阻塞// 根据处理速度差异设置buffer// 如果生产者快于消费者设置buffer避免阻塞ch:make(chanData,100)5.3 合并错误处理typeResultstruct{Valueinterface{}Errerror}funcprocessWithErrors(in-chanTask)-chanResult{out:make(chanResult)gofunc(){deferclose(out)fortask:rangein{val,err:process(task)out-Result{Value:val,Err:err}}}()returnout}六、全文总结Pipeline将任务分解为顺序处理的阶段Fan-out将工作分发到多个goroutine并行处理Fan-in将多个channel合并为一个发送方关闭channel是最安全的模式context贯穿整个Pipeline实现取消传播七、技术进阶展望errgroup在Pipeline中的集成有状态Pipeline与无状态Pipeline响应式编程在Go中的实现参考文献Go Blog - Go Concurrency Patterns: Pipelines and cancellationGo Blog - Advanced Go Concurrency PatternsKatherine Cox-Buday - Concurrency in GoArdan Labs - Ultimate Go: ConcurrencyGo by Example - Worker Pools