Kafka Consumer的poll()方法被阻止

ara*_*ran 9 polling apache-kafka kafka-consumer-api kafka-producer-api

我是Kafka 0.9的新手并测试了一些功能,我在Java实现的Consumer(KafkaConsumer)中发现了一个奇怪的行为.

Kafka经纪人位于Ambari外部机器中.

即使你我可以实现一个Producer并开始向外部代理发送消息,我也不知道为什么当消费者试图读取事件(民意调查)时,它会被卡住.

我知道生产者工作得很好,因为我可以通过控制台消费者(在ambari本地工作)消费消息.但是当我执行Java Consumer时,什么都没发生,只是卡住了.调试代码我可以看到它在该poll()行被阻止:

    ConsumerRecords<String, String> records = consumer.poll(100);
Run Code Online (Sandbox Code Playgroud)

顺便说一句,超时没有任何作用.如果你输入0,100或1000毫秒无关紧要,消费者在这一行被阻止并且不会超时也不会抛出异常.

我尝试了所有类型的替代属性,例如advertised.host.name,advertised.listener,...等等,运气不好.

任何帮助将受到高度赞赏.提前致谢!

sud*_*anu 6

原因可能是运行消费者代码的计算机无法连接到zookeeper.尝试在安装了Kafka的机器上运行相同的消费者代码(我试过这个并为我工作).我也通过在server.properties文件中提到以下属性来解决问题: advertised.host.name="ip address which you want to expose" //在我的情况下它是ec2机器的公共IP,我在相同的ec2上安装了kafka和zookeeper. advertised.port=9092 ConsumerRecords<String, String> records = consumer.poll(100); 以上陈述并不意味着消费者将在100毫秒后超时,这是投票期.无论在100毫秒内捕获的数据是什么,都会被读入记录集.