【问题标题】:kafka-python: producer is not able to connectkafka-python:生产者无法连接
【发布时间】:2016-06-11 21:12:02
【问题描述】:

kafka-python (1.0.0) 在连接代理时抛出错误。 同时 /usr/bin/kafka-console-producer 和 /usr/bin/kafka-console-consumer 工作正常。

以前python应用也可以,但是重启zookeeper后就连接不上了。

我正在使用文档中的简单示例:

from kafka import KafkaProducer
from kafka.common import KafkaError

producer = KafkaProducer(bootstrap_servers=['hostname:9092'])

# Asynchronous by default
future = producer.send('test-topic', b'raw_bytes')

我收到此错误:

Traceback (most recent call last):   File "pp.py", line 4, in <module>
    producer = KafkaProducer(bootstrap_servers=['hostname:9092'])   File "/usr/lib/python2.6/site-packages/kafka/producer/kafka.py", line 246, in __init__
    self.config['api_version'] = client.check_version()   File "/usr/lib/python2.6/site-packages/kafka/client_async.py", line 629, in check_version
    connect(node_id)   File "/usr/lib/python2.6/site-packages/kafka/client_async.py", line 592, in connect
    raise Errors.NodeNotReadyError(node_id) kafka.common.NodeNotReadyError: 0 Exception AttributeError: "'KafkaProducer' object has no attribute '_closed'" in <bound method KafkaProducer.__del__ of <kafka.producer.kafka.KafkaProducer object at 0x7f6171294c50>> ignored

在单步执行 (/usr/lib/python2.6/site-packages/kafka/client_async.py) 时,我注意到第 270 行的评估结果为 false:

270         if not self._metadata_refresh_in_progress and not self.cluster.ttl() == 0:
271             if self._can_send_request(node_id):
272                 return True
273         return False

在我的例子中 self._metadata_refresh_in_progress 是 False,但是 ttl() = 0;

与此同时,kafka-console-* 正在愉快地推送消息:

/usr/bin/kafka-console-producer --broker-list hostname:9092 --topic test-topic
hello again
hello2

有什么建议吗?

【问题讨论】:

    标签: apache-kafka kafka-python


    【解决方案1】:

    我遇到了类似的问题。就我而言,代理主机名在客户端无法解析。尝试在配置文件中显式设置advertised.host.name。

    【讨论】:

      【解决方案2】:

      一个主机可以有多个 dns 别名。它们中的任何一个都可以用于 ssh 或 ping 测试。但是,kafka 连接应该使用与代理的server.properties 文件中的advertised.host.name 匹配的别名。

      我在 bootstrap_servers 参数中使用了不同的别名。因此出现错误。一旦我将呼叫更改为使用advertised.hostname,问题就解决了

      【讨论】:

        【解决方案3】:

        我遇到了类似的问题,从 bootstrap_servers 中删除端口有所帮助。

        consumer = KafkaConsumer('my_topic',
                             #group_id='x',
                             bootstrap_servers='kafka.com')
        

        【讨论】:

          【解决方案4】:

          我遇到了同样的问题,上面的解决方案都没有奏效。然后我阅读了异常消息,似乎必须指定api_version,所以

          producer = KafkaProducer(bootstrap_servers=['localhost:9092'],api_version=(0,1,0))
          

          注意:元组(1,0,0)匹配kafka版本1.0.0

          工作正常(至少无例外地完成,现在必须说服它接受消息;))

          【讨论】:

          • 但是在发送第一条消息后它就停止了
          • 这对我来说就像一个魅力。否则,Kafka 0.10 为我破坏了 kafka-python。
          • 如果有人在第一条消息后没有发送更多数据时遇到问题,您必须执行producer.flush() 以清除发送缓冲区。
          • api_version 在 1.3.5 版本中不是强制性的。但是,不使用 keyed 参数将使其尝试自动检测版本,该版本在我当前的代理上失败(docker image wurstmeister/kafka-docker)
          【解决方案5】:

          在您的 server.properties 文件中,确保将侦听器 IP 设置为远程机器可访问的框 IP 地址。默认是本地主机

          在你的 server.properties 中更新这一行:

          listeners=PLAINTEXT://<Your-IP-address>:9092
          

          还要确保您没有可能阻止其他 IP 地址与您联系的防火墙。如果你有 sudo 特权。尝试禁用防火墙。

          sudo systemctl stop firewalld
          

          【讨论】:

            【解决方案6】:

            我遇到了同样的问题。

            我用 user3503929 的提示解决了这个问题。

            kafka 服务器安装在 windows 上。

            server.properties

            ...
            host.name = 0.0.0.0
            ...
            

            .

            producer = KafkaProducer(bootstrap_servers='192.168.1.3:9092',         
                                                     value_serializer=str.encode)
            producer.send('test', value='aaa')
            producer.close()
            print("DONE.")
            

            windows kafka客户端处理没有问题。 但是,当我在 ubuntu 中使用 kafka-python 向主题发送消息时,会引发 NoBrokersAvailable 异常。

            将以下设置添加到 server.properties。

            ...
            advertised.host.name = 192.168.1.3
            ...
            

            它在相同的代码中成功运行。 为此,我花了三个小时。

            谢谢

            【讨论】:

              【解决方案7】:

              使用pip install kafka-python安装kafka-python

              创建 kafka 数据管道的步骤:-
              1.使用shell命令运行Zookeeper或使用安装zookeeperd

              sudo apt-get install zookeeperd 
              

              这会将 zookeeper 作为守护进程运行,默认监听 2181 端口

              1. 运行 kafka 服务器
              2. 在不同的控制台上运行带有 producer.py 和 consumer.py 的脚本以查看实时数据。

              以下是要运行的命令:-

              cd kafka-directory
              ./bin/zookeeper-server-start.sh  ./config/zookeeper.properties    
              ./bin/kafka-server-start.sh  ./config/server.properties
              

              现在你已经运行了 zookeeper 和 kafka 服务器,运行 producer.py 脚本和 consumer.py

              生产者.py:

              从 kafka 导入 KafkaProducer 进口时间

              producer = KafkaProducer(bootstrap_servers=['localhost:9092'])
              topic = 'test'
              lines = ["1","2","3","4","5","6","7","8"]
              for line in lines:
                try:
                  producer.send(topic, bytes(line, "UTF-8")).get(timeout=10)
                except IndexError as e:
                  print(e)
                continue
              

              Consumer.py:-

              from kafka import KafkaConsumer
              topic = 'test'
              consumer = KafkaConsumer(topic, bootstrap_servers=['localhost:9092'])
              for message in consumer:
                  # message value and key are raw bytes -- decode if necessary!
                  # e.g., for unicode: `message.value.decode('utf-8')`
                  # print ("%s:%d:%d: key=%s value=%s" % (message.topic, message.partition,
                  #                                       message.offset, message.key,
                  #                                       message.value))
                  print(message)
              

              现在在不同的终端运行 producer.py 和 consumer.py 以查看实时数据..!

              注意:上面的 producer.py 脚本只运行一次就可以永久运行,使用 while 循环和使用 time 模块。

              【讨论】:

              • 在kafka-python 2.0.2 的新Mac M1 上完美运行。主要是因为.get(timeout=10)
              猜你喜欢
              • 1970-01-01
              • 2013-02-19
              • 2018-06-16
              • 2016-12-18
              • 2020-06-25
              • 2019-10-23
              • 2019-06-15
              • 1970-01-01
              • 1970-01-01
              相关资源
              最近更新 更多