【问题标题】:How to send Custom Object to Kafka Topic with Producer如何使用 Producer 将自定义对象发送到 Kafka 主题
【发布时间】:2018-09-21 20:04:22
【问题描述】:

我想将我的带有 Producer 的 Account 类发送到我的 Kafka 主题,然后我将与 Kafka Stream 聚合。但是,我无法发送我收到错误的对象:

原因:org.apache.kafka.common.KafkaException:bank.Account 不是 org.apache.kafka.common.serialization.Serializer 的实例

我的生产者类:

 public static void main(String[] args) {

        DataAccess dataAccess = new DataAccess();
        List<Account> accountList = dataAccess.read();

        final Logger logger = LoggerFactory.getLogger(Producer.class);
        Properties properties = new Properties();

        properties.setProperty(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,"127.0.0.1:9092");
        properties.setProperty(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,LongSerializer.class.getName());
        properties.setProperty(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,Account.class.getName());


        KafkaProducer<Long,Account> producer = new KafkaProducer<Long, Account>(properties);



        for (Account account : accountList) {

            ProducerRecord<Long,Account> record = new ProducerRecord<Long, Account>("bank_account",account.getFromId(),account);


            producer.send(record, new Callback() {
                public void onCompletion(RecordMetadata recordMetadata, Exception e) {
                    if (e == null) {
                        logger.info("Record sent successfully. \n "+ "Topic : "+recordMetadata.topic() +"\n"+
                                "Partition : " + recordMetadata.partition() + "\n"+
                                "Offset : " +recordMetadata.offset() +"\n"+
                                "Timestamp: " +recordMetadata.timestamp() +"\n");
                        try {
                            Thread.sleep(1000);
                        } catch (InterruptedException e1) {
                            e1.printStackTrace();
                        }

                    }
                    else{
                        logger.info("Error sending producer");
                    }
                }
            });
        }


        producer.flush();
        producer.close();
    }

这行出错了:

KafkaProducer<Long,Account> producer = new KafkaProducer<Long, Account>(properties);

我的帐户类:

public class Account {

    private long fromId;
    private long amount;
    private long toId;
    private ZonedDateTime time;
}

所以我的问题是,我们如何将自定义对象发送到 kafka 主题?之后,我当然想使用该消息。

【问题讨论】:

    标签: java serialization apache-kafka kafka-producer-api


    【解决方案1】:

    这一行

    properties.setProperty(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,Account.class.getName());

    您必须实现自己的Serializer class。它不能是一个普通的类。


    有些人使用 JSON 进行序列化,有些人使用 Avro 或 Protobuf。但是,您将数据放入byte[] 只是一个实现细节。

    【讨论】:

    • 感谢回答的朋友。这是正确的,现在我可以将我的对象发送给生产者。我现在卡在 kafkastreams 中。我不能 Serde 我的对象,也不能实时流式传输。如果你有时间,请看看这个:stackoverflow.com/questions/52455670/…
    【解决方案2】:
    //1 
    prop.setProperty(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
            prop.setProperty(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, Employee.class.getName());
    
    //2
     KafkaProducer<String, Employee> producer = new KafkaProducer(prop);
    
            Employee emp = new Employee(1, "Arun");
    
            ProducerRecord prodRecord = new ProducerRecord("aryan_topic", emp);
    //3
    import org.apache.kafka.common.header.Headers;
    import org.apache.kafka.common.serialization.Serializer;
    
    import java.io.Serializable;
    import java.util.Map;
    
    //Developed by Arun Singh
    public class Employee implements Serializable, Serializer {
        Integer empId;
        String empName;
        Address add;
    
        public Employee() {
        }
    
        public Employee(Integer empId, String empName, Address add) {
            this.empId = empId;
            this.empName = empName;
            this.add = add;
        }
    
        public Integer getEmpId() {
            return empId;
        }
    
        public String getEmpName() {
            return empName;
        }
    
        public Address getAdd() {
            return add;
        }
    
        public void setEmpId(Integer empId) {
            this.empId = empId;
        }
    
        public void setEmpName(String empName) {
            this.empName = empName;
        }
    
        public void setAdd(Address add) {
            this.add = add;
        }
    
        public void configure(Map configs, boolean isKey) {
    
        }
    
        public byte[] serialize(String s, Object o) {
            return new byte[0];
        }
    
        public byte[] serialize(String topic, Headers headers, Object data) {
            return new byte[0];
        }
    
        public void close() {
    
        }
    }
    

    【讨论】:

      猜你喜欢
      • 2017-04-29
      • 2019-11-03
      • 2018-10-27
      • 2019-02-14
      • 2020-05-26
      • 2020-08-18
      • 1970-01-01
      • 2021-11-25
      • 2019-04-23
      相关资源
      最近更新 更多