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集群磁盘健康、网络稳定性。
- 使用Kafka命令行工具导出异常偏移量的二进制消息:
4. 反序列化器配置问题
- 原因:消费端未正确配置Schema Registry地址,或Schema兼容性规则设置不当,导致部分消息解析失败。
- 解决:
- 确认消费端
schema.registry.url配置正确,且能正常访问Schema Registry; - 检查Schema Registry中该Schema的兼容性设置,确保消费端使用的Schema与生产者兼容。
- 确认消费端
内容的提问来源于stack exchange,提问作者amulet
相关产品推荐
相关产品推荐

