【发布时间】: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