【问题标题】:BrokerNotAvailableError: Could not find the leader Exception while Spark StreamingBrokerNotAvailableError:Spark Streaming 时找不到领导者异常
【发布时间】:2015-07-23 13:13:51
【问题描述】:

我用 NodeJS 编写了一个 Kafka Producer,用 Java Maven 编写了一个 Kafka Consumer。我的主题是由以下命令创建的“测试”:

bin/kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 1 --topic test

NodeJS 中的生产者:

var kafka = require('kafka-node');
var Producer = kafka.Producer;
var Client = kafka.Client;
var client = new Client('localhost:2181');
var producer = new Producer(client);

producer.on('ready', function () {
    producer.send([
        { topic: 'test', partition: 0, messages: ["This is the zero message I am sending from Kafka to Spark"], attributes: 0},
        { topic: 'test', partition: 1, messages: ["This is the first message I am sending from Kafka to Spark"], attributes: 0},
        { topic: 'test', partition: 2, messages: ["This is the second message I am sending from Kafka to Spark"], attributes: 0}
        ], function (err, result) {
        console.log(err || result);
        process.exit();
    });
});

当我从 NodeJS 生产者发送两条消息时,它成功地被 Java 消费者消费了。但是当我从 NodeJS 生产者发送三个或更多消息时,它给了我以下错误:

{ [BrokerNotAvailableError: 找不到领导者] 消息:'找不到领导者' }

我想问一下,如何将 LEADER 设置为主题“测试”中的任何消息。或者该问题的解决方案是什么。

【问题讨论】:

  • 为了获得更高的可靠性,您可以运行更多数量的代理并在主题上启用复制,这样如果领导者代理停止,其他跟随者代理将占据位置,并且您将在领导者不可用的情况下运行更少..

标签: apache-kafka


【解决方案1】:

使用partitions(复数键名)代替partition

例如:

producer.on('ready', function () {
  producer.send([
    { topic: 'test', partitions: 0, messages: ["This is the zero message I am sending from Kafka to Spark"], attributes: 0},
    { topic: 'test', partitions: 1, messages: ["This is the first message I am sending from Kafka to Spark"], attributes: 0},
    { topic: 'test', partitions: 2, messages: ["This is the second message I am sending from Kafka to Spark"], attributes: 0}
  ], function (err, result) {
    console.log(err || result);
    process.exit();
  });
});

【讨论】:

  • 这解决了我的问题,但您能告诉我们为什么将partition 更改为partitions 吗?
  • 2020 年,情况似乎不再如此。 partition(单数)工作正常。
【解决方案2】:

主题是用 1 个分区创建的,但是在生产者端,您尝试将消息发送到 3 个分区,从逻辑上讲,Kafka 不应该为其他分区找到领导者,应该抛出这个异常。

