Spring Kafka消费者Snappy解压失败问题求助
Spring-Kafka消费Snappy压缩消息时解压失败问题排查与解决
问题描述
使用Spring-Kafka消费Kafka消息时,遭遇Snappy解压错误java.io.IOException: FAILED_TO_UNCOMPRESS(5),以及上层抛出的org.apache.kafka.common.KafkaException: Failed to decompress record stream异常,导致消息消费中断。
代码片段
Consumer.java
@Autowired @Qualifier("batchAckEventProcessor") private BatchEventProcessor eventProcessor; @KafkaListener( topics = "MyAckTopic", containerFactory = "kafkaListenerContainerFactory") public void consume(Message<Object> message) { eventProcessor.process(message.getPayload()); }
ConsumerConfig.java
public ConsumerFactory<String, JSONObject> batchConsumerConfiguration() { Map<String, Object> config = new HashMap<>(); config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, env.getProperty("batch_kafka_bootstrap_servers_list")); config.put(ConsumerConfig.GROUP_ID_CONFIG, env.getProperty("batch_consumer_group_id")); config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class); return new DefaultKafkaConsumerFactory<>( config, new StringDeserializer(), new JsonDeserializer<>(JSONObject.class)); } @Bean public ConcurrentKafkaListenerContainerFactory<String, JSONObject> kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, JSONObject> factory = new ConcurrentKafkaListenerContainerFactory(); factory.setConsumerFactory(batchConsumerConfiguration()); return factory; }
环境信息
- Topic压缩类型:snappy
- Spring Boot版本:2.1.1.RELEASE
- 核心依赖版本:
[INFO] +- org.springframework.kafka:spring-kafka:jar:2.2.15.RELEASE:compile [INFO] | \- org.apache.kafka:kafka-clients:jar:2.0.1:compile [INFO] | +- org.lz4:lz4-java:jar:1.4.1:compile [INFO] | \- org.xerial.snappy:snappy-java:jar:1.1.7.1:compile
错误日志
2022-12-20 18:40:01 18:39:58 20-Dec [ERROR] [o.s.kafka.listener.LoggingErrorHandler ] [handle :37 ] [] []-> Error while processing: null 2022-12-20 18:40:01 Caused by: org.apache.kafka.common.KafkaException: Failed to decompress record stream 2022-12-20 18:40:01 org.apache.kafka.common.KafkaException: Received exception when fetching the next record from MyAckTopic-0. If needed, please seek past the record to continue consumption. 2022-12-20 18:40:01 at org.apache.kafka.common.record.DefaultRecordBatch$1.readNext(DefaultRecordBatch.java:268) 2022-12-20 18:40:01 at org.apache.kafka.common.record.DefaultRecordBatch$RecordIterator.next(DefaultRecordBatch.java:563) 2022-12-20 18:40:01 at org.apache.kafka.clients.consumer.internals.Fetcher$PartitionRecords.fetchRecords(Fetcher.java:1204) 2022-12-20 18:40:01 at org.apache.kafka.common.record.DefaultRecordBatch$RecordIterator.next(DefaultRecordBatch.java:532) 2022-12-20 18:40:01 at org.apache.kafka.clients.consumer.internals.Fetcher$PartitionRecords.access$1400(Fetcher.java:1072) 2022-12-20 18:40:01 at org.apache.kafka.clients.consumer.internals.Fetcher$PartitionRecords.nextFetchedRecord(Fetcher.java:1183) 2022-12-20 18:40:01 at org.apache.kafka.clients.consumer.internals.Fetcher.fetchRecords(Fetcher.java:562) 2022-12-20 18:40:01 at org.apache.kafka.clients.consumer.internals.Fetcher$PartitionRecords.fetchRecords(Fetcher.java:1218) 2022-12-20 18:40:01 at org.apache.kafka.clients.consumer.internals.Fetcher.fetchedRecords(Fetcher.java:523) 2022-12-20 18:40:01 at org.apache.kafka.clients.consumer.internals.Fetcher$PartitionRecords.access$1400(Fetcher.java:1072) 2022-12-20 18:40:01 at org.apache.kafka.clients.consumer.KafkaConsumer.pollForFetches(KafkaConsumer.java:1230) 2022-12-20 18:40:01 at org.apache.kafka.clients.consumer.internals.Fetcher.fetchRecords(Fetcher.java:562) 2022-12-20 18:40:01 at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1187) 2022-12-20 18:40:01 at org.apache.kafka.clients.consumer.internals.Fetcher.fetchedRecords(Fetcher.java:523) 2022-12-20 18:40:01 at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1154) 2022-12-20 18:40:01 at org.apache.kafka.clients.consumer.KafkaConsumer.pollForFetches(KafkaConsumer.java:1259) 2022-12-20 18:40:01 at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:757) 2022-12-20 18:40:01 ... 7 common frames omitted 2022-12-20 18:40:01 Caused by: java.io.IOException: FAILED_TO_UNCOMPRESS(5) 2022-12-20 18:40:01 at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:708) 2022-12-20 18:40:01 at org.xerial.snappy.SnappyNative.throw_error(SnappyNative.java:98) 2022-12-20 18:40:01 at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) 2022-12-20 18:40:01 at org.xerial.snappy.SnappyNative.rawUncompress(Native Method) 2022-12-20 18:40:01 at java.util.concurrent.FutureTask.run(FutureTask.java:266) 2022-12-20 18:40:01 at org.xerial.snappy.Snappy.rawUncompress(Snappy.java:474) 2022-12-20 18:40:01 at java.lang.Thread.run(Thread.java:748) 2022-12-20 18:40:01 at org.xerial.snappy.Snappy.uncompress(Snappy.java:513) 2022-12-20 18:40:01 Caused by: org.apache.kafka.common.KafkaException: Failed to decompress record stream 2022-12-20 18:40:01 at org.xerial.snappy.SnappyInputStream.hasNextChunk(SnappyInputStream.java:439) 2022-12-20 18:40:01 at org.apache.kafka.common.record.DefaultRecordBatch$1.readNext(DefaultRecordBatch.java:268) 2022-12-20 18:40:01 at org.xerial.snappy.SnappyInputStream.read(SnappyInputStream.java:466) 2022-12-20 18:40:01 at java.io.DataInputStream.readByte(DataInputStream.java:265) 2022-12-20 18:40:01 at org.apache.kafka.common.utils.ByteUtils.readVarint(ByteUtils.java:168) 2022-12-20 18:40:01 at org.apache.kafka.common.record.DefaultRecord.readFrom(DefaultRecord.java:292) 2022-12-20 18:40:01 at org.apache.kafka.common.record.DefaultRecordBatch$1.readNext(DefaultRecordBatch.java:264) 2022-12-20 18:40:01 ... 15 common frames omitted
解决方案
1. 升级Snappy依赖版本
当前使用的snappy-java:1.1.7.1存在已知解压兼容性问题,建议升级到1.1.8.4及以上版本。在pom.xml中显式指定版本:
<dependency> <groupId>org.xerial.snappy</groupId> <artifactId>snappy-java</artifactId> <version>1.1.8.4</version> </dependency>
2. 配置错误处理器跳过损坏消息
日志提示可跳过异常记录继续消费,通过配置SeekToCurrentErrorHandler实现自动跳过损坏消息:
@Bean public ConcurrentKafkaListenerContainerFactory<String, JSONObject> kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, JSONObject> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(batchConsumerConfiguration()); // 设置错误处理器,自动跳过异常消息 factory.setErrorHandler(new SeekToCurrentErrorHandler()); return factory; }
注:Spring Boot 2.1.x中SeekToCurrentErrorHandler默认最多重试10次,重试失败后会跳过该消息。
3. 验证生产者与消费者压缩格式一致性
确认生产者端是否正确配置了snappy压缩,避免出现生产者使用其他压缩算法但topic属性显示snappy的情况。生产者配置示例:
props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "snappy");
4. 检查Kafka版本兼容性
确保Kafka集群版本与客户端版本(kafka-clients:2.0.1)兼容,避免因版本差异导致的压缩格式不兼容问题。
内容的提问来源于stack exchange,提问作者AnV
相关产品推荐
相关产品推荐

