【发布时间】:2015-05-19 14:56:54
【问题描述】:
我有一个接收任务并将它们放入通道的函数。每个任务都有 ID、一些属性和一个放置结果的通道。看起来是这样的
task.Result = make(chan *TaskResult)
queue <- task
result := <-task.Result
sendReponse(result)
另一个 goroutine 从通道中获取一个任务,处理它并将结果放入任务的通道中
task := <-queue
task.Result <- doExpensiveComputation(task)
此代码运行良好。但现在我想在queue 中合并任务。任务处理是一项非常昂贵的操作,因此我希望一次处理队列中具有相同 ID 的所有任务。我看到了两种方法。
第一个是不要将具有相同 ID 的任务放入队列中,因此当现有任务到达时,它将等待其副本完成。这是伪代码
if newTask in queue {
existing := queue.getById(newTask.ID)
existing.waitForComplete()
sendResponse(existing.ProcessingResult)
} else {
queue.enqueue(newTask)
}
所以,我可以使用 go channel 和 map 来实现它以进行随机访问 + 一些同步方式,如互斥锁。我不喜欢这种方式的是我必须在代码周围同时携带地图和频道并保持它们的内容同步。
第二种方式是将所有任务放入队列,但是当结果到达时从队列中提取任务和所有具有相同ID的任务,然后将结果发送给所有任务。这是伪代码
someTask := queue.dequeue()
result := doExpensiveComputation(someTask)
someTask.Result <- result
moreTasks := queue.getAllWithID(someTask.ID)
for _,theSameTask := range moreTasks {
theSameTask.Result <- result
}
我知道如何使用 chan + map + mutex 以与上述相同的方式实现这一点。
问题来了:是否有一些内置/现有的数据结构可用于解决此类问题?还有其他(更好的)方法吗?
【问题讨论】:
-
如果你只想用一个唯一ID处理每个任务一次,那为什么不在发送端检查唯一性呢? IE。
if _, sent := alreadySentIDs[ID]; !sent { queue <- task }。当然,它周围有一个互斥锁。 -
@Ainar-G 我也想得到处理结果,所以你的代码应该有 else 分支
标签: multithreading go queue