【讨论】:

    【解决方案3】:

    在当前版本的kafka-node 中存在一个可能导致此问题的错误

    https://github.com/SOHU-Co/kafka-node/issues/354

    带有 KeyedPartitioner 的 HighLevelProducer 在第一次发送时失败 #354 将 KeyedParitioner 与 HighLevelProducer 一起使用时,第一次发送失败 BrokerNotAvailableError:找不到领导者 连续发送完美。

    另见https://www.npmjs.com/package/kafka-node#highlevelproducer-with-keyedpartitioner-errors-on-first-send

    推荐

    在发送第一条消息之前调用 client.refreshMetadata()。

    我就是这样做的

        // Refresh metadata required for the first message to go through
        // https://github.com/SOHU-Co/kafka-node/pull/378
        client.refreshMetadata([topic], (err) => {
            if (err) {
                console.warn('Error refreshing kafka metadata', err);
            }
        });
    

    【讨论】:

      【解决方案4】:

      我遇到了类似的问题。当我推送消息或读取分区不存在的消息时。请使用以下命令手动修改分区计数为您需要的数量。

      > sh kafka-topics.sh --alter --bootstrap-server localhost:9092 --partitions 2 --topic test
      

      在此之后它会正常工作。

      “Miguel Veces”的回答是有效的,但它是错误的,因为随着该更改它开始忽略您传递的分区参数。

      【讨论】:

        【解决方案5】:

        终于解决了。我犯了一个错误。默认情况下 Apache Kafka 创建 3 个复制,这意味着默认情况下它创建 3 个代理。然而,上面 YAML 只创建一个代理,并寻找其他 2 个未创建的代理。因此我们得到了同样的错误。

        修复

        version: '3'
        services:
          zookeeper:
            container_name: zookeeper
            image: confluentinc/cp-zookeeper
            ports:
              - "32181:32181"
            environment:
              ZOOKEEPER_CLIENT_PORT: 32181
              ZOOKEEPER_TICK_TIME: 2000
              ZOOKEEPER_SYNC_LIMIT: 2
            
          kafka:
            container_name: kafka
            image: confluentinc/cp-kafka
            ports:
              - "9094:9094"
            environment:
              KAFKA_ZOOKEEPER_CONNECT: zookeeper:32181
              KAFKA_LISTENERS: INTERNAL://:9092,OUTSIDE://:9094
              KAFKA_ADVERTISED_LISTENERS: INTERNAL://:9092,OUTSIDE://localhost:9094
              KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,OUTSIDE:PLAINTEXT
              KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL
              KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
              ES_JAVA_OPTS: "-Xms512m -Xmx3000m"
        

        然后KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 告诉 Apache kafka 我们只想创建一个复制及其工作。

        使用以下代码发布和消费kafka消息

        var intCounter = 104;
        
            setInterval(() => {
                setMessageInTopic('topic'+intCounter, function(response){
                    echo("Producer Response"+response);
                    if(response){
                        echo("Consumer init");
                        getMessageFromKafak(response,function(pStrMessage){
                            echo ("------------------------Consumer message --------------------");
                            echo(pStrMessage);
                        });
                    }else{
                        echo("Error");
                    }
                });
                intCounter++; 
            },5000);
            
            
            
        function setMessageInTopic(pStrTopicName, callback){
            const kafka = require('kafka-node');
        
            try {
              const Producer = kafka.Producer;
              const client = new kafka.KafkaClient({kafkaHost:"localhost:9094"});
              const producer = new Producer(client);
              const kafka_topic =  pStrTopicName;
              console.log(kafka_topic);
              let payloads = [
                {
                  topic: kafka_topic,
                  messages: {'name':'vipin'}
                }
              ];
        
              producer.on('ready', async function() {
                let push_status = producer.send(payloads, (err, data) => {
                  if (err) {
                    console.log('[kafka-producer -> '+kafka_topic+']: broker update failed');
                    callback(false);
                  } else {
                    console.log('[kafka-producer -> '+kafka_topic+']: broker update success');
                    callback(pStrTopicName);
                  }
                });
              });
        
              producer.on('error', function(err) {
                console.log(err);
                console.log('[kafka-producer -> '+kafka_topic+']: connection errored');
                callback(false);
                ///throw err;
              });
            }
            catch(e) {
              console.log(e);
              callback(false);
            }
        }
        
        function getMessageFromKafak(pStrTopicName, callback){
            try{
                var kafka = require('kafka-node');
                //var HighLevelProducer = kafka.HighLevelProducer;
                var Consumer = kafka.Consumer;
                var client = new kafka.KafkaClient({kafkaHost:"localhost:9094"});
                
                let consumer = new Consumer(
                    client,
                    [{ topic: pStrTopicName}],
                    {
                      autoCommit: true,
                      fetchMaxWaitMs: 1000,
                      fetchMaxBytes: 1024 * 1024,
                      encoding: 'utf8',
                      fromOffset: false
                    }
                  );
              
                echo("------------------RESPONE MESSAGE PROCESS "+pStrTopicName+" ----------------------")
                consumer.on('message', function (message) {
                    echo("------------------RESPONE MESSAGE  ----------------------")
                    console.log(message);
                    consumer.close();
                    client.close();
                    return callback(message);
                }).on('error', function (message) {
                    echo("------------------RESPONE MESSAGE ERROR ----------------------")
                    console.log(message);
                    consumer.close();
                    client.close();
                    return callback(message);
                });
            }catch(e) {
              console.log(e);
            }
        }
        

        编码愉快:)

        感谢和问候 Jaiswar Vipin Kumar R.

        【讨论】:

          猜你喜欢
          • 2016-03-21
          • 2017-01-09
          • 2020-07-26
          • 2015-10-25
          • 1970-01-01
          • 1970-01-01
          • 2017-07-12
          • 1970-01-01
          • 2017-06-07
          相关资源
          最近更新 更多