【问题标题】:How to switch to custom encoder in kafka producer?如何在 kafka 生产者中切换到自定义编码器?
【发布时间】:2018-05-24 22:13:58
【问题描述】:

当我尝试在 kafka 主题中使用字符串作为键时出现以下错误。

18/05/14 17:08:26 ERROR async.DefaultEventHandler: Error serializing message for topic my_topic
java.lang.ClassCastException: java.lang.String cannot be cast to [B
    at kafka.serializer.DefaultEncoder.toBytes(Encoder.scala:34)
    at kafka.producer.async.DefaultEventHandler$$anonfun$serialize$1.apply(DefaultEventHandler.scala:130)
    at kafka.producer.async.DefaultEventHandler$$anonfun$serialize$1.apply(DefaultEventHandler.scala:127)
    at scala.collection.mutable.ResizableArray$class.foreach(ResizableArray.scala:59)
    at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:47)
    at kafka.producer.async.DefaultEventHandler.serialize(DefaultEventHandler.scala:127)
    at kafka.producer.async.DefaultEventHandler.handle(DefaultEventHandler.scala:53)
    at kafka.producer.async.ProducerSendThread.tryToHandle(ProducerSendThread.scala:105)
    at kafka.producer.async.ProducerSendThread$$anonfun$processEvents$3.apply(ProducerSendThread.scala:88)
    at kafka.producer.async.ProducerSendThread$$anonfun$processEvents$3.apply(ProducerSendThread.scala:68)
    at scala.collection.immutable.Stream.foreach(Stream.scala:547)
    at kafka.producer.async.ProducerSendThread.processEvents(ProducerSendThread.scala:67)
    at kafka.producer.async.ProducerSendThread.run(ProducerSendThread.scala:45)

问题似乎在于默认编码器

public class DefaultEncoder implements Encoder<byte[]> 

不支持字符串转字节

public byte[] toBytes(byte[] value) {
    return value;
}

向生产者提供自定义编码器的正确方法是什么?

我是否也必须在消费者方面做出改变?

【问题讨论】:

  • 代码的哪一部分试图获取字符串?是的,生产者和消费者必须就主题数据格式达成一致。另外,ByteArraySerializer 已经存在

标签: java apache-kafka


【解决方案1】:

在 Producer/Consumer 上,您可以指定不同的序列化程序而不是默认的 ByteArraySerializer,Here 您可以找到版本 1.1.0 的可用序列化程序,或者您可以指定自己的序列化程序。

一般来说,如果您发送/接收字符串或 Json,您可以使用默认的 org.apache.kafka.common.serialization.StringSerializerorg.apache.kafka.common.serialization.StringDeserializer。对于生产者/消费者,属性是:

private void configureProducer() {
    Properties props = new Properties();
    props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
    props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
}

private void configureConsumer() {
    Properties props = new Properties();
    props.put("key.serializer", "org.apache.kafka.common.serialization.StringDeserializer");
    props.put("value.serializer", "org.apache.kafka.common.serialization.StringDeserializer");
}

Properties producerProps = configureProducer();
KafkaProducer<String, String> producer = new KafkaProducer<>(producerProps);
Properties consumerProps = configureConsumer();
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerProps);

如果您想发送您自己的自定义 bean,请将其转换为 json 并使用默认的 StringSerializer/Deserializer,或者构建您自己的序列化/反序列化类。你可以看到一个使用带有spring boot的json字符串的例子here

【讨论】:

    猜你喜欢
    • 2015-10-03
    • 2016-09-25
    • 2019-12-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-03-22
    • 1970-01-01
    相关资源
    最近更新 更多