【问题标题】:How can I produce messages with Kafka 8.2 API in Java?如何使用 Java 中的 Kafka 8.2 API 生成消息?
【发布时间】:2015-10-26 09:07:10
【问题描述】:

我正在尝试在 java 中使用 kafka API。我正在使用以下 Maven 依赖项:

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>0.8.2.0</version>
</dependency>

我无法连接到远程 kafka 服务器。 我将 kafka 'server.properties' 文件端口属性更改为端口 8080。 我可以同时启动zookeeper和kafka服务器没问题。 我还可以使用 kafka 下载附带的控制台生产者和消费者应用程序。 (Scala 2.10 版本)

我正在使用以下客户端代码创建远程 KafkaProducer

Properties propsProducer = new Properties();

propsProducer.put("bootstrap.servers", "172.xx.xx.xxx:8080");
propsProducer.put("key.serializer", org.apache.kafka.common.serialization.ByteArraySerializer.class);
propsProducer.put("value.serializer", org.apache.kafka.common.serialization.ByteArraySerializer.class);
propsProducer.put("topic.metadata.refresh.interval.ms", "0");

KafkaProducer<byte[], byte[]> m_kafkaProducer = new KafkaProducer<byte[], byte[]>(propsProducer);

创建生产者后,我可以运行以下行并返回有效的主题信息,授予 strTopic 是现有主题名称。

List<PartitionInfo> partitionInfo = m_kafkaProducer.partitionsFor(strTopic);

当我尝试发送消息时,我会执行以下操作:

ProducerRecord<byte[], byte[]> prMessage = new ProducerRecord<byte[],byte[]>(strTopic, strMessage.getBytes());

RecordMetadata futureData = m_kafkaProducer.send(prMessage).get();

对 send() 的调用无限期阻塞,当我手动终止进程时,我看到 ERROR Closing socket because of error on kafka server(IOException, Connection Reset by Peer) 错误。

此外,host.name、adverted.host.name 和adverted.port 属性仍然在“server.properties”文件中被注释掉,这毫无价值。哦,如果我换行:

propsProducer.put("bootstrap.servers", "172.xx.xx.xxx:8080");

propsProducer.put("bootstrap.servers", "127.0.0.1:8080");

并在安装了 kafka 服务器的同一台服务器上运行它,它可以工作,但我正在尝试远程使用它。

感谢任何帮助,如果我能澄清一下,请告诉我。

【问题讨论】:

  • 您是在使用172.xx.xx.xxx 作为主机IP 地址吗?
  • 不,这是一个完整的 IP,x 只是掩码。
  • Kk。也许是防火墙问题?您可以使用 netcat 验证端口 8080 上的网络连接吗?
  • 验证命令 'nc -vz localhost 8080' 在主机服务器上成功。 “netstat -plunt”还将端口发布为打开并在所有接口上侦听。
  • @Baron我很抱歉地说我对 Kafka 的了解还不够,无法在这一点上提供更多帮助。我偷看了文档,没有看到任何关于 IP 白名单等的信息。也许我明天会想点什么。

标签: java sockets maven apache-kafka


【解决方案1】:

经过大量挖掘,我决定实现此处找到的示例:Kafka Producer Example。我缩短了代码并且没有实现分区器类。我用列出的依赖项更新了我的 pom,但我仍然遇到同样的问题。最终,我进行了一些配置更改,一切正常。

最后一个难题是在服务器和客户端机器的 /etc/hosts 中定义 Kafka 服务器。我在两个文件中都添加了以下内容。

172.xx.xx.xxx     serverHost1

同样,x 只是掩码。然后,我将 server.properties 文件中的 Advertisementd.host.name 设置为 serverHost1。注意:我在服务器机器上运行 ifconfig 后获得了该 IP。

我换行了

propsProducer.put("metadata.broker.list", "172.xx.xx.xxx:8080");

propsProducer.put("metadata.broker.list", "serverHost1:8080");

Kafka API 不喜欢我将 IP 定义为字符串这一事实。相反,它是从 etc/hosts 文件中查找 IP,尽管文档说:

“代理将向生产者和消费者通告的主机名。如果未设置,它将使用“host.name”的值(如果已配置)。否则,它将使用从 java.net.InetAddress.getCanonicalHostName() 返回的值。 "

如果没有在客户端机器的 etc/hosts 中定义,它只会以字符串形式返回我之前使用的 IP,否则它会返回与 IP 配对的名称(在我的例子中是 serverHost1)。另外,我也从未设置过 host.name 的值。

【讨论】:

  • bootstrap.servers 是 metadata.broker.list 的替代品吗?
  • 是的,我相信。在 0.8.2.0 版本中,该字段为“metadata.broker.list”,但在较新版本中为“boostrap.servers”
  • 是的!那是真的。使用新的 ProducerAPI,这是新的配置。
猜你喜欢
  • 2021-02-16
  • 1970-01-01
  • 2018-10-16
  • 2021-05-08
  • 2018-03-09
  • 1970-01-01
  • 2015-10-28
  • 1970-01-01
  • 2020-09-20
相关资源
最近更新 更多