【发布时间】:2017-11-17 19:03:20
【问题描述】:
我正在使用 Kafka 0.11.0.0。我有一个发布到 Kafka 主题的测试程序;如果 zookeeper 和 Kafka 服务器宕机(这在我的开发环境中是正常的;我会根据需要启动它们)然后对 KafkaProducer.send() 的调用将无限期挂起。
我要么需要 send() 返回,最好是指示错误;或者我需要一种方法来检查服务器是启动还是关闭。基本上,我希望我的测试工具能够告诉我,“嘿,笨蛋,启动 Kafka!”而不是挂起。
我的生产者任务有没有办法确定服务器是启动还是关闭?
我这样调用 send():
kafkaProducer.send(new ProducerRecord<>(KAFKA_TOPIC, KAFKA_KEY,
message), (rm, ex) -> {
System.out.println("**** " + rm + "\n**** " +ex);
});
我有 linger.ms = 1;我尝试了 retries=0、1 和 2,但 send() 仍然阻塞。我从未见过回调被调用。
较早的消息建议将 metadata.fetch.timeout.ms 设置为一个较小的值,但在 0.11 中已取消。其他人建议调用命令行实用程序以查看服务器是否正常......但引用的实用程序似乎也不见了。
完成这项工作的优雅方式是什么?
【问题讨论】:
-
这很奇怪。它应该返回一个错误,说明“更新元数据失败”或“正在过期 x 条记录”。检查生产者的 request.timeout.ms 和 max.block.ms 设置。默认 request.timeout.ms 为 60 秒
-
request.timeout.ms 是 30000,max.block.ms 是 60000。我会尝试减少这些。 (以交互方式工作时,30 秒还不如无限期——我的错。)
-
好的,是的,这解决了我眼前的问题;谢谢!
-
不过,我仍然想知道如何主动检查服务器是否已启动。
-
如果对您有帮助,请您接受我的回答 :-)
标签: apache-kafka kafka-producer-api