【问题标题】:Golang: consume items for multiple workers from channel in `case` statementGolang:在“case”语句中为多个工作人员消费项目
【发布时间】:2021-12-14 22:45:57
【问题描述】:

我的消费者(从main 运行)支持上下文取消和通过case 语句从通道读取。我可以使用上下文关闭消费者,效果很好。但是,当我在一个案例语句中生成多个工人时,每个工人都从jobsChan 获得相同的工作(消息),这不是我想要的:

func (app *App) consumer() {
    for {
        select {
        case <-app.ctx.Done():
            app.infoLog.Print("Caught SIGINT, stopping.")
            app.wg.Wait()
            app.doneChan <- struct{}{} # main uses this channel to block itself until all goroutines are stopped
            app.infoLog.Print("Shutting down the consumer...")
            return
        case job := <-app.jobsChan:
            // PROBLEM here: wrong, each worker is given the same job
            for workerNumber := 0; workerNumber < app.config.workers; workerNumber++ {
                app.wg.Add(1)
                go app.workerFunc(workerNumber, job)
            }
        }
    }
}

func (app *App) workerFunc(id int, job Job) {
    defer app.wg.Done()
    
    ... actual worker code here ...
}

如何重写此代码,以便我可以为app.ctx.Done 频道保留select,同时可以生成工人,以便每个工人从频道中选择下一条消息作为工作?我需要保留for/select 以收听ctx 的取消,但同时我需要在消费者中生成X 工作人员读取来自jobsChan 的消息。这可能吗?

想到的唯一选择是将通道直接传递到生成的workerFunc,并在workerFunc 中添加另一个for job := range app.jobsChan。但随后消费者中的整个case job := &lt;-app.jobsChan: 变得毫无意义,我不知道如何重写它。

澄清一下:当我运行应用程序时,我希望每个工作人员都有一个从 jobsChan 中提取的新工作 ID - 但它们的处理方式都相同,例如1,然后他们都处理下一个,例如2

#wrong
Worker 0: start processing item 1
Worker 2: start processing item 1
Worker 1: start processing item 1

【问题讨论】:

  • 确保我关注,你想为你完成的每个工作创建一个app.config.workers 数量的app.workerFunc goroutines?如果你只拉一个job 为什么要生成多个 goroutine?或者您是否尝试使用固定数量的workers 来处理作业?展示您的问题的示例会有所帮助。
  • 是的@sberry,我的频道中有几十个项目,我想在几个同时运行的 goroutine 中处理。在我目前的设置中,我无法做到这一点,因为所有 goroutine 都被赋予了相同的工作。我在这里学习教程rodrigoaraujo.me/posts/…
  • 如果您删除 for 循环。如果没有通过通道多次发送作业,则每个工人应该得到不同的作业。 for 循环的目的尚不清楚。也许你会发现这个包很有用:github.com/MicahParks/ctxerrpool
  • @MicahParks 没有帮助。仍在使用相同的作业调用工人。作业被添加到处理程序的队列中。图像正在通过 Postman 上传,一旦图像到达,它就会被添加到缓冲通道中。所以我看不出它会如何多次进入队列(频道)。

标签: go


【解决方案1】:

您现有的代码明确地将相同的工作分配给所有工作人员。如果您有固定数量的工作人员,请为他们创建 goroutine(在初始化期间),并让他们收听频道:

for workerNumber:0;workerNumber<app.config.workers;workerNumber++ {
   go app.workerFunc(ctx,workerNumber,app.jobsChan)
}

在每个worker中,只需检查jobQueue和上下文取消。

换句话说,您不需要consumer,直接将工作传递给工人。

【讨论】:

  • 谢谢@Burak Serdar。这个解决方案奏效了。当服务器完成时,我还需要关闭main 中的通道,并且我还从consumer 中删除了所有for, select, case 语句,以便只阻止&lt;-app.ctx.Done() 和之后的命令。
  • 这听起来不对。我的意思是你不需要consumer 函数。您所需要的只是工作人员和他们从中阅读的公共频道。
猜你喜欢
  • 2014-06-05
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-11-11
  • 1970-01-01
  • 2018-11-15
相关资源
最近更新 更多