【问题标题】:KAFKA REMOTE AWS consumer.pollKAFKA REMOTE AWS 消费者.poll
【发布时间】: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


【解决方案1】:

你确定它挂在 poll() 吗?还是poll() 只是返回一个空的ConsumerRecords 并且它在while(true) 中循环?

默认情况下,如果您没有为组提交任何偏移量,则消费者从主题的末尾开始,因此它只会接收新消息。在这种情况下,如果您想使用该主题中已经存在的消息,则需要将 auto.offset.reset 设置为 earliest(就像您在控制台消费者中使用 --from-beginning 所做的那样)

编辑:

如果它实际上卡在poll() 中,则可能是连接问题。要找出答案,最好的方法是在启用日志记录的情况下运行您的客户端。创建一个包含:

log4j.rootLogger=DEBUG, stdout
log4j.appender.stdout=org.apache.log4j.ConsoleAppender
log4j.appender.stdout.layout=org.apache.log4j.PatternLayout
log4j.appender.stdout.layout.ConversionPattern=[%d] %p %m (%c)%n

并使用 -Dlog4j.configuration=file:PATH_TO_FILE 启动您的客户端

【讨论】:

  • 是的 Mickael,它在投票中挂起,基本上它不会移动到下一行,我已经验证了这一点。
  • 嗨 Mickael,这绝对是网络问题。我早些时候在我的办公室网络中尝试过这个,它没有工作并且在不同的网络中完美地工作。事实上,我按照您的建议在启用调试日志的情况下运行了我的程序,并且没有记录任何提示任何网络问题的记录。但这现在可以在不同的网络中使用,您对网络问题的建议使我考虑通过不同的网络运行它。非常感谢。
猜你喜欢
  • 2020-09-19
  • 2017-11-16
  • 1970-01-01
  • 2020-07-10
  • 2016-07-19
  • 1970-01-01
  • 2017-01-04
  • 2017-09-23
  • 1970-01-01
相关资源
最近更新 更多