Go语言同步与异步执行多个任务封装详解(Runner和RunnerAsync)

 更新时间:2018年01月12日 11:00:42   作者:雪山飞猪  
这篇文章主要给大家介绍了关于Go语言同步与异步执行多个任务封装(Runner和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>]

总结

以上就是这篇文章的全部内容了,希望本文的内容对大家的学习或者工作具有一定的参考学习价值,如果有疑问大家可以留言交流,谢谢大家对脚本之家的支持。

相关文章

  • 深入理解gorm如何和数据库建立连接

    深入理解gorm如何和数据库建立连接

    这篇文章主要为大家详细介绍了gorm如何和数据库建立连接,文中的示例代码讲解详细,对我们深入了解GO语言有一定的帮助,需要的小伙伴可以参考下
    2023-11-11
  • 一文搞懂Golang中的内存逃逸

    一文搞懂Golang中的内存逃逸

    我们都知道go语言中内存管理工作都是由Go在底层完成的,这样我们可以不用过多的关注底层的内存问题。本文主要总结一下 Golang内存逃逸分析,需要的朋友可以参考以下内容,希望对大家有帮助
    2022-09-09
  • Go语言基础数组用法及示例详解

    Go语言基础数组用法及示例详解

    这篇文章主要为大家介绍了Go语言基础Go语言数组的用法及示例详解,有需要的朋友可以借鉴参考下,希望能够有所帮助,祝大家多多进步,早日升职加薪
    2021-11-11
  • Go语言中读取命令参数的几种方法总结

    Go语言中读取命令参数的几种方法总结

    部署golang项目时难免要通过命令行来设置一些参数,那么在golang中如何操作命令行参数呢?那么下面这篇文章就来给大家介绍了关于Go语言中读取命令参数的几种方法,文中通过示例代码介绍的非常详细,需要的朋友可以参考借鉴,下面随着小编来一起看看吧。
    2017-11-11
  • Golang中的强大Web框架Fiber详解

    Golang中的强大Web框架Fiber详解

    在不断发展的Web开发领域中,选择正确的框架可以极大地影响项目的效率和成功,介绍一下Fiber,这是一款令人印象深刻的Golang(Go语言)Web框架,在本文中,我们将深入了解Fiber的世界,探讨其独特的特性,并理解为什么它在Go生态系统中引起了如此大的关注
    2023-10-10
  • 使用Go语言创建WebSocket服务的实现示例

    使用Go语言创建WebSocket服务的实现示例

    这篇文章主要介绍了使用Go语言创建WebSocket服务的实现示例,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来一起学习学习吧
    2020-03-03
  • Go1.18新特性工作区模糊测试及泛型的使用详解

    Go1.18新特性工作区模糊测试及泛型的使用详解

    这篇文章主要为大家介绍了Go 1.18新特性中的工作区 模糊测试 泛型使用进行详细讲解,有需要的朋友可以借鉴参考下,希望能够有所帮助,祝大家多多进步,早日升职加薪
    2022-07-07
  • Go语言学习教程之反射的示例详解

    Go语言学习教程之反射的示例详解

    这篇文章主要通过记录对reflect包的简单使用,来对反射有一定的了解。文中的示例代码讲解详细,对我们学习Go语言有一定帮助,需要的可以参考一下
    2022-09-09
  • Go 语言 json解析框架与 gjson 详解

    Go 语言 json解析框架与 gjson 详解

    这篇文章主要介绍了Go语言json解析框架与gjson,JSON 解析是我们不可避免的常见问题,在Go语言中,我们可以借助gjson库来方便的进行json属性的提取与解析,需要的朋友可以参考一下
    2022-07-07
  • golang gorm的Callbacks事务回滚对象操作示例

    golang gorm的Callbacks事务回滚对象操作示例

    这篇文章主要为大家介绍了golang gorm的Callbacks事务回滚对象操作示例,有需要的朋友可以借鉴参考下,希望能够有所帮助,祝大家多多进步早日升职加薪
    2022-04-04

最新评论