【问题标题】:Coalescing items in channel合并通道中的项目
【发布时间】: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 &lt;- task }。当然,它周围有一个互斥锁。
  • @Ainar-G 我也想得到处理结果,所以你的代码应该有 else 分支

标签: multithreading go queue


【解决方案1】:

如果我正确理解了这个问题,我想到的最简单的解决方案是在任务发送者(放入queue)和工作人员(取自queue)之间添加一个中间层。这可能是例行公事,负责存储当前任务(按 ID)并将结果广播到每个匹配的任务。

伪代码:

go func() {
    active := make(map[TaskID][]Task)

    for {
        select {
        case task := <-queue:
            tasks := active[task.ID]
            // No tasks with such ID, start heavy work
            if len(tasks) == 0 {
                worker <- task
            }
            // Save task for the result
            active[task.ID] = append(active[task.ID], task)
        case r := <-response:
            // Broadcast to all tasks
            for _, task := range active[r.ID] {
                task.Result <- r.Result
            }
        }
    }
}()

不需要互斥体,也可能不需要携带任何东西,工作人员只需将所有结果放入这个中间层,然后正确路由响应。如果冲突的 ID 有可能在一段时间后到达,您甚至可以轻松地在此处添加缓存。

编辑:我做了一个梦,上面的代码导致了死锁。如果您一次发送大量请求并阻塞 worker 通道,则会出现严重问题 - 这个中间层例程卡在 worker &lt;- task 等待工作人员完成,但所有工作人员可能会在发送到响应通道时被阻塞(因为我们的例程无法收集它)。 Playable proof.

可以考虑在通道中添加一些缓冲区,但这不是一个合适的解决方案(除非您可以设计系统以使缓冲区永远不会填满)。有几种方法可以解决这个问题;例如,您可以运行一个单独的例程来收集响应,但是您需要使用互斥锁来保护active 映射。可行的。您还可以将worker &lt;- task 放入一个选择中,该选择将尝试向工作人员发送任务、接收新任务(如果没有要发送)或收集响应。可以利用 nil 通道从未准备好进行通信(被选择忽略)这一事实,因此您可以在单个选择中交替接收和发送任务。示例:

go func() {
    var next Task // received task which needs to be passed to a worker
    in := queue // incoming channel (new tasks) -- active
    var out chan Task // outgoing channel (to workers) -- inactive
    for {
        select {
        case t := <-in:
            next = t // store task, so we can pass to worker
            in, out = nil, worker // deactivate incoming channel, activate outgoing
        case out <- next:
            in, out = queue, nil // deactivate outgoing channel, activate incoming
        case r := <-response:
            collect <- r
        }
    }
}()

play

【讨论】:

  • 乐于助人。我不知道您的流程到底是什么样子,但已经意识到您可以通过这样一个幼稚的示例轻松陷入僵局。我已经编辑了我的答案以进一步讨论这个问题。
猜你喜欢
  • 2011-02-06
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-04-05
  • 2019-02-27
  • 2019-02-22
  • 1970-01-01
相关资源
最近更新 更多