Golang被证明非常适合并发编程,goroutine比异步编程更易读、优雅、高效。本文提出一个适合由Golang实现的Pipeline执行模型,适合批量处理大量数据(ETL)的情景。
想象这样的应用情景:
(1)从数据库A(Cassandra)加载用户评论(量巨大,例如10亿条);
(2)根据每条评论的用户ID、从数据库B(MySQL)关联用户资料;
(3)调用NLP服务(自然语言处理),处理每条评论;
(4)将处理结果写入数据库C(ElasticSearch)。
由于应用中遇到的各种问题,归纳出这些需求:
需求一:应分批处理数据,例如规定每批100条。出现问题时(例如任意一个数据库故障)则中断,下次程序启动时使用checkpoint从中断处恢复。
需求二:每个流程设置合理的并发数、让数据库和NLP服务有合理的负载(不影响其它业务的基础上,尽可能占用更多资源以提高ETL性能)。例如,步骤(1)-(4)分别设置并发数1、8、32、2。
这就是一个典型的Pipeline(流水线)执行模型。把每一批数据(例如100条)看作流水线上的产品,4个步骤对应流水线上4个处理工序,每个工序处理完毕后就把半成品交给下一个工序。每个工序可以同时处理的产品数各不相同。
你可能首先想到启用1+8+32+2个goroutine,使用channel来传递半成品。我也曾经这么干,结论就是这么干会让程序员疯掉:流程并发控制代码非常复杂,特别是你得处理异常、执行时间超出预期、可控中断等问题,你不得不加入一堆channel,直到你自己都不记得有什么用。
为了更高效完成ETL工作,我将Pipeline抽象成模块。我先把代码粘贴出来,再解析含义。模块可以直接使用,主要使用的接口是:NewPipeline、Async、Wait。
1package main 2 3import "sync" 4 5func HasClosed(c <-chan struct{}) bool { 6 select { 7 case <-c: return true 8 default: return false 9 } 10} 11 12type SyncFlag interface{ 13 Wait() 14 Chan() <-chan struct{} 15 Done() bool 16} 17 18func NewSyncFlag() (done func(), flag SyncFlag) { 19 f := &syncFlag{ 20 c : make(chan struct{}), 21 } 22 return f.done, f 23} 24 25type syncFlag struct { 26 once sync.Once 27 c chan struct{} 28} 29 30func (f *syncFlag) done() { 31 f.once.Do(func(){ 32 close(f.c) 33 }) 34} 35 36func (f *syncFlag) Wait() { 37 <-f.c 38} 39 40func (f *syncFlag) Chan() <-chan struct{} { 41 return f.c 42} 43 44func (f *syncFlag) Done() bool { 45 return HasClosed(f.c) 46} 47 48type pipelineThread struct { 49 sigs []chan struct{} 50 chanExit chan struct{} 51 interrupt SyncFlag 52 setInterrupt func() 53 err error 54} 55 56func newPipelineThread(l int) *pipelineThread { 57 p := &pipelineThread{ 58 sigs : make([]chan struct{}, l), 59 chanExit : make(chan struct{}), 60 } 61 p.setInterrupt, p.interrupt = NewSyncFlag() 62 63 for i := range p.sigs { 64 p.sigs[i] = make(chan struct{}) 65 } 66 return p 67} 68 69type Pipeline struct { 70 mtx sync.Mutex 71 workerChans []chan struct{} 72 prevThd *pipelineThread 73} 74 75//创建流水线,参数个数是每个任务的子过程数,每个参数对应子过程的并发度。 76func NewPipeline(workers ...int) *Pipeline { 77 if len(workers) < 1 { panic("NewPipeline need aleast one argument") } 78 79 workersChan := make([]chan struct{}, len(workers)) 80 for i := range workersChan { 81 workersChan[i] = make(chan struct{}, workers[i]) 82 } 83 84 prevThd := newPipelineThread(len(workers)) 85 for _,sig := range prevThd.sigs { 86 close(sig) 87 } 88 close(prevThd.chanExit) 89 90 return &Pipeline{ 91 workerChans : workersChan, 92 prevThd : prevThd, 93 } 94} 95 96//往流水线推入一个任务。如果第一个步骤的并发数达到设定上限,这个函数会堵塞等待。 97//如果流水线中有其它任务失败(返回非nil),任务不被执行,函数返回false。 98func (p *Pipeline) Async(works ...func()error) bool { 99 if len(works) != len(p.workerChans) { 100 panic("Async: arguments number not matched to NewPipeline(...)") 101 } 102 103 p.mtx.Lock() 104 if p.prevThd.interrupt.Done() { 105 p.mtx.Unlock() 106 return false 107 } 108 prevThd := p.prevThd 109 thisThd := newPipelineThread(len(p.workerChans)) 110 p.prevThd = thisThd 111 p.mtx.Unlock() 112 113 lock := func(idx int) bool { 114 select { 115 case <-prevThd.interrupt.Chan(): return false 116 case <-prevThd.sigs[idx]: //wait for signal 117 } 118 select { 119 case <-prevThd.interrupt.Chan(): return false 120 case p.workerChans[idx]<-struct{}{}: //get lock 121 } 122 return true 123 } 124 if !lock(0) { 125 thisThd.setInterrupt() 126 <-prevThd.chanExit 127 thisThd.err = prevThd.err 128 close(thisThd.chanExit) 129 return false 130 } 131 go func() { //watch interrupt of previous thread 132 select { 133 case <-prevThd.interrupt.Chan(): 134 thisThd.setInterrupt() 135 case <-thisThd.chanExit: 136 } 137 }() 138 go func() { 139 var err error 140 for i,work := range works { 141 close(thisThd.sigs[i]) //signal next thread 142 if work != nil { 143 err = work() 144 } 145 if err != nil || (i+1 < len(works) && !lock(i+1)) { 146 thisThd.setInterrupt() 147 break 148 } 149 <-p.workerChans[i] //release lock 150 } 151 152 <-prevThd.chanExit 153 if prevThd.interrupt.Done() { 154 thisThd.setInterrupt() 155 } 156 if prevThd.err != nil { 157 thisThd.err = prevThd.err 158 } else { 159 thisThd.err = err 160 } 161 close(thisThd.chanExit) 162 }() 163 return true 164} 165 166//等待流水线中所有任务执行完毕或失败,返回第一个错误,如果无错误则返回nil。 167func (p *Pipeline) Wait() error { 168 p.mtx.Lock() 169 lastThd := p.prevThd 170 p.mtx.Unlock() 171 <-lastThd.chanExit 172 return lastThd.err 173}
使用这个Pipeline组件,我们的ETL程序将会简单、高效、可靠,让程序员从繁琐的并发流程控制中解放出来:
1package main 2 3import "log" 4 5func main() { 6 checkpoint := loadCheckpoint() 7 8 //工序(1)在pipeline外执行,最后一个工序是保存checkpoint 9 pipeline := NewPipeline(8, 32, 2, 1) 10 for { 11 //(1) 12 //加载100条数据,并修改变量checkpoint 13 //data是数组,每个元素是一条评论,之后的联表、NLP都直接修改data里的每条记录。 14 data, err := extractReviewsFromA(&checkpoint, 100) 15 if err != nil { 16 log.Print(err) 17 break 18 } 19 curCheckpoint := checkpoint 20 21 ok := pipeline.Async(func() error { 22 //(2) 23 return joinUserFromB(data) 24 }, func() error { 25 //(3) 26 return nlp(data) 27 }, func() error { 28 //(4) 29 return loadDataToC(data) 30 }, func() error { 31 //(5)保存checkpoint 32 log.Print("done:", curCheckpoint) 33 return saveCheckpoint(curCheckpoint) 34 }) 35 if !ok { break } 36 37 if len(data) < 100 { break } //处理完毕 38 } 39 err := pipeline.Wait() 40 if err != nil { log.Print(err) } 41}
Pipeline执行模型的特性:
1、Pipeline分别控制每一个工序的并发数,如果(4)的并发数已满,某个线程的(3)即使完成都会堵塞等待,直到(4)有一个线程完成。
2、在上面的情景中,Pipeline最多同时处理1+8+32+2+1=44个线程共4400条记录,内存开销可控。
3、每个线程的每个工序的调度,不早于上一个线程同一个工序的调度。
例如:有两个线程正在执行,<1>先执行、<2>后执行。如果<2>(4)早于<1>(4)完成,那<2>必须堵塞等待,直到<1>(4)完成、<1>(5)开始执行,那<2>(5)才会开始。又因为(5)的最大并发数是1,所以实际上<2>(5)必须等待<1>(5)完成才会开始。这个机制保证checkpoint的执行顺序一定是按照Async的顺序,避免中断、继续时漏处理数据。
4、如果某个线程的某个工序处理失败(例如数据库故障),那之后的线程都会中止执行,下一次调用Async返回false,pipeline.Wait()返回第一个错误,整个流水线作业可控中断。
例如:有三个线程正在执行:<1>、<2>、<3>。如果<2>(4)失败(loadDataToC返回error非nil),那<3>无论正在执行到哪一个工序,都不会进入下一个工序而中断。<1>不会受到影响,会一直执行完毕。Wait()等待<1><2><3>全部完成或中止,返回loadDataToC的错误。
5、无法避免中断过程中有checkpoint后的数据写入。下次重启程序将重新写入、覆盖这些数据。
例如:<2>(4)失败、<3>(4)执行成功(已写入数据),那<2>(5)和<3>(5)都不会被执行,checkpoint的最新状态是<1>写入的,下次重启程序将重新执行<2>和<3>,其中<3>的数据会再次写入,所以写入应该按照记录ID作覆盖写入。
6、你可以随时Ctrl+C、重启程序,所有事情都会继续有序执行。死机?毫无压力。
总结:Pipeline执行模型除了限制并发数,也能限制内存开销,对失败恢复有充足的考虑,让程序员从繁琐的并发编程中解放出来。
吐槽:Python程序员没办法用几百行代码就漂亮地完成这个任务。