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:[]
总结
以上就是这篇文章的全部内容了,希望本文的内容对大家的学习或者工作具有一定的参考学习价值,如果有疑问大家可以留言交流,谢谢大家对毛票票的支持。