Kafka 0.9.0.1 Java Consumer застрял в awaitMetadataUpdate ()
Я пытаюсь заставить простого Kafka Consumer работать с использованием Java API v0.9.0.1. Сервер kafka, который я использую, является док-контейнером, также работающим с версией 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, я могу отправлять и получать сообщения из очереди.
Кто-нибудь знает, что здесь происходит?