Go语言同步与异步执行多个任务封装详解(Runner和RunnerAsync)
程序员文章站
2022-07-05 12:05:30
前言
同步适合多个连续执行的,每一步的执行依赖于上一步操作,异步执行则和任务执行顺序无关(如从10个站点抓取数据)
同步执行类runnerasync
支持返回超时...
前言
同步适合多个连续执行的,每一步的执行依赖于上一步操作,异步执行则和任务执行顺序无关(如从10个站点抓取数据)
同步执行类runnerasync
支持返回超时检测,系统中断检测
错误常量定义
//超时错误 var errtimeout = errors.new("received timeout") //操作系统系统中断错误 var errinterrupt = errors.new("received interrupt")
实现代码如下
package task import ( "os" "time" "os/signal" "sync" ) //异步执行任务 type runner struct { //操作系统的信号检测 interrupt chan os.signal //记录执行完成的状态 complete chan error //超时检测 timeout <-chan time.time //保存所有要执行的任务,顺序执行 tasks []func(id int) error waitgroup sync.waitgroup lock sync.mutex errs []error } //new一个runner对象 func newrunner(d time.duration) *runner { return &runner{ interrupt: make(chan os.signal, 1), complete: make(chan error), timeout: time.after(d), waitgroup: sync.waitgroup{}, lock: sync.mutex{}, } } //添加一个任务 func (this *runner) add(tasks ...func(id int) error) { this.tasks = append(this.tasks, tasks...) } //启动runner,监听错误信息 func (this *runner) start() error { //接收操作系统信号 signal.notify(this.interrupt, os.interrupt) //并发执行任务 go func() { this.complete <- this.run() }() select { //返回执行结果 case err := <-this.complete: return err //超时返回 case <-this.timeout: return errtimeout } } //异步执行所有的任务 func (this *runner) run() error { for id, task := range this.tasks { if this.gotinterrupt() { return errinterrupt } this.waitgroup.add(1) go func(id int) { this.lock.lock() //执行任务 err := task(id) //加锁保存到结果集中 this.errs = append(this.errs, err) this.lock.unlock() this.waitgroup.done() }(id) } this.waitgroup.wait() return nil } //判断是否接收到操作系统中断信号 func (this *runner) gotinterrupt() bool { select { case <-this.interrupt: //停止接收别的信号 signal.stop(this.interrupt) return true //正常执行 default: return false } } //获取执行完的error func (this *runner) geterrs() []error { return this.errs }
使用方法
add添加一个任务,任务为接收int类型的一个闭包
start开始执行伤,返回一个error类型,nil为执行完毕, errtimeout代表执行超时,errinterrupt代表执行被中断(类似ctrl + c操作)
测试示例代码
package task import ( "testing" "time" "fmt" "os" "runtime" ) func testrunnerasync_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("同步执行任务") //开始执行任务 if err := runner.start(); err != nil { switch err { case errtimeout: fmt.println("执行超时") os.exit(1) case errinterrupt: fmt.println("任务被中断") os.exit(2) } } t.log("执行结束") } //创建要执行的任务 func createtaskasync() func(id int) { return func(id int) { 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
支持返回超时检测,系统中断检测
实现代码如下
package task import ( "os" "time" "os/signal" "sync" ) //异步执行任务 type runner struct { //操作系统的信号检测 interrupt chan os.signal //记录执行完成的状态 complete chan error //超时检测 timeout <-chan time.time //保存所有要执行的任务,顺序执行 tasks []func(id int) error waitgroup sync.waitgroup lock sync.mutex errs []error } //new一个runner对象 func newrunner(d time.duration) *runner { return &runner{ interrupt: make(chan os.signal, 1), complete: make(chan error), timeout: time.after(d), waitgroup: sync.waitgroup{}, lock: sync.mutex{}, } } //添加一个任务 func (this *runner) add(tasks ...func(id int) error) { this.tasks = append(this.tasks, tasks...) } //启动runner,监听错误信息 func (this *runner) start() error { //接收操作系统信号 signal.notify(this.interrupt, os.interrupt) //并发执行任务 go func() { this.complete <- this.run() }() select { //返回执行结果 case err := <-this.complete: return err //超时返回 case <-this.timeout: return errtimeout } } //异步执行所有的任务 func (this *runner) run() error { for id, task := range this.tasks { if this.gotinterrupt() { return errinterrupt } this.waitgroup.add(1) go func(id int) { this.lock.lock() //执行任务 err := task(id) //加锁保存到结果集中 this.errs = append(this.errs, err) this.lock.unlock() this.waitgroup.done() }(id) } this.waitgroup.wait() return nil } //判断是否接收到操作系统中断信号 func (this *runner) gotinterrupt() bool { select { case <-this.interrupt: //停止接收别的信号 signal.stop(this.interrupt) return true //正常执行 default: return false } } //获取执行完的error func (this *runner) geterrs() []error { return this.errs }
使用方法
add添加一个任务,任务为接收int类型,返回类型error的一个闭包
start开始执行伤,返回一个error类型,nil为执行完毕, errtimeout代表执行超时,errinterrupt代表执行被中断(类似ctrl + c操作)
geterrs获取所有的任务执行结果
测试示例代码
package task import ( "testing" "time" "fmt" "os" "runtime" ) func testrunner_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("异步执行任务") //开始执行任务 if err := runner.start(); err != nil { switch err { case errtimeout: fmt.println("执行超时") os.exit(1) case errinterrupt: fmt.println("任务被中断") os.exit(2) } } t.log("执行结束") t.log(runner.geterrs()) } //创建要执行的任务 func createtask() func(id int) error { return func(id int) error { fmt.printf("正在执行%v个任务\n", id) //模拟任务执行,sleep //time.sleep(1 * time.second) return nil } }
执行结果
异步执行任务 正在执行2个任务 正在执行1个任务 正在执行4个任务 正在执行3个任务 正在执行6个任务 正在执行5个任务 正在执行9个任务 正在执行7个任务 正在执行10个任务 正在执行13个任务 正在执行8个任务 正在执行11个任务 正在执行12个任务 正在执行0个任务 runner_test.go:49: 执行结束 runner_test.go:51: [<nil> <nil> <nil> <nil> <nil> <nil> <nil> <nil> <nil> <nil> <nil> <nil> <nil> <nil>]
总结
以上就是这篇文章的全部内容了,希望本文的内容对大家的学习或者工作具有一定的参考学习价值,如果有疑问大家可以留言交流,谢谢大家对的支持。
上一篇: jdk下httpserver源码解析
下一篇: Linux安装Redis