You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.07 19:05:20