【问题标题】:Throttle number of concurrent executing processes via buffered channels (Golang)通过缓冲通道(Golang)限制并发执行进程的数量
【发布时间】:2018-06-11 17:52:33
【问题描述】:

意图:

我正在寻找一种方法来并行运行操作系统级别的 shell 命令,但要小心不要破坏 CPU,我想知道缓冲通道是否适合这种用例。

已实现:

创建一系列具有模拟运行时持续时间的Jobs。将这些作业发送到一个队列,该队列将通过EXEC_THROTTLE 限制的缓冲通道将它们dispatch 发送到run。

观察:

这个“有效”(在它编译和运行的范围内),但我想知道缓冲区是否按规定工作(参见:“意图”)以限制并行运行的进程数。

免责声明:

现在,我知道新手倾向于过度使用渠道,但我觉得这种洞察力的要求是诚实的,因为我至少克制了使用sync.WaitGroup。请原谅这个有点玩具的例子,但所有的见解都会受到赞赏。

Playground

package main

import (
    // "os/exec"
    "log"
    "math/rand"
    "strconv"
    "sync"
    "time"
)

const (
    EXEC_THROTTLE = 2
)

type JobsManifest []Job

type Job struct {
    cmd     string
    result  string
    runtime int // Simulate long-running task
}

func (j JobsManifest) queueJobs(logChan chan<- string, runChan chan Job, wg *sync.WaitGroup) {
    go dispatch(logChan, runChan)
    for _, job := range j {
        wg.Add(1)
        runChan <- job
    }
}

func dispatch(logChan chan<- string, runChan chan Job) {
    for j := range runChan {
        go run(j, logChan)
    }
}

func run(j Job, logChan chan<- string) {
    time.Sleep(time.Second * time.Duration(j.runtime))
    j.result = strconv.Itoa(rand.Intn(10)) // j.result = os.Exec("/bin/bash", "-c", j.cmd).Output()
    logChan <- j.result
    log.Printf("   ran: %s\n", j.cmd)
}

func logger(logChan <-chan string, wg *sync.WaitGroup) {
    for {
        res := <-logChan
        log.Printf("logged: %s\n", res)
        wg.Done()
    }
}

func main() {

    jobs := []Job{
        Job{
            cmd:     "ps -p $(pgrep vim) | tail -n 1 | awk '{print $3}'",
            runtime: 1,
        },
        Job{
            cmd:     "wc -l /var/log/foo.log | awk '{print $1}'",
            runtime: 2,
        },
        Job{
            cmd:     "ls -l ~/go/src/github.com/ | wc -l | awk '{print $1}'",
            runtime: 3,
        },
        Job{
            cmd:     "find /var/log/ -regextype posix-extended -regex '.*[0-9]{10}'",
            runtime: 4,
        },
    }

    var wg sync.WaitGroup
    logChan := make(chan string)
    runChan := make(chan Job, EXEC_THROTTLE)
    go logger(logChan, &wg)

    start := time.Now()
    JobsManifest(jobs).queueJobs(logChan, runChan, &wg)
    wg.Wait()
    log.Printf("finish: %s\n", time.Since(start))
}

