【问题标题】:Kafka 0.9.0.1 Java Consumer stuck in awaitMetadataUpdate()Kafka 0.9.0.1 Java Consumer 卡在 awaitMetadataUpdate()
【发布时间】:2016-10-12 16:58:50
【问题描述】:

我正在尝试使用 Java API v0.9.0.1 让一个简单的 Kafka Consumer 工作。我使用的 kafka 服务器是一个 docker 容器,也运行版本 0.9.0.1。下面是消费者代码:

public class Consumer {
    public static void main(String[] args) throws IOException {

        KafkaConsumer<String, String> consumer;
        try (InputStream props = Resources.getResource("consumer.props").openStream()) {
            Properties properties = new Properties();
            properties.load(props);
            consumer = new KafkaConsumer<>(properties);
        }

        consumer.subscribe(Arrays.asList("messages"));
        try {
            while (true) {
                ConsumerRecords<String, String> records = consumer.poll(100);
                for (ConsumerRecord<String, String> record : records)
                    System.out.println("Message received: " + record.value());
            }
        }catch(WakeupException ex){
            System.out.println("Exception caught " + ex.getMessage());
        }finally{
            consumer.close();
            System.out.println("After closing KafkaConsumer");
        }
    }
}

但是,当启动消费者时,它会调用上面的 poll(100) 方法并且永远不会返回。调试,看起来它在 org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient 中永远运行以下方法:

public void awaitMetadataUpdate() {
    int version = this.metadata.requestUpdate();

    do {
        this.poll(9223372036854775807L);
    } while(this.metadata.version() == version);

}

(版本和 this.metadata.version() 似乎总是 == 2)。此外,尽管它没有引发任何错误,但来自我的 java 生产者的消息从未出现在队列中。我已经验证使用命令行 kafka 工具,我可以发送和接收来自队列的消息。

有人知道这里发生了什么吗?

【问题讨论】:

  • 这通常是与您的代理相关的问题,该问题未通告生产者和消费者可访问的端点。您的经纪人是否在从外部监听可访问的地址?检查代理属性 advertized.listeners
  • 卢西亚诺,你很准。在服务器上设置环境变量 ADVERTISED_PORT 和 ADVERTISED_HOST 解决了这个问题。有点令人困惑的是,没有这些,命令行消费者/生产者可以正常工作,但 java 实现不能。
  • 似乎是一个已知问题:issues.apache.org/jira/browse/KAFKA-3727
  • @cacois 当您说要在服务器上设置这些变量时,您是指哪个服务器?是与消费者一起运行 Java 代码的机器吗?还是 docker 容器?
  • 另外,您使用的是 Mac、Linux 还是 Win? :-)

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


【解决方案1】:

如果这有助于其他有类似问题的人,我的解决方案是设置以下环境变量:

ADVERTISED_HOST=localhost
ADVERTISED_PORT=9092

(当然,这里的值可能会根据您的安装而改变)

显然,命令行消费者和生产者脚本可以在不设置这些环境变量的情况下正确地找到代理并与之通信,但 Java API 实现却不能。也不会抛出任何错误,只是在第一次轮询时尝试更新元数据时出现无限循环。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2023-03-02
    • 1970-01-01
    • 1970-01-01
    • 2021-11-21
    • 1970-01-01
    • 2017-03-02
    • 2020-05-19
    相关资源
    最近更新 更多