Kafka反序列化异常时Offset仍自动提交的问题求助
Kafka反序列化异常时Offset仍自动提交的问题求助
大家好,我最近在测试Kafka的反序列化错误处理逻辑时遇到了一个棘手的问题,想请各位帮忙排查下:
我故意使用错误的反序列化器来测试Kafka对反序列化异常的处理,结果发现即便触发了反序列化错误,Kafka主题的Offset依然会自动提交,完全检测不到消费滞后。我已经尝试配置了ErrorHandlingDeserializer和DefaultErrorHandler,但问题还是没有解决,Offset还是会被自动提交。
以下是我的相关配置和代码:
application.yaml配置
server: port: 9292 spring: kafka: consumer: group-id: consumer-group-1 enable-auto-commit: false bootstrap-servers: localhost:9092 key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer properties: spring.json.type.mapping: com.test.kafka_producer.dto.BankTransferEvent:com.test.kafka_consumer.dto.BankTransferEvent spring.deserializer.value.delegate.class: org.apache.kafka.common.serialization.StringDeserializer spring.json.trusted.packages: com.test.kafka_consumer.dto,com.test.kafka_producer.dto
Kafka消费者配置类
package com.test.kafka_consumer.configuration; import com.test.kafka_consumer.dto.BankTransferEvent; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.listener.ContainerProperties; import org.springframework.kafka.listener.DefaultErrorHandler; import org.springframework.util.backoff.FixedBackOff; @Configuration public class KafkaConsumerConfig { @Bean public DefaultErrorHandler errorHandler() { // 设置重试策略:间隔1秒,重试2次 DefaultErrorHandler errorHandler = new DefaultErrorHandler(new FixedBackOff(1000L, 2L)); // 恢复后不提交Offset errorHandler.setCommitRecovered(false); return errorHandler; } @Bean public ConcurrentKafkaListenerContainerFactory<String, BankTransferEvent> kafkaListenerContainerFactory( ConsumerFactory<String, BankTransferEvent> consumerFactory) { ConcurrentKafkaListenerContainerFactory<String, BankTransferEvent> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); // 启用手动确认模式 factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL); // 给容器绑定错误处理器 factory.setCommonErrorHandler(errorHandler()); return factory; } }
Kafka监听类
package com.test.kafka_consumer.service; import com.test.kafka_consumer.dto.BankTransferEvent; import lombok.extern.slf4j.Slf4j; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.support.Acknowledgment; import org.springframework.stereotype.Service; @Service @Slf4j public class KafkaConsumer { @KafkaListener(topics = "topic-json",groupId = "consumer-group-2", containerFactory = "kafkaListenerContainerFactory") public void consumeJsonEvent1(BankTransferEvent bankTransferEvent, Acknowledgment ack){ log.info("Kafka Consumer aufgerufen"); try { log.info("consumer-JSON-1 consume the event: {" + bankTransferEvent.toString() + "}"); ack.acknowledge(); } catch (Exception e){ log.error("Error bei der Event-Verarbeitung: " + e.getMessage()); } } }
我已经确认enable-auto-commit设为了false,也开启了手动确认模式,并且在错误处理器里禁用了恢复后的Offset提交,但反序列化异常发生时,Offset还是会被自动提交,导致异常消息直接被跳过,既没有重试也没有保留消费滞后。有没有朋友能帮我看看哪里配置出问题了?
内容来源于stack exchange
相关产品推荐
相关产品推荐

