【问题标题】:Getting Serialization exception while consuming from Kafka queue using Spring Boot使用 Spring Boot 从 Kafka 队列消费时获取序列化异常
【发布时间】:2019-05-26 02:47:32
【问题描述】:

我正在使用 Spring Boot 和 Kafka 来使用来自 Kafka 队列的消息。

但是,如果 Jackson 解析有效载荷时出现任何错误。

消息一直卡住,不断重新尝试消费,一直解析异常。

我尝试在 Kafka 配置中使用 ErrorHandlingDeserializer2 并映射错误处理程序,但问题仍然存在。

@Configuration
@EnableKafka
public class KafkaConfig 
{

    @Value("${spring.kafka.bootstrap-servers}")
    private String bootstrapServers;

    @Bean
    public Map<String, Object> consumerConfigs() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);

        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer2.class);
        props.put(ErrorHandlingDeserializer2.KEY_DESERIALIZER_CLASS, StringDeserializer.class);

        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer2.class);
        props.put(ErrorHandlingDeserializer2.VALUE_DESERIALIZER_CLASS, JsonDeserializer.class);

        props.put(ConsumerConfig.GROUP_ID_CONFIG, "${connecto.group-id}");
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        props.put(JsonSerializer.ADD_TYPE_INFO_HEADERS, false);

        return props;
    }

    @Bean
    public ConsumerFactory<String, UserDto> consumerFactory() {
        return new DefaultKafkaConsumerFactory<>(
                consumerConfigs(),
                new StringDeserializer(),
                new JsonDeserializer<>(UserDto.class));
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, UserDto> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, UserDto> factory =
                new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());

        factory.setErrorHandler(new MosaicKafkaErrorHandler());

        return factory;
    }

    @Bean
    public KafkaTemplate kafkaTemplate() 
    {
        return new KafkaTemplate<>(producerFactory());
    }

    @Bean
    public ProducerFactory<String, UserSearchDto> producerFactory() 
    {
        Map<String, Object> config = new HashMap<>();
        config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        config.put(JsonSerializer.ADD_TYPE_INFO_HEADERS, false);
        config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
        return new DefaultKafkaProducerFactory<>(config);
    }
}

public class MosaicKafkaErrorHandler implements ContainerAwareErrorHandler{

@Override
public void handle(Exception thrownException, List<ConsumerRecord<?, ?>> records, Consumer<?, ?> consumer,
        MessageListenerContainer container) {
    thrownException.printStackTrace();

    }

}

我希望以这样一种方式实现它,即每当发生配对异常时,它应该只记录该消息并继续下一条记录,而不是陷入困境。

以下是错误-

Caused by: com.fasterxml.jackson.databind.exc.MismatchedInputException: Cannot construct instance of `com.gxxx.mxxx.common.dto.UserDetailsDto` (although at least one Creator exists): no String-argument constructor/factory method to deserialize from String value ('{"field1":20000,"field2":20000,"type":""}')

at org.apache.kafka.clients.consumer.KafkaConsumer.pollForFetches(KafkaConsumer.java:1230)
at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1187)
at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1154)
at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:741)
at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:698)
at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
at java.util.concurrent.FutureTask.run(FutureTask.java:266)
at java.lang.Thread.run(Thread.java:748)
org.apache.kafka.common.errors.SerializationException: Error deserializing key/value for partition kafakaqueue.identify-1 at offset 326355379. If needed, please seek past the record to continue consumption.

【问题讨论】:

  • 我们遇到了同样的问题。解决方案取决于您的 kafka 版本。请告诉我您使用的是哪个版本。
  • 显示此MosaicKafkaErrorHandler 代码
  • public class MosaicKafkaErrorHandler implements ContainerAwareErrorHandler{ @Override public void handle(Exception thrownException, List&lt;ConsumerRecord&lt;?, ?&gt;&gt; records, Consumer&lt;?, ?&gt; consumer, MessageListenerContainer container) { thrownException.printStackTrace(); } }
  • @Deadpool - 只记录错误
  • 你能更新帖子@DhruvSaksena中的代码吗

标签: java spring-boot spring-kafka


【解决方案1】:

为了解决这个问题,version 2.2 引入了ErrorHandlingDeserializer2。这个反序列化器委托给一个真正的反序列化器(键或值)。如果委托未能反序列化记录内容,则ErrorHandlingDeserializer2 在包含原因和原始字节的标头中返回空值和DeserializationException。当您使用记录级别 MessageListener 时,如果 ConsumerRecord 包含键或值的 DeserializationException 标头,则使用失败的 ConsumerRecord 调用容器的 ErrorHandler。记录不会传递给监听器。

但是,由于您使用的是早期版本,因此您可以使用以下 hack。

     @Override
        public void handle(Exception thrownException, List<ConsumerRecord<?, ?>> records, Consumer<?, ?> consumer,
                MessageListenerContainer container) {
            thrownException.printStackTrace();
            if (thrownException instanceOf SerializationException){
                String s = thrownException.getMessage().split("Error deserializing key/value for partition ")[1].split(". If needed, please seek past the record to continue consumption.")[0];
                String topics = s.split("-")[0];
                int offset = Integer.valueOf(s.split("offset ")[1]);
                int partition = Integer.valueOf(s.split("-")[1].split(" at")[0]);

                TopicPartition topicPartition = new TopicPartition(topics, partition);
                consumer.seek(topicPartition, offset + 1);  
               }
            }

【讨论】:

  • 我尝试了您的解决方案,它已到达处理程序,但记录集合为空-
  • java.lang.IndexOutOfBoundsException:索引:0 at java.util.Collections$EmptyList.get(Collections.java:4454) ~[na:1.8.0_212] at com.girnarsoft.mosaic.config .MosaicKafkaErrorHandler.handle(MosaicKafkaErrorHandler.java:22) ~[classes!/:0.0.1-SNAPSHOT] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.handleConsumerException(KafkaMessageListenerContainer.java:852) [spring-kafka-2.2 .5.RELEASE.jar!/:2.2.5.RELEASE] 在
  • @DhruvSaksena 我已经修改了代码。它看起来有点hacky,但你能试一试吗。
  • 当然我也这样做了。请继续进行更改。
  • 完成。再次感谢!!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-07-10
  • 2022-01-04
  • 1970-01-01
  • 2016-03-23
  • 2011-02-17
  • 1970-01-01
  • 2015-01-26
相关资源
最近更新 更多