【问题标题】:Send Custom Java Objects to Kafka Topic将自定义 Java 对象发送到 Kafka 主题
【发布时间】:2017-04-29 17:50:20
【问题描述】:

我有我的自定义 Java 对象并希望利用 JVM 的内置序列化将其发送到 Kafka 主题,但序列化失败并出现以下错误

org.apache.kafka.common.errors.SerializationException: 无法转换 com.spring.kafka.Payload 类的值到类 org.apache.kafka.common.serialization.ByteArraySerializer 中指定 value.serializer

Payload.java

public class Payload implements Serializable {

    private static final long serialVersionUID = 123L;

    private String name="vinod";

    private int anInt = 5;

    private Double aDouble = new Double("5.0");

    public String getName() {
        return name;
    }

    public void setName(String name) {
        this.name = name;
    }

    public int getAnInt() {
        return anInt;
    }

    public void setAnInt(int anInt) {
        this.anInt = anInt;
    }

    public Double getaDouble() {
        return aDouble;
    }

    public void setaDouble(Double aDouble) {
        this.aDouble = aDouble;
    }

}

在创建生产者期间,我设置了以下属性

<entry key="key.serializer"
                       value="org.apache.kafka.common.serialization.ByteArraySerializer" />
                <entry key="value.serializer"
                       value="org.apache.kafka.common.serialization.ByteArraySerializer" />

我的发送调用如下

kafkaProducer.send(new ProducerRecord<String, Payload>("test", new Payload()));

通过生产者将自定义 java 对象发送到 kafka 主题的正确方法是什么?

【问题讨论】:

  • 其他选项是转换为 JSON 格式并发送

标签: java serialization apache-kafka


【解决方案1】:

由于您使用的是ByteArraySerializer,因此您需要实例化一个 byte[] 生产者。

Producer<byte[],byte[]> producer = new KafkaProducer<>(props);

然后在产生序列化或其他方法后传递字节[],例如,

producer.send(new ProducerRecord<byte[],byte[]>("test", new Payload().toString().getBytes()));

如果您只是将一个有效负载对象传递给生产者,那么最好将键序列化器和值序列化器作为您打算传递的任何内容,并且在读取时您需要从该数据中读取。

最好使用 Serializable 和 ByteArraySerializer/ByteArrayDeserializer。

【讨论】:

  • 反序列化是否有类似的东西(即,不创建自定义反序列化器)?
【解决方案2】:

我们有下面列出的 2 个选项

1) 如果我们打算将自定义 java 对象发送给生产者,我们需要创建一个实现 org.apache.kafka.common.serialization.Serializer 的序列化器,并在创建过程中传递该序列化器类你的制片人

下面的代码参考

public class PayloadSerializer implements org.apache.kafka.common.serialization.Serializer {

    public void configure(Map map, boolean b) {

    }

    public byte[] serialize(String s, Object o) {

       try {
            ByteArrayOutputStream baos = new ByteArrayOutputStream();
            ObjectOutputStream oos = new ObjectOutputStream(baos);
            oos.writeObject(o);
            oos.close();
            byte[] b = baos.toByteArray();
            return b;
        } catch (IOException e) {
            return new byte[0];
        }
    }

    public void close() {

    }
}

并相应地设置值序列化器

<entry key="value.serializer"
                       value="com.spring.kafka.PayloadSerializer" />

2) 无需创建自定义序列化程序类。使用现有的 ByteArraySerializer,但在发送过程中遵循流程

Java 对象 -> 字符串(最好是 JSON 表示而不是 toString)->byteArray

【讨论】:

  • for 2) 反序列化是否有类似的东西,而不是创建自定义反序列化器?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2019-11-03
  • 1970-01-01
  • 2012-04-07
  • 2019-12-28
  • 2016-06-23
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多