【发布时间】:2018-02-02 23:24:17
【问题描述】:
我正在在线阅读管道教程并尝试构建一个像这样运行的阶段--
- 将传入事件以每批 10 个为一组,然后将它们发送到输出通道
- 如果我们在 5 秒内没有看到 10 个事件,则合并我们收到的所有事件并发送它们,关闭 out chan 并返回。
但是,我不知道第一个选择案例会是什么样子。尝试了多种方法,但无法通过这个。 任何指针都非常感谢!
func BatchEvents(inChan <- chan *Event) <- chan *Event {
batchSize := 10
comboEvent := Event{}
go func() {
defer close(out)
i = 0
for event := range inChan {
select {
case -WHAT GOES HERE?-:
if i < batchSize {
comboEvent.data = append(comboEvent.data, event.data)
i++;
} else {
out <- &comboEvent
// reset for next batch
comboEvent = Event{}
i=0;
}
case <-time.After(5 * time.Second):
// process whatever we have seen so far if the batch size isn't filled in 5 secs
out <- &comboEvent
// stop after
return
}
}
}()
return out
}
【问题讨论】:
标签: go concurrency pipeline channel