【发布时间】: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