【问题标题】:processing jobs from a neverending queue with a fixed number of workers处理具有固定数量工人的无休止队列中的作业
【发布时间】:2014-03-23 21:31:08
【问题描述】:

这是我的头,我不知道如何解决它;

  • 我想让固定数量 N 个 goroutine 并行运行
  • 我将从永无止境的队列中获取有关要处理的作业的 X 条消息
  • 我想让 N 个 goroutine 处理这些 X 个作业,一旦其中一个例程无事可做,我想从永无止境的队列中获取另一个 X 个作业

下面答案中的代码(请参阅 url)可以出色地处理任务,但是一旦该任务列表为空,工作人员就会死去,我希望他们保持活力并以某种方式通知主代码他们已经失业所以我可以获取更多的工作来填充任务列表

How would you define a pool of goroutines to be executed at once in Golang?

使用下面的 user:Jsor 示例代码,我尝试创建一个简单的程序,但我很困惑。

import (
    "fmt"
    "strconv"
)

//workChan - read only that delivers work
//requestChan - ??? what is this
func Worker(myid string, workChan <- chan string, requestChan chan<- struct{}) {
    for {
        select {
        case work := <-workChan:
            fmt.Println("Channel: " + myid + " do some work: " + work)
        case requestChan <- struct{}{}:
            //hm? how is the requestChan used?
        }
    }
}

func Logic(){

    workChan := make(chan string)
    requestChan := make(chan struct{})

    //Create the workers
    for i:=1; i < 5; i++ {
        Worker( strconv.Itoa( i), workChan, requestChan)
    }

    //Give the workers some work
    for i:=100; i < 115; i++ {
        workChan<- "workid"+strconv.Itoa( i)
    }

}

【问题讨论】:

    标签: concurrency go


    【解决方案1】:

    这就是select 语句的用途。

    func Worker(workChan chan<- Work, requestChan chan<- struct{}) {
        for {
            select {
            case work := <-workChan:
                // Do work
            case requestChan <- struct{}{}:
            }
        }
    }
    

    这个工人将永远运行下去。如果工作可用,它将从工作通道中提取。如果没有剩余,它将发送请求。

    不是这样,因为它永远运行,如果你想杀死一个工人,你需要做其他事情。一种可能性是始终使用 workChan 检查ok,如果该通道已关闭,则退出该功能。另一种选择是为每个工作人员使用单独的退出通道。

    【讨论】:

    • 另请注意,select 中的 casestatement 变为非阻塞,如果相应的通道已关闭。
    • 谢谢!我仍然很困惑它是如何工作的;什么是 requestChan 以及它是如何使用的? (参见我上面的示例程序)
    【解决方案2】:

    other solution you posted相比,您只需要(首先)关闭频道,并继续向其喂食。

    那么您需要回答以下问题:(a) 是否绝对有必要从队列中获取下一个 X 项“无事可做”(或者,相同的是,一旦前 X 个项目被完全处理或分配给一个工人);或者 (b) 是否可以将第二组 X 项保留在内存中,并在需要新工作项时将它们提供给工人?

    据我了解,只有 (a) 需要您想知道的 requestChan(见下文)。对于 (b),如下所示的简单内容就足够了:

    # B version
    
    type WorkItem int
    
    const (
      N = 5  // Number of workers
      X = 15 // Number of work items to get from the infinite queue at once
    )
    
    func Worker(id int, workChan <-chan WorkItem) {
      for {
        item := <-workChan
        doWork(item)
        fmt.Printf("Worker %d processes item #%v\n", id, item)
      }
    }
    
    func Dispatch(workChan chan<- WorkItem) {
      for {
        items := GetNextFromQueue(X)
    
        for _, item := range items {
          workChan <- item
          fmt.Printf("Dispatched item #%v\n", item)
        }
      }
    }
    
    func main() {
      workChan := make(chan WorkItem) // Shared amongst all workers; could make it buffered if GetNextFromQueue() is slow.
    
      // Start N workers.
      for i := 0; i < N; i++ {
        go Worker(i, workChan)
      }
    
      // Dispatch items to the workers.
      go Dispatch(workChan)
    
      time.Sleep(20 * time.Second) // Ensure main(), and our program, finish.
    }
    

    (我已将 full working solution for (b) 上传到 Playground。)

    至于(a),工人改为说:做工作,如果没有更多工作,告诉调度员通过获得更多reqChan 沟通渠道。 “或”是通过select实现的。然后,调度程序等待 reqChan,然后再次调用GetNextFromQueue()。它的代码更多,但确保了您可能感兴趣的语义。(不过,以前的版本总体上更简单。)

    # A version
    
    func Worker(id int, workChan <-chan WorkItem, reqChan chan<- int) {
      for {
        select {
        case item := <-workChan:
          doWork(item)
          fmt.Printf("Worker %d processes item #%v\n", id, item)
        case reqChan <- id:
          fmt.Printf("Worker %d thinks they requested more work\n", id)
        }
      }
    }
    
    func Dispatch(workChan chan<- WorkItem, reqChan <-chan int) {
      for {
        items := GetNextFromQueue(X)
    
        for _, item := range items {
          workChan <- item
          fmt.Printf("Dispatched item #%v\n", item)
        }
    
        id := <-reqChan
        fmt.Printf("Polling the queue in Dispatch() at the request of worker %d\n", id)
      }
    }
    

    (我还向 Playground 上传了 full working solution for (a)。)

    【讨论】:

    • 再想一想,第二个版本也可能完全被划掉:它几乎没有添加任何内容,因为 select 不能保证按自上而下的顺序成功:“如果多个案例可以继续,a统一的伪随机选择来决定将执行哪个单一通信。” (link) 因此,一旦调度员等待reqChan,其中一名工作人员就可以写信给它。 — 是否有一些更简单的版本不足以满足您的需求?
    猜你喜欢
    • 2018-09-27
    • 2013-06-23
    • 1970-01-01
    • 1970-01-01
    • 2016-03-15
    • 1970-01-01
    • 1970-01-01
    • 2011-05-11
    • 1970-01-01
    相关资源
    最近更新 更多