【问题标题】:Not able to publish Avro messages to Kafka Topic无法将 Avro 消息发布到 Kafka 主题
【发布时间】:2018-12-12 03:10:22
【问题描述】:

我使用以下命令启动了 Kafka

docker run -p 2181:2181 -p 9092:9092 -p 8081:8081 --env ADVERTISED_HOST=`docker-machine ip \`docker-machine active\`` --env ADVERTISED_PORT=9092 spotify/kafka

现在,我编写了一个简单的程序,将字符串发布到 kafka 主题中。它没有任何问题。

props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "192.168.99.100:9092")
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer")
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer")
val producer = new KafkaProducer[String, String](props)
val inputRecord = new ProducerRecord[String, String]("test", "key2", "Hello World")
producer.send(inputRecord)
producer.close()

所以现在我修改了这个程序并尝试向 kafka 主题发送一条 avro 消息

val props = new Properties()
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "192.168.99.100:9092")
props.put("schema.registry.url", "http://192.168.99.100:8081")
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer")
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "io.confluent.kafka.serializers.KafkaAvroSerializer")
val producer = new KafkaProducer[String, Object](props)
val inputRecord = createAvroRecord(schemaStr, "test1", "test1")
val producer: KafkaProducer[String, Object] = CreateProducerAvro
val producerAvroRecord = new ProducerRecord[String, Object]("test", "key1", inputRecord)
producer.send(producerAvroRecord)
producer.close()

但我得到错误

[error] (run-main-0) org.apache.kafka.common.errors.SerializationException: Error serializing Avro message
org.apache.kafka.common.errors.SerializationException: Error serializing Avro message
Caused by: java.net.ConnectException: Connection refused
    at java.net.PlainSocketImpl.socketConnect(Native Method)
    at java.net.AbstractPlainSocketImpl.doConnect(AbstractPlainSocketImpl.java:339)
    at java.net.AbstractPlainSocketImpl.connectToAddress(AbstractPlainSocketImpl.java:200)
    at java.net.AbstractPlainSocketImpl.connect(AbstractPlainSocketImpl.java:182)
    at java.net.SocksSocketImpl.connect(SocksSocketImpl.java:392)
    at java.net.Socket.connect(Socket.java:579)
    at java.net.Socket.connect(Socket.java:528)
    at sun.net.NetworkClient.doConnect(NetworkClient.java:180)
    at sun.net.www.http.HttpClient.openServer(HttpClient.java:432)
    at sun.net.www.http.HttpClient.openServer(HttpClient.java:527)
    at sun.net.www.http.HttpClient.<init>(HttpClient.java:211)
    at sun.net.www.http.HttpClient.New(HttpClient.java:308)
    at sun.net.www.http.HttpClient.New(HttpClient.java:326)
    at sun.net.www.protocol.http.HttpURLConnection.getNewHttpClient(HttpURLConnection.java:997)
    at sun.net.www.protocol.http.HttpURLConnection.plainConnect(HttpURLConnection.java:933)
    at sun.net.www.protocol.http.HttpURLConnection.connect(HttpURLConnection.java:851)
    at sun.net.www.protocol.http.HttpURLConnection.getOutputStream(HttpURLConnection.java:1092)
    at io.confluent.kafka.schemaregistry.client.rest.utils.RestUtils.httpRequest(RestUtils.java:128)
    at io.confluent.kafka.schemaregistry.client.rest.utils.RestUtils.registerSchema(RestUtils.java:174)
    at io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient.registerAndGetId(CachedSchemaRegistryClient.java:51)
    at io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient.register(CachedSchemaRegistryClient.java:89)
    at io.confluent.kafka.serializers.AbstractKafkaAvroSerializer.serializeImpl(AbstractKafkaAvroSerializer.java:49)
    at io.confluent.kafka.serializers.KafkaAvroSerializer.serialize(KafkaAvroSerializer.java:67)
    at org.apache.kafka.clients.producer.KafkaProducer.send(KafkaProducer.java:424)
    at org.apache.kafka.clients.producer.KafkaProducer.send(KafkaProducer.java:339)
    at KafkaPublisher$.SendAvroMessage(KafkaPublisher.scala:35)
    at KafkaPublisher$.main(KafkaPublisher.scala:20)
    at KafkaPublisher.main(KafkaPublisher.scala)
    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:57)
    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.lang.reflect.Method.invoke(Method.java:606)
[trace] Stack trace suppressed: run last compile:run for the full output.

【问题讨论】:

    标签: apache-kafka kafka-producer-api


    【解决方案1】:

    将你的 key.serializer 和 value.serializer 更改为 avro,如下所示。

    props.put("key.serializer", io.confluent.kafka.serializers.KafkaAvroSerializer.class);
    props.put("value.serializer", io.confluent.kafka.serializers.KafkaAvroSerializer.class);
    

    如果您使用的是架构注册表,请设置架构注册表 url

    props.put("schema.registry.url", "http://localhost:8081");
    

    注意: 将您的架构注册表 url 设置为您正在运行架构注册表的服务器,主题名称为“Kafka-value” 如果你不使用上面的,你可以丢弃这个。

    【讨论】:

      【解决方案2】:

      我假设您正在使用 Windows 或至少使用 docker 机器,并且您在 bootstrap_serversschema.registry.url 上设置的 ip 地址是 docker 机器 ip 地址。

      尝试删除两个选项:ADVERTISED_HOSTADVERTISED_PORT,通常在这种情况下您的端口绑定就足够了。

      【讨论】:

        猜你喜欢
        • 2018-01-12
        • 1970-01-01
        • 1970-01-01
        • 2018-10-28
        • 2016-01-17
        • 2018-10-27
        • 2020-01-18
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多