【发布时间】:2015-10-28 02:38:01
【问题描述】:
我在 CentOS 环境中设置了 Apache Kafka 0.8.2.1,创建了一个主题并通过命令行生产者/消费者发送/接收了一些虚拟消息。
正如您在屏幕截图中看到的那样,效果很好。不,我正在编写一个自定义生产者来将我的消息从 Java 发送到主题。
package de.jofre.kafka;
import java.util.Properties;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;
public class TestProducer {
public static void main(String[] args) {
Properties props = new Properties();
props.put("bootstrap.servers", "192.168.145.130:9092");
props.put("key.serializer", StringSerializer.class.getName());
props.put("value.serializer", StringSerializer.class.getName());
KafkaProducer<String, String> prod = new KafkaProducer<String, String>(props);
ProducerRecord<String, String> record = new ProducerRecord<String, String>("test", "Kafka is great");
prod.send(record);
prod.close();
}
}
调用显示的 main 不会产生关于 Kafka 主题的消息,也不会打印任何错误消息。
有人知道为什么消息没有到达我的主题吗?
【问题讨论】:
-
请您尝试删除
props.put("key.serializer", StringSerializer.class.getName()); -
仅更改 bootstrap.servers 属性和要写入的主题,您的代码可以在我的机器上完美运行。关于正在发生的事情,您可以提供更多信息吗?因为问题似乎不在于您的代码,假设您的 ipaddress 和主题名称是正确的。
-
@user2720864:当我删除该属性时,我得到一个 ConfigException,告诉我 key.serializer 属性是必需的。
-
@morganw09dev:感谢您的测试。您使用版本 0.8.2.1 吗?您使用什么版本的 kafka java 库?我在我们的测试环境以及本地 VM 中尝试了代码。它们都不起作用。
-
是的,我在 0.8.2.1 和 0.8.2.0 版本上进行了测试。我在 pom.xml 中的唯一依赖项是
org.apache.kafka kafka_2.10 0.8.2.1 有一些排除项。
标签: java apache-kafka