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

Kafka Avro消息生产成功但消费遇EOFException求助

Kafka Avro反序列化EOFException排查求助

问题现象

Kafka Avro消息生产成功,但消费端反序列化部分消息时抛出EOFException。该Topic存在多生产者,多数消息可正常消费,仅部分消息无法处理。

已排查内容

  • 确认Schema ID 13212无异常,本地调试可正常获取该Schema
  • 使用StringDeserializer验证对应消息非空

POM依赖

<dependency>
    <groupId>io.confluent</groupId>
    <artifactId>kafka-json-serializer</artifactId>
    <version>4.1.1</version>
    <exclusions>
        <exclusion>
            <groupId>org.slf4j</groupId>
            <artifactId>slf4j-log4j12</artifactId>
        </exclusion>
    </exclusions>
</dependency>

<dependency>
    <groupId>io.confluent</groupId>
    <artifactId>kafka-avro-serializer</artifactId>
    <version>5.3.0</version>
    <exclusions>
        <exclusion>
            <groupId>org.slf4j</groupId>
            <artifactId>slf4j-log4j12</artifactId>
        </exclusion>
    </exclusions>
</dependency>

异常栈信息

2024-05-09 10:03:31.702 ERROR [main] - [SchemaConsumerTest.java:100] - Error deserializing key/value for partition xxxx.topic-3 at offset 480960. If needed, please seek past the record to continue consumption.
org.apache.kafka.common.errors.RecordDeserializationException: Error deserializing key/value for partition xxxx.topic-3 at offset 480960. If needed, please seek past the record to continue consumption.
    at org.apache.kafka.clients.consumer.internals.CompletedFetch.parseRecord(CompletedFetch.java:331)
    at org.apache.kafka.clients.consumer.internals.CompletedFetch.fetchRecords(CompletedFetch.java:283)
    at org.apache.kafka.clients.consumer.internals.FetchCollector.fetchRecords(FetchCollector.java:168)
    at org.apache.kafka.clients.consumer.internals.FetchCollector.collectFetch(FetchCollector.java:134)
    at org.apache.kafka.clients.consumer.internals.Fetcher.collectFetch(Fetcher.java:145)
    at org.apache.kafka.clients.consumer.internals.LegacyKafkaConsumer.pollForFetches(LegacyKafkaConsumer.java:693)
    at org.apache.kafka.clients.consumer.internals.LegacyKafkaConsumer.poll(LegacyKafkaConsumer.java:617)
    at org.apache.kafka.clients.consumer.internals.LegacyKafkaConsumer.poll(LegacyKafkaConsumer.java:585)
    at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:827)
    at com.citi.gsp.kafka.consumer.SchemaConsumerTest.receiveMsg(SchemaConsumerTest.java:92)
    at com.citi.gsp.kafka.consumer.SchemaConsumerTest.testConsumer(SchemaConsumerTest.java:83)
    at com.citi.gsp.kafka.consumer.SchemaConsumerTest.main(SchemaConsumerTest.java:55)
Caused by: org.apache.kafka.common.errors.SerializationException: Error deserializing Avro message for id 13212
    at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer.deserialize(AbstractKafkaAvroDeserializer.java:156)
    at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer.deserialize(AbstractKafkaAvroDeserializer.java:79)
    at io.confluent.kafka.serializers.KafkaAvroDeserializer.deserialize(KafkaAvroDeserializer.java:55)
    at org.apache.kafka.common.serialization.Deserializer.deserialize(Deserializer.java:62)
    at org.apache.kafka.common.serialization.Deserializer.deserialize(Deserializer.java:73)
    at org.apache.kafka.clients.consumer.internals.CompletedFetch.parseRecord(CompletedFetch.java:321)
    ... 11 common frames omitted
Caused by: java.io.EOFException: null
    at org.apache.avro.io.BinaryDecoder.ensureBounds(BinaryDecoder.java:514)
    at org.apache.avro.io.BinaryDecoder.readInt(BinaryDecoder.java:155)
    at org.apache.avro.io.BinaryDecoder.readIndex(BinaryDecoder.java:465)
    at org.apache.avro.io.ResolvingDecoder.readIndex(ResolvingDecoder.java:282)
    at org.apache.avro.generic.GenericDatumReader.readWithoutConversion(GenericDatumReader.java:188)
    at org.apache.avro.generic.GenericDatumReader.read(GenericDatumReader.java:161)
    at org.apache.avro.generic.GenericDatumReader.readField(GenericDatumReader.java:260)
    at org.apache.avro.generic.GenericDatumReader.readRecord(GenericDatumReader.java:248)
    at org.apache.avro.generic.GenericDatumReader.readWithoutConversion(GenericDatumReader.java:180)
    at org.apache.avro.generic.GenericDatumReader.read(GenericDatumReader.java:161)
    at org.apache.avro.generic.GenericDatumReader.read(GenericDatumReader.java:154)
    at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer.deserialize(AbstractKafkaAvroDeserializer.java:125)
    ... 16 common frames omitted

可能原因及解决办法

1. 生产者端序列化异常(消息体不完整)

  • 原因:多生产者场景下,部分生产者未正确使用Confluent的KafkaAvroSerializer,或序列化过程中出现中断(如IO异常、内存不足),导致写入Kafka的消息体缺失部分Avro二进制数据,触发消费端EOF校验失败。
  • 解决:
    • 检查所有生产者代码,确保统一使用KafkaAvroSerializer,禁止自定义Avro序列化逻辑;
    • 定位异常消息对应的生产者实例,排查其序列化环节是否存在未捕获的异常;
    • 校验生产者输出的二进制数据格式,确认符合Confluent Avro规范:1字节魔术位 + 4字节Schema ID + 完整Avro二进制数据。

2. 依赖版本不兼容

  • 原因:kafka-json-serializer(4.1.1)与kafka-avro-serializer(5.3.0)版本跨度较大,Confluent组件跨版本混用可能导致序列化/反序列化逻辑不兼容,引发部分消息解析失败。
  • 解决:
    • 将所有Confluent相关依赖统一到同一稳定版本(建议升级至5.3.0或更高);
    • 升级后重新部署并测试消费逻辑,确认异常是否消失。

3. Kafka消息损坏

  • 原因:Kafka Broker在存储或传输过程中出现消息损坏,导致消费端读取的消息长度不足,触发EOFException。
  • 解决:
    • 使用Kafka命令行工具导出异常偏移量的二进制消息:
      kafka-console-consumer.sh --bootstrap-server <broker地址> --topic <topic名称> --offset 480960 --partition 3 --value-deserializer org.apache.kafka.common.serialization.ByteArrayDeserializer --max-messages 1
      
    • 分析二进制数据结构是否符合Confluent Avro格式,若确认损坏,通过消费端seek跳过该偏移量,同时检查Kafka集群磁盘健康、网络稳定性。

4. 反序列化器配置问题

  • 原因:消费端未正确配置Schema Registry地址,或Schema兼容性规则设置不当,导致部分消息解析失败。
  • 解决:
    • 确认消费端schema.registry.url配置正确,且能正常访问Schema Registry;
    • 检查Schema Registry中该Schema的兼容性设置,确保消费端使用的Schema与生产者兼容。

内容的提问来源于stack exchange,提问作者amulet

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 14:40:33