【问题标题】:How to consume from Kafka topic in multiple goroutines, using Sarama?如何使用 Sarama 从多个 goroutine 中的 Kafka 主题中消费?
【发布时间】:2019-07-05 03:02:50
【问题描述】:

我使用https://github.com/Shopify/sarama 与 Kafka 进行交互。我有一个主题,例如 100 个分区。我有应用程序,它部署在 1 台主机 上。所以,我想在多个 goroutines 中使用这个主题。

我看到了这个例子 - https://github.com/Shopify/sarama/blob/master/examples/consumergroup/main.go ,我们可以在其中看到如何在特定的消费者组中创建消费者。

所以,我的问题是,我应该创建多个这样的消费者,还是在 Sarama 中有一些设置,我可以在其中设置所需数量的消费者 goroutine。

附:我看到了这个问题 - https://github.com/Shopify/sarama/issues/140 - 但没有答案,如何创建 MultiConsumer。

【问题讨论】:

  • 我认为您可以使用 100 个消费者在属于同一消费者组的 100 个 goroutine 中运行。但是,如果您从单个主机执行此操作,则使用单个消费者进行消费可能不会有太大的加速。是的,如果你想使用 N 个消费者,你需要设置 N 个消费者。

标签: go apache-kafka kafka-consumer-api goroutine sarama


【解决方案1】:

这个例子展示了一个完全工作的控制台应用程序,它可以消费一个主题中的所有分区,为每个分区创建一个 goroutine:

https://github.com/Shopify/sarama/blob/master/tools/kafka-console-consumer/kafka-console-consumer.go

链接在您在问题中发布的主题的末尾。

它基本上创建了一个消费者:

c, err := sarama.NewConsumer(strings.Split(*brokerList, ","), config)

然后获取所需主题的所有分区:

func getPartitions(c sarama.Consumer) ([]int32, error) {
    if *partitions == "all" {
        return c.Partitions(*topic)
    }
...

然后为每个分区创建一个 PartitionConsumer 并在不同的 goroutine 中从每个分区消费:

for _, partition := range partitionList {
    pc, err := c.ConsumePartition(*topic, partition, initialOffset)
    ....

    wg.Add(1)
    go func(pc sarama.PartitionConsumer) {
        defer wg.Done()
        for message := range pc.Messages() {
            messages <- message
        }
    }(pc)

}

【讨论】:

  • 谢谢,这个例子对应用程序的 1 个节点非常有效。有意思,consumer group API 有这样的东西吗?
  • 当您需要将消费者拆分到不同的进程中时,通常会使用消费者组(相同的组名表示所有这些单独的进程实际上是同一组的一部分,而不是独立的消费者)。在这种情况下,您可能希望按照以下方式实现一些东西:godoc.org/github.com/Shopify/sarama#ConsumerGroup,然后在不同的进程/主机中运行相同的代码。
猜你喜欢
  • 2020-03-19
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-09-26
  • 2017-11-24
  • 1970-01-01
相关资源
最近更新 更多