【问题标题】:ERROR Error when sending message to topicERROR 向主题发送消息时出错
【发布时间】:2016-11-22 07:11:48
【问题描述】:

在 kafka 中生成消息时,我收到以下错误:

$ bin/kafka-console-producer.sh --broker-list localhost:9092 --topic nil_PF1_P1
hi
hello

[2016-07-19 17:06:34,542] ERROR Error when sending message to topic nil_PF1_P1 with key: null, value: 2 bytes with error: (org.apache.kafka.clients.producer.internals.ErrorLoggingCallback)
org.apache.kafka.common.errors.TimeoutException: Failed to update metadata after 60000 ms.
[2016-07-19 17:07:34,544] ERROR Error when sending message to topic nil_PF1_P1 with key: null, value: 5 bytes with error: (org.apache.kafka.clients.producer.internals.ErrorLoggingCallback)
org.apache.kafka.common.errors.TimeoutException: Failed to update metadata after 60000 ms.


$ bin/kafka-topics.sh --describe --zookeeper localhost:2181 --topic nil_PF1_P1
Topic:nil_PF1_P1    PartitionCount:1    ReplicationFactor:1 Configs:
Topic: nil_PF1_P1   Partition: 0    Leader: 2   Replicas: 2 Isr: 2

对此有什么想法吗??

【问题讨论】:

  • 你能把下面命令的结果贴出来bin/kafka-topics.sh --describe --zookeeper localhost:2181 --topic nil_PF1_P1
  • 改版后转贴!!
  • 从 zk 端看来一切都很好。 Kafka代理绑定地址和端口似乎有问题!你能在 localhost:9092 检查你的 kafka 是否可以访问吗?
  • 顺便说一句,您使用的是什么版本的 Kafka?如果是 0.9 或更高版本,您是否使用 ssl 对其进行了配置?如果那么您需要在生产和消费时提供 ssl 信息!
  • 是的,它是 0.10.0.0 。什么是 ssl?可以简单解释一下。链接也行!!提前致谢!

标签: apache-kafka


【解决方案1】:

不要更改server.properties,而是在代码本身中包含地址0.0.0.0。 而不是

/usr/bin/kafka-console-producer --broker-list Hostname:9092 --topic MyFirstTopic1

使用

/usr/bin/kafka-console-producer --broker-list 0.0.0.0:9092 --topic MyFirstTopic1

