【问题标题】:What is the proper way to read from the errors channel in Sarama?从 Sarama 的错误通道中读取的正确方法是什么?
【发布时间】:2018-02-01 15:07:55
【问题描述】:

当我生成消息时,我正在使用用 Go 编写的 Sarama 库从错误通道中读取。整个代码如下所示,包含在一个函数中:

producer.AsyncProducer.Input() <- &sarama.ProducerMessage{Topic: topic, Key: nil, Value: sarama.ByteEncoder(message)}
go func() {
    for err := range saramaProducer.Errors() {
        if producer.callbacks.OnError != nil {
            producer.callbacks.OnError(err)
        }
    }
}()

根据我对 goroutine 的理解,我的 goroutine 会不断迭代 Errors() 频道,直到它收到一个。有没有办法让它在我的函数执行完成后停止侦听错误?

【问题讨论】:

  • 如果 api 希望你在一个频道上进行范围,它应该在完成后关闭频道。

标签: go goroutine sarama


【解决方案1】:

您可以使用另一个频道和select 使循环返回。

var quit chan struct{}
go func() {
    for {
        select {
        case err:=<-saramaProducer.Errors():
            //handle errors
        case <-quit:
            return
        }
    }
}
defer func() { quit<-struct{}{} }()

原始的for ... range 循环在得到一个之前不会迭代通道。相反,它会一直阻塞,直到出现错误,处理它,然后再次等待新的错误,直到通道关闭或main 返回。

上面的代码有个小问题,当quit和error channel都准备好时,select会随机选择一个,可能会导致单个error丢失。如果这值得处理,只需将另一个 switchdefault 放在一起以获取该错误,然后再添加 return

【讨论】:

  • 你不需要处理这种情况,当生产者关闭时,包会关闭通道(应该)。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2010-10-15
  • 2021-08-25
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多