ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

并发模式:Pipeline与流式处理

并发模式:Pipeline与流式处理 第29篇 并发模式Pipeline与流式处理摘要Pipeline将数据处理拆分为多个阶段每个阶段用goroutine和通道连接实现流式处理和背压控制。本文从一个日志分析ETL场景说起讲清楚多阶段流水线的搭建和背压机制。一个日志分析需求上个月做了个日志分析平台需求是把应用产生的原始日志做三步处理解析JSON、过滤无效日志、聚合统计。日志量每天几十亿条串行处理吞吐量完全不够。最初的想法是分三批服务器第一批解析第二批过滤第三批聚合每批之间用消息队列连接。但这就过度设计了其实在一个进程里用 Pipeline 模式就能搞定。Pipeline 的思路是把处理流程拆成多个阶段每个阶段是一个或多个 goroutine阶段之间用通道连接。数据从第一阶段流向最后阶段每个阶段并行处理。好处是各阶段同时工作吞吐量取决于最慢的阶段天然支持流式处理不用等全部数据读完再处理。单阶段流水线先从一个简单的两阶段流水线开始数据生成和处理。packagemainimport(fmttime)// generate 日志生成阶段产生原始数据写入通道funcgenerate(countint)-chanstring{out:make(chanstring,10)// 带缓冲吸收速率差gofunc(){deferclose(out)// 生成完毕关闭通道fori:rangecount{// 模拟JSON格式的原始日志out-fmt.Sprintf({id:%d,msg:hello},i)time.Sleep(10*time.Millisecond)// 模拟产生间隔}}()returnout}// parse 解析阶段从输入通道读数据处理后写入输出通道funcparse(in-chanstring)-chanstring{out:make(chanstring,10)gofunc(){deferclose(out)forraw:rangein{// 持续读取直到上游关闭// 模拟解析这里简化处理parsed:parsed:raw time.Sleep(20*time.Millisecond)// 解析比生成慢out-parsed}}()returnout}funcmain(){start:time.Now()// 流水线连接generate的输出是parse的输入parsed:parse(generate(20))// 最终消费forp:rangeparsed{fmt.Println(p)}fmt.Printf(总耗时 %v\n,time.Since(start))}这里有个细节generate 每条间隔10毫秒parse 每条处理20毫秒。parse 是瓶颈总耗时约400毫秒。启动多个 goroutine 并行解析就能提升吞吐量。多阶段并行流水线真实场景流水线会有三四个阶段每个阶段的处理速度不同慢的阶段需要多个 worker 并行才能匹配上游速度。下面是一个完整的日志ETL流水线解析、过滤、聚合三阶段过滤阶段用4个worker并行加速。packagemainimport(contextfmtstringssynctime)// LogEntry 日志结构体typeLogEntrystruct{IDintMsgstringValidbool}// generateLog 生成原始日志字符串funcgenerateLog(countint)-chanstring{out:make(chanstring,20)gofunc(){deferclose(out)fori:rangecount{// 模拟原始JSON日志out-fmt.Sprintf({id:%d,msg:log-%s},i,strings.Repeat(x,3))}}()returnout}// stage1Parse 解析阶段把原始字符串解析为LogEntryfuncstage1Parse(ctx context.Context,in-chanstring)-chanLogEntry{out:make(chanLogEntry,20)gofunc(){deferclose(out)forraw:rangein{// 简化解析逻辑提取id字段id:0fmt.Sscanf(raw,{id:%d,id)entry:LogEntry{ID:id,Msg:raw}// 输出时监听取消信号select{caseout-entry:case-ctx.Done():return}}}()returnout}// fanOutStage 对慢阶段做并行扇出// n个worker同时处理同一个输入通道结果合并输出funcfanOutStage(ctx context.Context,nint,in-chanLogEntry,processfunc(LogEntry)LogEntry)-chanLogEntry{outs:make([]-chanLogEntry,n)fori:rangen{out:make(chanLogEntry,10)gofunc(){deferclose(out)forentry:rangein{result:process(entry)// 每个worker独立处理select{caseout-result:case-ctx.Done():return}}}()outs[i]out}returnfanInLogEntry(ctx,outs)// 合并多个worker的输出}// fanInLogEntry 合并多个LogEntry通道为一个funcfanInLogEntry(ctx context.Context,channels[]-chanLogEntry)-chanLogEntry{varwg sync.WaitGroup out:make(chanLogEntry,20)for_,ch:rangechannels{wg.Add(1)gofunc(c-chanLogEntry){deferwg.Done()forentry:rangec{select{caseout-entry:case-ctx.Done():return}}}(ch)}// 等所有输入读完再关闭输出通道gofunc(){wg.Wait()close(out)}()returnout}// stage3Aggregate 聚合阶段统计有效日志数量funcstage3Aggregate(ctx context.Context,in-chanLogEntry)-chanstring{out:make(chanstring,5)gofunc(){deferclose(out)count:0forentry:rangein{ifentry.Valid{// 只统计有效日志count}}// 输出最终统计结果select{caseout-fmt.Sprintf(有效日志总数 %d,count):case-ctx.Done():}}()returnout}funcmain(){ctx:context.Background()start:time.Now()// 阶段1: 解析原始日志parsed:stage1Parse(ctx,generateLog(100))// 阶段2: 过滤4个worker并行处理// id为偶数的视为有效日志filtered:fanOutStage(ctx,4,parsed,func(e LogEntry)LogEntry{time.Sleep(5*time.Millisecond)// 模拟过滤耗时e.Valide.ID%20returne})// 阶段3: 聚合统计有效日志result:stage3Aggregate(ctx,filtered)forr:rangeresult{fmt.Println(r)}fmt.Printf(总耗时 %v\n,time.Since(start))}这个流水线有三个阶段解析、过滤、聚合。过滤阶段用了4个worker并行因为过滤是耗时操作。所有阶段通过通道串联数据从左到右流动各阶段同时工作。独家踩坑背压失效导致内存暴涨这个坑差点让我背P0故障。流水线上线后跑了一段时间监控发现内存从2G涨到16G。排查发现是某个阶段的处理速度突然变慢上游还在高速生产通道缓冲区被打满后数据堆积在通道里。Go 的带缓冲通道在缓冲区满时会阻塞发送方这本来是天然的背压机制。问题出在缓冲设置不合理上游通道的缓冲设了100008个worker消费不过来10000条数据全积压在通道里每条数据还带着原始日志体内存就爆了。// 问题代码缓冲设置不合理funcstageBuggy(in-chanint)-chanint{workers:8// 缓冲设了10000太大掩盖了下游慢的问题out:make(chanint,10000)// 积压数据全堆在这fori:rangeworkers{gofunc(){forv:rangein{time.Sleep(50*time.Millisecond)// 慢处理out-v}}()}returnout}缓冲不是越大越好太大会掩盖下游慢的问题导致积压。我的经验是缓冲设为 worker 数量的2到4倍就够了让背压快速传导到上游。// 修复后的合理缓冲funcstageFixed(in-chanint)-chanint{workers:8// 缓冲设为worker数的2倍既有一点吸收能力又能及时背压out:make(chanint,workers*2)varwg sync.WaitGroupfori:rangeworkers{wg.Add(1)gofunc(idint){deferwg.Done()forv:rangein{time.Sleep(50*time.Millisecond)out-v}}(i)}gofunc(){wg.Wait()close(out)// 所有worker完成后关闭输出}()returnout}修复后内存曲线立刻平稳。背压的核心是让压力沿流水线向上游传导下游慢了上游就别生产缓冲只是吸收短时波动的减震器。对比分析维度Pipeline流水线批处理Worker Pool数据处理流式边读边处理攒批后处理单阶段并发内存占用低背压控制高全量加载中等首条延迟低快速产出高等攒批中等阶段并行各阶段同时工作无无适用场景ETL、流式日志报表统计任务队列Pipeline 的核心优势是流式处理和阶段并行。数据不需要全部加载到内存第一条数据在后续阶段还在处理时就能被消费。批处理适合需要全量数据的场景Worker Pool 适合单阶段任务分发。总结预告Pipeline 把数据处理拆成多个阶段用通道连接实现流式处理和天然背压。核心要点是每个阶段用 goroutine 独立处理通道缓冲要合理设置让背压及时传导配合 Context 做优雅退出。Go 并发编程系列到此讲完。从 goroutine 和 channel 基础到 sync 工具箱、Context、Worker Pool、Fan-in/Fan-out、Pipeline覆盖了日常开发80%的并发场景。下一篇进入模块三聊 Go 在云原生中的工程实践。
返回列表