【问题标题】:Kafka Listener method could not be invoked with the incoming message无法使用传入消息调用 Kafka 侦听器方法
【发布时间】:2020-01-17 02:12:56
【问题描述】:

我通过使用 Spring Boot 应用程序将 JSON 数组转换为 Kafka Producer 中的 toString() 来发送 JSON 数组,但在 Consumer 中出现以下错误:

org.springframework.kafka.listener.ListenerExecutionFailedException: 无法使用传入消息调用侦听器方法 端点处理程序详细信息: 方法 [public void com.springboot.service.KafkaReciever.recieveData(com.springboot.model.Student,java.lang.String) 抛出 java.lang.Exception] 豆 [com.springboot.service.KafkaReciever@5bb3d42d];嵌套异常是 org.springframework.messaging.converter.MessageConversionException: 无法处理消息;嵌套异常是 org.springframework.messaging.converter.MessageConversionException: 无法从 [java.lang.String] 转换为 [com.springboot.model.Student] 用于 GenericMessage [有效载荷=[com.springboot.model.Student@5e40dc31, com.springboot.model.Student@235e68b8], headers={kafka_offset=45, kafka_receivedMessageKey=null, kafka_receivedPartitionId=0, kafka_receivedTopic=myTopic-kafkasender}], 失败消息=通用消息 [有效载荷=[com.springboot.model.Student@5e40dc31, com.springboot.model.Student@235e68b8], headers={kafka_offset=45, kafka_receivedMessageKey=null, kafka_receivedPartitionId=0, kafka_receivedTopic=myTopic-kafkasender}]

配置文件:

@Configuration
@EnableKafka
public class KafkaConsumerConfig {

    @Value("${kafka.boot.server}")
    private String kafkaServer;

    @Value("${kafka.consumer.group.id}")
    private String kafkaGroupId;

    @Bean
    public ConsumerFactory<String, String> consumerConfig() {

         Properties props = new Properties();

         props.put("bootstrap.servers", "localhost:9092");
         props.put(ConsumerConfig.GROUP_ID_CONFIG, kafkaGroupId);
         props.put("message.assembler.buffer.capacity", 33554432);
         props.put("max.tracked.messages.per.partition", 24);
         props.put("exception.on.message.dropped", true);
         props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
         props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
         props.put("segment.deserializer.class", DefaultSegmentDeserializer.class.getName());

         return new DefaultKafkaConsumerFactory(props, null, new StringDeserializer());
    }

    @Bean
    public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String> listener = new ConcurrentKafkaListenerContainerFactory<>();
        listener.setConsumerFactory(consumerConfig());
        return listener;
    }
}

接收文件:

@Service
public class KafkaReciever {

    private static final Logger LOGGER = LoggerFactory.getLogger(KafkaReciever.class);

    @KafkaListener(topics = "${kafka.topic.name}", group = "${kafka.consumer.group.id}")
    public void recieveData(@Payload Student student, @Header(KafkaHeaders.MESSAGE_KEY) String messageKey) throws Exception{
        LOGGER.info("Data - " + student + " recieved");
    }
}

POST json:

 [{
        "studentId": "Q45678123",
        "firstName": "Anderson",
        "lastName": "John",
        "age": "12",
        "address": {
          "apartment": "apt 123",
          "street": "street Info",
          "state": "state",
          "city": "city",
          "postCode": "12345"
        }
    },
    {
        "studentId": "Q45678123",
        "firstName": "abc",
        "lastName": "xyz",
        "age": "12",
        "address": {
          "apartment": "apt 123",
          "street": "street Info",
          "state": "state",
          "city": "city",
          "postCode": "12345"
        }
    }]

我得到以下消费者输出:

[com.springboot.model.Student@5e40dc31, com.springboot.model.Student@235e68b8]

【问题讨论】:

  • 您的消费者正在等待字符串(反序列化器),但它正在接收一个学生。您可以访问生产者的序列化程序吗?你需要这样的东西:stackoverflow.com/questions/40154086/…
  • 我的模型数据是:public class Student implements Serializable { private static final long serialVersionUID = 1L;私人字符串学生ID;私人字符串名;私人字符串姓氏;私有字符串年龄;私人地址地址;获取/设置}
  • 我将 VALUE_DESERIALIZER_CLASS_CONFIG 更改为 JsonDeserializer 然后我收到以下错误:org.apache.kafka.common.errors.SerializationException: Error deserializing key/value for partition myTopic-kafkasender-0 at offset 47 Caused by: com .fasterxml.jackson.databind.JsonMappingException:无法反序列化 com.springboot.model.Student 的实例 [来源:[B@304eaefb;行:1,列:1]
  • ConsumerFactory 请尝试使用您的有效负载对象类型,例如如果您的班级名称为 student 那么 ConsumerFactory ____________________ props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);如果您需要适当的 ans ping 我,您可以使用 json desrialzer 我会帮助您
  • @harkesh 谢谢.. 添加这个得到相同的错误.. 但在服务中更改为 '@Payload List student' 错误已解决但数据未转换为 json 数组。

标签: spring spring-boot apache-kafka kafka-consumer-api


【解决方案1】:

无法从 START_ARRAY 中反序列化 com.springboot.model.Student 的实例

如果使用 json 脱轨器,你有一个列表,而不是一个学生

@Payload List<Student> student

或者如果使用字符串脱轨器,你有一个 JSON 字符串,你必须手动解析它

receiveData(@Payload String student ... ) { 
    JsonNode data = new ObjectMapper().readTree(student); // for example, but should extract ObjectMapper to a field
}

关于您的其他输出,请参阅How do I print my Java object without getting "SomeType@2f92e0f4"?

【讨论】:

  • 谢谢。我已经在生产者中解析了 JSON 数据并将其发送给消费者,之后我在消费者中获得了正确的 json 数据。只是想检查一下,我这样做是否正确?
  • 只要没有例外,我就说没有错
猜你喜欢
  • 2019-02-25
  • 2018-01-15
  • 2022-01-19
  • 1970-01-01
  • 1970-01-01
  • 2022-06-13
  • 1970-01-01
  • 2014-07-09
  • 1970-01-01
相关资源
最近更新 更多