【问题标题】:Python kafka consumer wont consume messages from producerPython kafka消费者不会消费来自生产者的消息
【发布时间】:2019-07-01 20:47:44
【问题描述】:

所以我对卡夫卡还很陌生。我试图运行一个简单的 Kafka 消费者和生产者。当我运行我的消费者时,它会在 for 循环之前打印 hello 。但是在 for 循环中什么都没有打印,这让我相信它从一开始就不会进入 for 循环,并且消费者不会消费来自生产者的消息。我在linux系统上运行它。

任何人都可以就生产者或消费者可能出现的问题提供建议吗?我已经展示了我的生产者和消费者代码,它们都只有几行代码。

这是我的制作人:

from kafka import KafkaProducer
producer = KafkaProducer(bootstrap_servers='localhost:2181',api_version=(1,0,1))
producer.send('MyFirstTopic1', 'Hello, World!')

这是我的消费者:

from kafka import KafkaConsumer,KafkaProducer,TopicPartition,OffsetAndMetadata
consumer = KafkaConsumer(
 bootstrap_servers=['localhost:2181'],api_version=(1,0,1),
 group_id=None,
 enable_auto_commit=False,
 auto_offset_reset='smallest'
)
consumer.subscribe('MyFirstTopic1',0)
print("hello")
for message in consumer:
 print(message)

因此,当运行我的生产者时,它最终会出错。任何人都知道这意味着什么,以及这是否可能是错误的。

File "producer.py", line 3, in <module>
    producer.send('MyFirstTopic1', 'Hello, World!')
  File "/usr/local/lib/python3.5/site-packages/kafka/producer/kafka.py", line 543, in send
    self._wait_on_metadata(topic, self.config['max_block_ms'] / 1000.0)
  File "/usr/local/lib/python3.5/site-packages/kafka/producer/kafka.py", line 664, in _wait_on_metadata
    "Failed to update metadata after %.1f secs." % max_wait)
kafka.errors.KafkaTimeoutError: KafkaTimeoutError: Failed to update metadata after 60.0 secs.

【问题讨论】:

  • 投票结束是一个错字。再次查看文档,发现示例中没有使用:2181 端口
  • 你的 Kafka 版本是多少?也许你不需要设置api_version

标签: python apache-kafka kafka-consumer-api kafka-producer-api


【解决方案1】:

您的客户端配置中似乎使用了错误的主机。 localhost:2181 通常是 Zookeeper 服务器。

为了让您的客户端正常工作,您需要将 bootstrap_servers 设置为 Kafka 代理主机名和端口。默认为localhost:9092

https://kafka-python.readthedocs.io/en/latest/apidoc/KafkaProducer.html

【讨论】:

  • 即使我将其切换到 localhost:9092 我仍然会收到相同的错误消息。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-09-19
  • 1970-01-01
  • 2023-01-21
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多