【问题讨论】:

    标签: go concurrency channel


    【解决方案1】:

    您还可以使用缓冲通道限制并发:

    concurrencyLimit := 2 // Number of simultaneous jobs.
    semaphore := make(chan struct{}, concurrencyLimit)
    for job := range jobs {
        job := job // Pin loop variable.
        semaphore <- struct{}{} // Reserve limiter slot.
        go func() {
            defer func() {
                <-semaphore // Free semaphore slot.
            }()
            
            do(job) // Do the job.
        }()
    }
    // Wait for goroutines to finish by filling full channel.
    for i := 0; i < cap(semaphore); i++ {
        semaphore <- struct{}{}
    }
    

    【讨论】:

      【解决方案2】:

      将 processItem 函数替换为所需的作业执行。

      下面将以正确的顺序执行作业。最多 EXEC_CONCURRENT 项将同时执行。

      package main
      
      import (
          "fmt"
          "sync"
          "time"
      )
      
      func processItem(i int, done chan int, wg *sync.WaitGroup) { 
          fmt.Printf("Async Start: %d\n", i)
          time.Sleep(100 * time.Millisecond * time.Duration(i))
          fmt.Printf("Async Complete: %d\n", i)
          done <- 1
          wg.Done()
      }
      
      func popItemFromBufferChannelWhenItemDoneExecuting(items chan int, done chan int) { 
          _ = <- done
          _ = <-items
      }
      
      
      func main() {
          EXEC_CONCURRENT := 3
      
          items := make(chan int, EXEC_CONCURRENT)
          done := make(chan int)
          var wg sync.WaitGroup
      
          for i:= 1; i < 11; i++ {
              items <- i
              wg.Add(1)   
              go processItem(i, done, &wg)
              go popItemFromBufferChannelWhenItemDoneExecuting(items, done)
          }
      
          wg.Wait()
      }
      

      下面将以随机顺序执行作业。最多 EXEC_CONCURRENT 项将同时执行。

      package main
      
      import (
          "fmt"
          "sync"
          "time"
      )
      
      func processItem(i int, items chan int, wg *sync.WaitGroup) { 
          items <- i
          fmt.Printf("Async Start: %d\n", i)
          time.Sleep(100 * time.Millisecond * time.Duration(i))
          fmt.Printf("Async Complete: %d\n", i)
          _ = <- items
          wg.Done()
      }
      
      func main() {
          EXEC_CONCURRENT := 3
      
          items := make(chan int, EXEC_CONCURRENT)
          var wg sync.WaitGroup
      
          for i:= 1; i < 11; i++ {
              wg.Add(1)   
              go processItem(i, items, &wg)
          }
      
          wg.Wait()
      }
      

      您可以根据自己的要求进行选择。

      【讨论】:

        【解决方案3】:

        如果我理解正确,您的意思是建立一种机制来确保在任何时候最多有多个EXEC_THROTTLE 作业在运行。如果这是您的意图,那么代码将不起作用。

        这是因为当您开始一项工作时,您已经消耗了频道 - 允许开始另一个工作,但没有完成任何工作。您可以通过添加一个计数器来调试它(您需要原子添加或互斥锁)。

        您可以通过简单地启动一组具有无缓冲通道的 goroutine 并在执行作业时阻塞来完成工作:

        func Run(j Job) r Result {
            //Run your job here
        }
        
        func Dispatch(ch chan Job) {
            for j:=range ch {
                wg.Add(1)
                Run(j)
                wg.Done()
            }
        }
        
        func main() {
            ch := make(chan Job)
            for i:=0; i<EXEC_THROTTLE; i++ {
                go Dispatch(ch)
            }
            //call dispatch according to the queue here.
        }
        

        它之所以有效,是因为只要有一个 goroutine 正在使用通道,这意味着至少有一个 goroutine 没有运行,并且最多有 EXEC_THROTTLE-1 作业在运行,所以最好再执行一个,它确实这样做了。

        【讨论】:

        • 感谢您的反馈@leafbebop,一旦我解决了这个问题,我会尝试建议并接受答案。
        • blog.golang.org/advanced-go-concurrency-patterns 您可能会发现这篇博文很有帮助。
        • 嗨@leafbebop。也许我对您的评论 //call dispatch according to the queue here. 感到困惑。我将队列更新为循环直到EXEC_THROTTLE 的长度,但正如您所见,如果EXEC_THROTTLE == 2 仅运行前两个Jobs。我错过了什么? play.golang.org/p/Ye35isSP8kv 另外,我想我需要等待logger 回电话wg.Done()。
        • 你确实错过了我的意思。我在这里更改了您的代码:play.golang.org/p/HhkDh3WYgBH 首先尝试通过我在答案中的解释来理解它。我发现很难解释更多,但如果仍有不清楚的地方,请随时询问。
        • 谢谢。我之前觉得有必要在 goroutine 中调用 run 的原因是我想阻塞 EXEC_THROTTLE 的值,而不是每个 Job。经编辑,似乎每个Job 都按顺序运行,而不是并行执行前两个,但现在通过研究您的答案,我意识到这是我的误解,因为它在run 返回后立即发送到logger .如果我编辑任何一对Jobs 使其具有相同的runtime 值,则很明显它是并行执行的。非常感谢您抽出宝贵时间来解决这个问题。
        【解决方案4】:

        我经常使用这个。 https://github.com/dustinevan/go-utils

        package async
        import (
            "context"
        
            "github.com/pkg/errors"
        )
        
        type Semaphore struct {
            buf    chan struct{}
            ctx    context.Context
            cancel context.CancelFunc
        }
        
        func NewSemaphore(max int, parentCtx context.Context) *Semaphore {
        
            s := &Semaphore{
                buf:    make(chan struct{}, max),
                ctx:    parentCtx,
            }
        
            go func() {
                <-s.ctx.Done()
                close(s.buf)
                drainStruct(s.buf)
            }()
        
            return s
        }
        
        var CLOSED = errors.New("the semaphore has been closed")
        
        func (s *Semaphore) Acquire() error {
            select {
            case <-s.ctx.Done():
                return CLOSED
            case s.buf <- struct{}{}:
                return nil
            }
        }
        
        func (s *Semaphore) Release() {
            <-s.buf
        }
        

        你会这样使用它:

        func main() {
        
            sem := async.NewSemaphore(10, context.Background())
            ...
            var wg sync.Waitgroup 
            for _, job := range jobs {
                go func() {
                    wg.Add(1)
                    err := sem.Acquire()
                    if err != nil {
                         // handle err, 
                    }
                    defer sem.Release()
                    defer wg.Done()
                    job()
            }
            wg.Wait()
        }
        

        【讨论】:

        • 您的代码的第二部分有无与伦比的手镯。我很难理解它是如何工作的,所以我不会为你编辑。
        • 好的,我明白了,并为您修好了。但我想说,在这样一个问题上它过于复杂了。
        • 好吧,如果你只是在做一个小游乐场应用程序,当然可以。但是,如果您尝试管理外部资源的使用,那么在使用资源时只需说:sem.Acquire() 然后sem.Release() 就非常好。如上所示,工人也是解决这个问题的一种方式。
        • 如果代码包含一些用于实际内存管理的选项,或用于分布式计算的功能,我将非常感谢这个摘要。但由于它只是围绕缓冲通道做一些基本的事情,所以代码对我来说似乎不自然和冗余。但这是一个非常个人的观点。
        猜你喜欢
        • 2017-07-22
        • 2016-08-30
        • 2014-09-26
        • 2020-09-03
        • 2021-07-01
        • 2018-07-03
        • 2021-10-13
        • 2019-01-17
        • 2019-10-26
        相关资源
        最近更新 更多