【问题标题】:Go Nonblocking multiple receive on channel在通道上进行非阻塞多个接收
【发布时间】:2014-10-09 02:26:43
【问题描述】:

似乎到处都在讨论从通道读取应该始终是阻塞操作。态度似乎是这就是Go的方式。这是有道理的,但我正试图弄清楚如何从渠道中聚合内容。

例如,发送 http 请求。假设我有一个生成数据流的管道设置,所以我有一个生成队列/点流的通道。然后我可以让一个 goroutine 监听这个通道并发送一个 HTTP 请求来将它存储在一个服务中。这可行,但我正在为每个点创建一个 http 请求。

我发送的端点也允许我批量发送多个数据点。我想做的,是

  1. 读取尽可能多的值,直到我阻塞通道。
  2. 组合它们/发送单个 http 请求。
  3. 然后阻塞频道直到我可以阅读 再来一张。

这就是我在 C 语言中使用线程安全队列和选择语句的方式。基本上在可能的情况下刷新整个/队列缓冲区。这是一种有效的围棋技术吗?

似乎 go select 语句确实给了我类似于 C 的 select 的东西,但我仍然不确定通道上是否存在“非阻塞读取”。

编辑:我也愿意接受我想要的可能不是 Go Way,但不断粉碎不间断的 http 请求对我来说似乎也是错误的,特别是如果它们可以聚合的话。如果有人有一个很酷的替代架构,但我想避免诸如神奇地缓冲 N 个项目或等待 X 秒直到发送之类的事情。

【问题讨论】:

  • 虽然它并没有真正超时。如果没有,我想阻止,然后在可能的情况下阅读多个。如果 multiple 由 1 组成,那很好。如果有 100 件事情,让我把它们都拿走并处理它们。这样,http 帖子就会适应负载。

标签: select go nonblocking


【解决方案1】:

这是批处理直到通道为空的方法。变量batch 是数据点类型的一部分。变量ch 是您的数据点类型的通道。

var batch []someType
for {
    select {
    case v := <-ch:
       batch = append(batch, v)
    default:
       if len(batch) > 0 {
           sendBatch(batch)
           batch := batch[:0]
       }
       batch = append(batch, <-ch)  // receiving a value here prevents busy waiting.
    }
}

您应该防止批次无限制地增长。这是一个简单的方法:

var batch []someType
for {
    select {
    case v := <-ch:
       batch = append(batch, v)
       if len(batch) >= batchLimit {
           sendBatch(batch)
           batch := batch[:0]
       }
    default:
       if len(batch) > 0 {
           sendBatch(batch)
           batch := batch[:0]
       }
       batch = append(batch, <-ch)
    }
}

【讨论】:

  • 我认为这可能类似于尝试阅读我在选择中等待的内容:)。这是一个很好的例子。读起来确实很奇怪:)
  • 如果你不想自己实现,这个模式也在我的频道包中的BatchingChannel中实现:godoc.org/github.com/eapache/channels#BatchingChannel
【解决方案2】:

Dewy Broto 为您的问题提供了一个很好的解决方案。这是一个直截了当的直接解决方案,但我想更广泛地评论一下您如何为不同的问题寻找解决方案。

Go 使用通信顺序进程代数 (CSP) 作为通道、选择和轻量级进程(“goroutines”)的基础。 CSP 保证事件的顺序;只有当您通过选择(又名select)来实现它时,它才会引入非确定性。有保证的排序有时被称为“先发生”——它使编码比替代的(广泛流行的)非阻塞样式简单得多。它还为创建组件提供了更多空间:通过渠道以可预测的方式与外部世界交互的长期功能单元。

也许在频道上讨论阻塞会给人们学习 Go 的方式带来心理障碍。我们在 I/O 上阻塞,但我们在通道上等待等待通道是不应该的,只要系统作为一个整体有足够的并行松弛(即其他活动的 goroutines)来保持 CPU 忙碌。

可视化组件

那么,回到你的问题。让我们从组件的角度来考虑它,您有许多需要探索的点来源。假设每个源都是一个 goroutine,然后它在您的设计中形成一个带有输出通道的组件。 Go 允许共享通道端,因此许多源可以安全地将它们的点按顺序交错到单个通道上。您无需执行任何操作 - 这就是渠道的运作方式。

Dewy Broto 描述的批处理功能本质上是另一个组件。作为一个学习练习,用这种方式表达它是一件好事。批处理组件有一个点输入通道和一个批处理输出通道。

最后,HTTP i/o 行为也可以是一个只有一个输入通道而没有输出通道的组件,仅用于接收整批点然后通过 HTTP 发送它们。

以只有一个来源的简单情况为例,可以这样描述:

+--------+     point     +---------+     batch     +-------------+
| source +------->-------+ batcher +------->-------+ http output |
+--------+               +---------+               +-------------+

这里的目的是描述不同的活动在它们的基本层面。这有点像数字电路图,这不是巧合。

你确实可以在 Go 中实现它并且它会起作用。它甚至可能工作得很好,但实际上您可能更喜欢通过组合成对的组件来优化它,必要时重复。在这种情况下,很容易将批处理程序和 http 输出结合起来,最终得到了 Dewy Broto 的解决方案。

重要的一点是Go并发最容易发生

  • (a) 不用担心阻塞;
  • (b) 描述了需要在相当细粒度的级别上发生的活动(在简单的情况下,您可以在脑海中执行此操作);
  • (c) 如有必要,通过组合功能进行优化。

我将把可视化移动通道末端(Pi-Calculus)的更高级主题作为挑战,其中通道用于将通道末端发送到其他 goroutine。

【讨论】:

  • 我的系统实际上是这样构建的。我有大约 7 个“管道组件”,每个组件都推到下一个。设计它很有趣。它或多或少地从不同的系统收集数据点,并将每个点转换为新的东西,然后再将其存储到最终系统中。我什至有一个单点 http 输出组件的工作实现,但是在 1 秒内抛出 3 天的数据,而且它很慢。 259200 http 连接相当昂贵,因此想要找出“组合”进入通道的事物的方式,以便对其采取行动:)。
  • 写得很棒。这对于其他解决这个问题的人来说非常有用。以“GO 方式”做事时要记住的好事情
猜你喜欢
  • 2020-06-20
  • 2016-07-16
  • 2020-01-16
  • 1970-01-01
  • 1970-01-01
  • 2012-08-24
  • 1970-01-01
  • 1970-01-01
  • 2019-07-10
相关资源
最近更新 更多