Go语言同步与异步执行多个任务封装详解(Runner和RunnerAsync)
前言
同步适合多个连续执行的,每一步的执行依赖于上一步操作,异步执行则和任务执行顺序无关(如从10个站点抓取数据)
同步执行类RunnerAsync
支持返回超时检测,系统中断检测
错误常量定义
//超时错误 varErrTimeout=errors.New("receivedtimeout") //操作系统系统中断错误 varErrInterrupt=errors.New("receivedinterrupt")
实现代码如下
packagetask import( "os" "time" "os/signal" "sync" ) //异步执行任务 typeRunnerstruct{ //操作系统的信号检测 interruptchanos.Signal //记录执行完成的状态 completechanerror //超时检测 timeout<-chantime.Time //保存所有要执行的任务,顺序执行 tasks[]func(idint)error waitGroupsync.WaitGroup locksync.Mutex errs[]error } //new一个Runner对象 funcNewRunner(dtime.Duration)*Runner{ return&Runner{ interrupt:make(chanos.Signal,1), complete:make(chanerror), timeout:time.After(d), waitGroup:sync.WaitGroup{}, lock:sync.Mutex{}, } } //添加一个任务 func(this*Runner)Add(tasks...func(idint)error){ this.tasks=append(this.tasks,tasks...) } //启动Runner,监听错误信息 func(this*Runner)Start()error{ //接收操作系统信号 signal.Notify(this.interrupt,os.Interrupt) //并发执行任务 gofunc(){ this.complete<-this.Run() }() select{ //返回执行结果 caseerr:=<-this.complete: returnerr //超时返回 case<-this.timeout: returnErrTimeout } } //异步执行所有的任务 func(this*Runner)Run()error{ forid,task:=rangethis.tasks{ ifthis.gotInterrupt(){ returnErrInterrupt } this.waitGroup.Add(1) gofunc(idint){ this.lock.Lock() //执行任务 err:=task(id) //加锁保存到结果集中 this.errs=append(this.errs,err) this.lock.Unlock() this.waitGroup.Done() }(id) } this.waitGroup.Wait() returnnil } //判断是否接收到操作系统中断信号 func(this*Runner)gotInterrupt()bool{ select{ case<-this.interrupt: //停止接收别的信号 signal.Stop(this.interrupt) returntrue //正常执行 default: returnfalse } } //获取执行完的error func(this*Runner)GetErrs()[]error{ returnthis.errs }
使用方法
Add添加一个任务,任务为接收int类型的一个闭包
Start开始执行伤,返回一个error类型,nil为执行完毕,ErrTimeout代表执行超时,ErrInterrupt代表执行被中断(类似Ctrl+C操作)
测试示例代码
packagetask import( "testing" "time" "fmt" "os" "runtime" ) funcTestRunnerAsync_Start(t*testing.T){ //开启多核 runtime.GOMAXPROCS(runtime.NumCPU()) //创建runner对象,设置超时时间 runner:=NewRunnerAsync(8*time.Second) //添加运行的任务 runner.Add( createTaskAsync(), createTaskAsync(), createTaskAsync(), createTaskAsync(), createTaskAsync(), createTaskAsync(), createTaskAsync(), createTaskAsync(), createTaskAsync(), createTaskAsync(), createTaskAsync(), createTaskAsync(), createTaskAsync(), ) fmt.Println("同步执行任务") //开始执行任务 iferr:=runner.Start();err!=nil{ switcherr{ caseErrTimeout: fmt.Println("执行超时") os.Exit(1) caseErrInterrupt: fmt.Println("任务被中断") os.Exit(2) } } t.Log("执行结束") } //创建要执行的任务 funccreateTaskAsync()func(idint){ returnfunc(idint){ fmt.Printf("正在执行%v个任务\n",id) //模拟任务执行,sleep两秒 //time.Sleep(1*time.Second) } }
执行结果
同步执行任务 正在执行0个任务 正在执行1个任务 正在执行2个任务 正在执行3个任务 正在执行4个任务 正在执行5个任务 正在执行6个任务 正在执行7个任务 正在执行8个任务 正在执行9个任务 正在执行10个任务 正在执行11个任务 正在执行12个任务 runnerAsync_test.go:49:执行结束
异步执行类Runner
支持返回超时检测,系统中断检测
实现代码如下
packagetask import( "os" "time" "os/signal" "sync" ) //异步执行任务 typeRunnerstruct{ //操作系统的信号检测 interruptchanos.Signal //记录执行完成的状态 completechanerror //超时检测 timeout<-chantime.Time //保存所有要执行的任务,顺序执行 tasks[]func(idint)error waitGroupsync.WaitGroup locksync.Mutex errs[]error } //new一个Runner对象 funcNewRunner(dtime.Duration)*Runner{ return&Runner{ interrupt:make(chanos.Signal,1), complete:make(chanerror), timeout:time.After(d), waitGroup:sync.WaitGroup{}, lock:sync.Mutex{}, } } //添加一个任务 func(this*Runner)Add(tasks...func(idint)error){ this.tasks=append(this.tasks,tasks...) } //启动Runner,监听错误信息 func(this*Runner)Start()error{ //接收操作系统信号 signal.Notify(this.interrupt,os.Interrupt) //并发执行任务 gofunc(){ this.complete<-this.Run() }() select{ //返回执行结果 caseerr:=<-this.complete: returnerr //超时返回 case<-this.timeout: returnErrTimeout } } //异步执行所有的任务 func(this*Runner)Run()error{ forid,task:=rangethis.tasks{ ifthis.gotInterrupt(){ returnErrInterrupt } this.waitGroup.Add(1) gofunc(idint){ this.lock.Lock() //执行任务 err:=task(id) //加锁保存到结果集中 this.errs=append(this.errs,err) this.lock.Unlock() this.waitGroup.Done() }(id) } this.waitGroup.Wait() returnnil } //判断是否接收到操作系统中断信号 func(this*Runner)gotInterrupt()bool{ select{ case<-this.interrupt: //停止接收别的信号 signal.Stop(this.interrupt) returntrue //正常执行 default: returnfalse } } //获取执行完的error func(this*Runner)GetErrs()[]error{ returnthis.errs }
使用方法
Add添加一个任务,任务为接收int类型,返回类型error的一个闭包
Start开始执行伤,返回一个error类型,nil为执行完毕,ErrTimeout代表执行超时,ErrInterrupt代表执行被中断(类似Ctrl+C操作)
getErrs获取所有的任务执行结果
测试示例代码
packagetask import( "testing" "time" "fmt" "os" "runtime" ) funcTestRunner_Start(t*testing.T){ //开启多核心 runtime.GOMAXPROCS(runtime.NumCPU()) //创建runner对象,设置超时时间 runner:=NewRunner(18*time.Second) //添加运行的任务 runner.Add( createTask(), createTask(), createTask(), createTask(), createTask(), createTask(), createTask(), createTask(), createTask(), createTask(), createTask(), createTask(), createTask(), createTask(), ) fmt.Println("异步执行任务") //开始执行任务 iferr:=runner.Start();err!=nil{ switcherr{ caseErrTimeout: fmt.Println("执行超时") os.Exit(1) caseErrInterrupt: fmt.Println("任务被中断") os.Exit(2) } } t.Log("执行结束") t.Log(runner.GetErrs()) } //创建要执行的任务 funccreateTask()func(idint)error{ returnfunc(idint)error{ fmt.Printf("正在执行%v个任务\n",id) //模拟任务执行,sleep //time.Sleep(1*time.Second) returnnil } }
执行结果
异步执行任务 正在执行2个任务 正在执行1个任务 正在执行4个任务 正在执行3个任务 正在执行6个任务 正在执行5个任务 正在执行9个任务 正在执行7个任务 正在执行10个任务 正在执行13个任务 正在执行8个任务 正在执行11个任务 正在执行12个任务 正在执行0个任务 runner_test.go:49:执行结束 runner_test.go:51:[]
总结
以上就是这篇文章的全部内容了,希望本文的内容对大家的学习或者工作具有一定的参考学习价值,如果有疑问大家可以留言交流,谢谢大家对毛票票的支持。