【发布时间】:2019-05-30 16:05:53
【问题描述】:
您好,我一直在尝试学习 KAFKA,但我的远程轮询器/消费者遇到了问题。
我已经在 AWS EC2 实例中使用私有和公共 ip 设置了 KAFKA。我的 server.properties 看起来像这样。
listeners=PLAINTEXT://172.31.31.58:9092 #AWS Private IP
advertised.listeners=PLAINTEXT://35.??.??.??:9092 #AWS Public IP Masked
我的 AWS EC2 安全组配置为允许通过任何端口上的任何 ip 进行流量以进行测试。
当我使用以下脚本在我的 EC2 实例中本地生成/使用消息时,它可以完美运行
bin/kafka-console-producer.sh --broker-list localhost:9092 --topic test
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test --from-beginning
但是,当我尝试从运行 java API 的远程笔记本电脑 Eclipse 代码连接到同一个 kafka 实例时,我的代码在 consumer.poll(100) 中永远挂起。我在这里做错了吗?
Properties props = new Properties();
props.put("bootstrap.servers", "35.??.??.??:9092");//my aws public ip configured in advertised.listeners
props.put("group.id", "test123");
props.put("enable.auto.commit", "false");
props.put("auto.commit.interval.ms", "1000");
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("test"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(100);
for (ConsumerRecord<String, String> record : records)
System.out.printf("offset = %d, key = %s, value = %s", record.offset(), record.key(), record.value());
}
}
【问题讨论】:
-
您可以为您的消费者发布日志文件吗?查看调试消息会很有帮助。
标签: amazon-web-services apache-kafka consumer producer