【问题标题】:Unable to consume messages from locally running Kafka server, using Golang Sarama Package无法使用 Golang Sarama 包从本地运行的 Kafka 服务器消费消息
【发布时间】:2017-11-02 09:07:03
【问题描述】:

我正在制作一个简单的 Telegram 机器人,它会从本地 Kafka 服务器读取消息并将其打印到聊天中。 zookeeper 和 kafka 服务器配置文件都是默认的。控制台消费者工作。当我尝试使用 Golang Sarama 包从代码中使用消息时,问题就出现了。在我添加这些行之前:

case err := <-pc.Errors(): log.Panic(err)

程序只打印一次消息,之后它就会停止。 现在它会将它打印到日志中: kafka: error while consuming test1/0: kafka: broker not connected

代码如下:

    type kafkaResponse struct {
        telega  *tgbotapi.Message
        message []byte
    }

    type kafkaRequest struct {
        telega *tgbotapi.Message
        topic  string
    }    
    var kafkaBrokers = []string{"localhost:9092"}
    func main() {
                //channels for request response
                var reqChan = make(chan kafkaRequest)
                var respChan = make(chan kafkaResponse)

                //starting kafka client routine to listen to topic channnel
                go consumer(reqChan, respChan, kafkaBrokers)

                //bot thingy here
                bot, err := tgbotapi.NewBotAPI(token)
                if err != nil {
                    log.Panic(err)
                }
                bot.Debug = true
                log.Printf("Authorized on account %s", bot.Self.UserName)
                u := tgbotapi.NewUpdate(0)
                u.Timeout = 60
                updates, err := bot.GetUpdatesChan(u)
                for {
                    select {
                    case update := <-updates:
                        if update.Message == nil {
                            continue
                        }
                        switch update.Message.Text {

                        case "Topic: test1":
                            topic := "test1"
                            reqChan <- kafkaRequest{update.Message, topic}
                        }
                    case response := <-respChan:
                        bot.Send(tgbotapi.NewMessage(response.telega.Chat.ID, string(response.message)))
                    }

                }

这里是消费者。去:

 func consumer(reqChan chan kafkaRequest, respChan chan kafkaResponse, brokers []string) {
            config := sarama.NewConfig()
            config.Consumer.Return.Errors = true

            // Create new consumer
            consumer, err := sarama.NewConsumer(brokers, config)
            if err != nil {
                panic(err)
            }
            defer func() {
                if err := consumer.Close(); err != nil {
                    panic(err)
                }
            }()

            select {
            case request := <-reqChan:
                //get all partitions on the given topic
                partitionList, err := consumer.Partitions(request.topic)
                if err != nil {
                    fmt.Println("Error retrieving partitionList ", err)
                }

                initialOffset := sarama.OffsetOldest
                for _, partition := range partitionList {
                    pc, _ := consumer.ConsumePartition(request.topic, partition, initialOffset)

                    go func(pc sarama.PartitionConsumer) {
                        for {
                            select {
                            case message := <-pc.Messages():
                                respChan <- kafkaResponse{request.telega, message.Value}
                            case err := <-pc.Errors():
                                log.Panic(err)
                            }
                        }
                    }(pc)
                }
            }
        }

【问题讨论】:

  • 错误信息非常清楚,您的代理(kafka 服务器)无法从您的客户端(无论它在哪里)访问
  • 问题是客户端能够获取消息,但只能获取一次

标签: go apache-kafka telegram-bot sarama


【解决方案1】:

在代码中设置完所有PartitionConsumers 后,您将关闭您的消费者

defer func() {
            if err := consumer.Close(); err != nil {
                panic(err)
            }
        }()

但是,文档指定您应该仅在所有 PartitionConsumer 都已关闭后才关闭使用者。

// Close shuts down the consumer. It must be called after all child
// PartitionConsumers have already been closed.
Close() error

我建议你在函数go func(pc sarama.PartitionConsumer) { 中添加一个sync.WaitGroup

【讨论】:

  • 这似乎不是问题,因为我删除了 defer 功能(只是为了看看它是否是),但程序仍然停止。
  • 它是否仍然给出问题中提到的错误消息kafka: error while consuming test1/0: kafka: broker not connected?
猜你喜欢
  • 2018-10-31
  • 2018-08-29
  • 1970-01-01
  • 2019-09-23
  • 2021-09-02
  • 2018-11-29
  • 2016-12-07
  • 1970-01-01
相关资源
最近更新 更多