【讨论】:

    【解决方案2】:

    可能是因为 Kafka 的 server.properties 文件中的一些参数。你可以找到更多信息here

    1. 停止 Kafka 服务器

      cd $KAFKA_HOME/bin  
      ./kafka-server-stop.sh
      
    2. 改变

      listeners=PLAINTEXT://hostname:9092   
      

      listeners=PLAINTEXT://0.0.0.0:9092

      $KAFKA_HOME/config/server.properties

    3. 使用

      重启 Kafka 服务器
      $KAFKA_HOME/bin/kafka-server-start.sh $KAFKA_HOME/config/server.properties  
      

    【讨论】:

      【解决方案3】:

      我知道这是旧的,但这可能适用于正在处理它的其他人: 我改变了两件事:
      1. 将 "bootstrap.servers" 属性或 --broker-list 选项更改为 0.0.0.0:9092
      2. 更改(在我的情况下取消注释并编辑)2 个属性中的 server.properties

      • listeners = PLAINTEXT://your.host.name:9092 到 listeners=PLAINTEXT://:9092
      • advertised.listeners=PLAINTEXT://your.host.name:9092 到 advertised.listeners=PLAINTEXT://localhost:9092

      【讨论】:

      • 0.0.0.0 是主要解决方案,因为 localhost127.0.0.1 那么它是如何工作的?
      • 不,这不是解决方案所在,只是您需要通过更改我提到的属性在服务器端正确配置服务的端点,否则它将为您返回默认的本地 ipv4 地址来自java.net.InetAddress.getCanonicalHostName(),所以为了确保您使用任何本地 ipv4 地址,您可以使用 0.0.0.0 但您也可以在正确配置后使用 127.0.0.1
      • 您是说0.0.0. 正在工作是因为其他问题吗?
      • 在我的情况下,它很容易默认绑定到 ipv6,所以只是一个注释 - 确保它在 ipv4 addy 上运行并在不需要时禁用 ipv6
      【解决方案4】:

      如果您正在运行 hortonworks 集群,请检查 ambari 中的侦听端口。

      在我的情况下 9092 不是我的端口。我去ambari发现监听端口设置为6667 它对我有用。 :)

      【讨论】:

      • 我知道你这样做是有原因的,很快你就会解决这个问题并更新代码和 hfd
      【解决方案5】:

      我遇到了类似的问题,我可以在localhost 上生产和消费,但不能从网络上的不同机器上生产和消费。根据几个答案,我得到的线索是,基本上我们需要将advertised.listener 暴露给生产者和消费者,但是提供 0.0.0.0 也不起作用。所以给出了针对advertised.listeners的确切IP

      advertised.listeners=PLAINTEXT://HOST.IP:9092

      我就这样离开了listener=PLAINTEXT://:9092

      因此,Spark 将广告的 ip 和端口暴露给生产者和消费者

      【讨论】:

      • 你是对的,如果你想从不同的机器生产或消费,你需要暴露advertised.listeners=PLAINTEXT://HOST.IP:9092
      • 谢谢,你帮我省了很多麻烦!为什么当我也将监听器设置为 HOST.IP 时它不起作用。我不明白。
      【解决方案6】:

      我今天遇到了与confluent_kafka 0.9.2 (0x90200)librdkafka 0.9.2 (0x90401) 相同的错误。就我而言,我在 tutorialpoints 示例中指定了错误的代理端口:

      $ kafka-console-producer.sh --broker-list localhost:9092 --topic tutorialpoint-basic-ops-01
      

      虽然我的代理是在 9094 端口上启动的:

      $ cat server-02.properties 
      broker.id=2
      port=9094
      log.dirs=/tmp/kafka-example-logs-02
      zookeeper.connect=localhost:2181
      

      虽然 9092 端口未打开 (netstat -tunap),但 kafka-console-producer.sh 需要 60 秒才能引发错误。看起来这个工具需要修复:

      • 更快地失败
      • 带有更明确的错误消息。

      【讨论】:

        【解决方案7】:

        我遇到了上述异常堆栈跟踪。我调查并找到了根本原因。我在使用两个节点建立 Kafka 集群时遇到了它。使用 server.properties 中的以下设置。这里我将 kafka 节点 1 和 2 的 server.properties 表示为 broker1.properties 和 broker2.properties

        broker1.properties 设置

            listeners=PLAINTEXT://A.B.C.D:9092
            zookeeper.connect=A.B.C.D:2181,E.F.G.H:2181
        

        broker2.properties 设置

            listeners=PLAINTEXT://E.F.G.H:9092
            zookeeper.connect=A.B.C.D:2181,E.F.G.H:2181
        

        我试图使用以下命令从 node1 或 node2 启动生产者: ./bin/kafka-console-producer.sh --broker-list localhost:9092 --topic OUR_TOPIC 尽管 Kafka 在两台机器上都运行,但我得到了上述超时异常堆栈跟踪。

        虽然生产者是从领导节点或追随者开始的,但我总是得到相同的结果。

        在使用来自任何代理的以下命令时,我能够获得生产者的消息。

           ./bin/kafka-console-producer.sh --broker-list A.B.C.D:9092 --topic OUR_TOPIC
            or
           ./bin/kafka-console-producer.sh --broker-list E.F.G.H:9092 --topic OUR_TOPIC
           or
           ./bin/kafka-console-producer.sh --broker-list A.B.C.D:9092,E.F.G.H:9092 --topic OUR_TOPIC
        

        所以根本原因是 Kafka 代理在启动生产者时在内部使用 listeners=PLAINTEXT://EFGH:9092 属性。此属性必须匹配才能在启动生产者时从任何节点启动 kafka 代理。转换此属性to listeners=PLAINTEXT://localhost:9092 将适用于我们的第一个命令。

        【讨论】:

          【解决方案8】:

          有这个问题: 使用 Hortonworks HDP 2.5。 启用 Kerberisation

          通过提供正确的安全协议和端口来修复。 示例命令:

          ./kafka-console-producer.sh --broker-list sand01.intranet:6667, san02.intranet:6667, san03.intranet:6667--topic test--security-protocol PLAINTEXTSASL
          
          
          ./kafka-console-consumer.sh --zookeeper sand01:2181 --topic test--from-beginning --security-protocol PLAINTEXTSASL
          

          【讨论】:

            【解决方案9】:

            在我的例子中,我将 Kafka docker 与 Openshift 一起使用。我遇到了同样的问题。当我传递值为PLAINTEXT://:9092 的环境变量KAFKA_LISTENERS 时,它得到了修复。这最终将在 server.properties 下添加创建条目 listeners=PLAINTEXT://:9092

            侦听器不必有主机名。

            【讨论】:

              【解决方案10】:

              这里是另一种情况。在我找到带有以下消息的 kafka 日志之前,不知道发生了什么:

              Caused by: java.lang.IllegalArgumentException: Invalid version for API key 3: 2
              

              显然,生产者使用的 kafka-client (java) 比 kafka 服务器更新,并且使用的 API 无效(客户端使用 1.1,服务器使用 10.0)。在客户端/生产者上,我得到了:

              Error producing to topic Failed to update metadata after 60000 ms.
              

              【讨论】:

                【解决方案11】:

                适用于 Apache Kafka v2.11-1.1.0

                启动zookeeper服务器:

                $ bin/zookeeper-server-start.sh config/zookeeper.properties
                

                启动kafka服务器:

                $ bin/kafka-server-start.sh config/server.properties
                

                创建一个主题名称“my_topic”:

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

                启动生产者:

                $ bin/kafka-console-producer.sh --broker-list localhost:9092 --topic my_topic
                

                启动消费者:

                $ bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic my_topic --from-beginning
                

                【讨论】:

                  【解决方案12】:

                  我在 Hortonworks(HDP 2.X 版本)安装上使用 Apache Kafka。遇到的错误消息意味着 Kafka 生产者无法将数据推送到段日志文件。从命令行控制台,这意味着两件事:

                  1. 您使用的代理端口不正确
                  2. 您在 server.properties 中的侦听器配置不起作用

                  如果您在通过 scala api 编写时遇到错误消息,请另外检查使用 telnet <cluster-host> <broker-port> 与 kafka 集群的连接

                  注意:如果您使用 scala api 创建主题,代理需要一些时间才能了解新创建的主题。因此,在创建主题后,生产者可能会立即失败并出现错误Failed to update metadata after 60000 ms.

                  为了解决这个问题,我做了以下检查:

                  我通过 Ambari 检查后的第一个区别是 Kafka 代理在 HDP 2.x 上侦听端口 6667(apache kafka 使用 9092)。

                  listeners=PLAINTEXT://localhost:6667
                  

                  接下来,使用 ip 而不是 localhost。 我执行了netstat -na | grep 6667

                  tcp        0      0 192.30.1.5:6667        0.0.0.0:*               LISTEN     
                  tcp        1      0 192.30.1.5:52242       192.30.1.5:6667        CLOSE_WAIT 
                  tcp        0      0 192.30.1.5:54454       192.30.1.5:6667        TIME_WAIT
                  

                  所以,我修改了生产者调用以使用 IP 而不是 localhost:

                  ./kafka-console-producer.sh --broker-list 192.30.1.5:6667 --topic rdl_test_2
                  

                  要监控是否有新记录正在写入,请监控/kafka-logs 文件夹。

                  cd /kafka-logs/<topic name>/
                  ls -lart
                  -rw-r--r--.  1 kafka hadoop        0 Feb 10 07:24 00000000000000000000.log
                  -rw-r--r--.  1 kafka hadoop 10485756 Feb 10 07:24 00000000000000000000.timeindex
                  -rw-r--r--.  1 kafka hadoop 10485760 Feb 10 07:24 00000000000000000000.index
                  

                  一旦生产者写入成功,segment log-file 00000000000000000000.log 就会变大。

                  请看下面的尺寸:

                  -rw-r--r--. 1 kafka hadoop 10485760 Feb 10 07:24 00000000000000000000.index
                  -rw-r--r--. 1 kafka hadoop       **45** Feb 10 09:16 00000000000000000000.log
                  -rw-r--r--. 1 kafka hadoop 10485756 Feb 10 07:24 00000000000000000000.timeindex
                  

                  此时,您可以运行consumer-console.sh:

                  ./kafka-console-consumer.sh --bootstrap-server 192.30.1.5:6667 --topic rdl_test_2 --from-beginning
                  response is hello world
                  

                  在这一步之后,如果您想通过 Scala API 生成消息,则更改 listeners 值(从 localhost 到公共 IP)并通过 Ambari 重新启动 Kafka 代理:

                  listeners=PLAINTEXT://192.30.1.5:6667 
                  

                  样品生产者如下:

                  package com.scalakafka.sample
                  import java.util.Properties
                  import java.util.concurrent.TimeUnit
                  
                  import org.apache.kafka.clients.producer.{ProducerRecord, KafkaProducer}
                  import org.apache.kafka.common.serialization.{StringSerializer, StringDeserializer}
                  
                  class SampleKafkaProducer {
                    case class KafkaProducerConfigs(brokerList: String = "192.30.1.5:6667") {
                      val properties = new Properties()
                      val batchsize :java.lang.Integer = 1
                  
                      properties.put("bootstrap.servers", brokerList)
                      properties.put("key.serializer", classOf[StringSerializer])
                      properties.put("value.serializer", classOf[StringSerializer])
                      //    properties.put("serializer.class", classOf[StringDeserializer])
                          properties.put("batch.size", batchsize)
                      //    properties.put("linger.ms", 1)
                      //    properties.put("buffer.memory", 33554432)
                    }
                  
                    val producer = new KafkaProducer[String, String](KafkaProducerConfigs().properties)
                  
                    def produce(topic: String, messages: Iterable[String]): Unit = {
                      messages.foreach { m =>
                        println(s"Sending $topic and message is $m")
                        val result = producer.send(new ProducerRecord(topic, m)).get()
                        println(s"the write status is ${result}")
                      }
                      producer.flush()
                      producer.close(10L, TimeUnit.MILLISECONDS)
                    }
                  }
                  

                  希望这对某人有所帮助。

                  【讨论】:

                    【解决方案13】:

                    在主题之后添加这样一行有助于解决同样的问题: ... --topic XXX --property "parse.key = true" --property "key.separator =:"

                    希望这对某人有所帮助。

                    【讨论】:

                      猜你喜欢
                      • 2019-02-27
                      • 2019-10-17
                      • 2019-10-13
                      • 1970-01-01
                      • 2017-12-06
                      • 1970-01-01
                      • 2017-09-22
                      • 2018-06-23
                      • 1970-01-01
                      相关资源
                      最近更新 更多