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