删除Kafka主题Avro Schema后无法反序列化消息,该如何解决?
问题描述
已永久删除Kafka主题对应的Avro Schema,现在无法反序列化该主题中的消息,报错信息如下:
错误堆栈
Internal Server Error 发生500错误:请求处理失败;嵌套异常为org.apache.kafka.common.errors.SerializationException: Error retrieving Avro schema for id 100869 堆栈跟踪详情: org.apache.kafka.common.errors.SerializationException: 无法获取ID为100869的Avro Schema Caused by: io.confluent.kafka.schemaregistry.client.rest.exceptions.RestClientException: Schema 100869未找到;错误码:40403 at io.confluent.kafka.schemaregistry.client.rest.RestService.sendHttpRequest(RestService.java:226) at io.confluent.kafka.schemaregistry.client.rest.RestService.httpRequest(RestService.java:252) at io.confluent.kafka.schemaregistry.client.rest.RestService.getId(RestService.java:482) at io.confluent.kafka.schemaregistry.client.rest.RestService.getId(RestService.java:475) at io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient.getSchemaByIdFromRegistry(CachedSchemaRegistryClient.java:153) at io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient.getBySubjectAndId(CachedSchemaRegistryClient.java:232) at io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient.getById(CachedSchemaRegistryClient.java:211) at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer.deserialize(AbstractKafkaAvroDeserializer.java:116) at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer.deserialize(AbstractKafkaAvroDeserializer.java:88) at io.confluent.kafka.serializers.KafkaAvroDeserializer.deserialize(KafkaAvroDeserializer.java:55) at kafdrop.util.AvroMessageDeserializer.deserializeMessage(AvroMessageDeserializer.java:22) at kafdrop.service.KafkaHighLevelConsumer.deserialize(KafkaHighLevelConsumer.java:199) at kafdrop.service.KafkaHighLevelConsumer.lambda$getLatestRecords$3(KafkaHighLevelConsumer.java:132) at java.base/java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:195) at java.base/java.util.ArrayList$SubList$2.forEachRemaining(ArrayList.java:1510) at java.base/java.util.stream.AbstractPipeline.copyInto(AbstractPipeline.java:484) at...
解决方法
1. 恢复Schema到注册表
如果有Schema的备份文件(.avsc格式),直接重新注册到Schema Registry:
- 使用curl命令注册:
注意替换实际的Schema Registry地址、subject名称和Schema内容。curl -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" \ --data '{"schema": "{\"type\": \"record\", \"name\": \"YourRecord\", \"fields\": [{\"name\": \"field1\", \"type\": \"string\"}]}"}' \ http://你的SchemaRegistry地址:8081/subjects/主题对应subject名/versions注:如果原Schema ID是自增生成的,重新注册后无法复用原ID,此时需要结合自定义反序列化器处理。
2. 自定义反序列化器绕过Schema Registry
编写自定义Avro反序列化器,直接加载本地Schema文件进行反序列化,无需从Registry获取:
import org.apache.avro.Schema; import org.apache.avro.generic.GenericDatumReader; import org.apache.avro.generic.GenericRecord; import org.apache.avro.io.DecoderFactory; import org.apache.kafka.common.serialization.Deserializer; import java.io.File; import java.io.IOException; import java.util.Map; public class LocalSchemaAvroDeserializer implements Deserializer<GenericRecord> { private Schema schema; @Override public void configure(Map<String, ?> configs, boolean isKey) { String schemaPath = (String) configs.get("avro.schema.path"); try { schema = new Schema.Parser().parse(new File(schemaPath)); } catch (IOException e) { throw new RuntimeException("加载本地Avro Schema失败", e); } } @Override public GenericRecord deserialize(String topic, byte[] data) { if (data == null) return null; // 跳过Confluent Avro序列化器添加的前5字节(魔术字节+Schema ID) byte[] avroData = new byte[data.length - 5]; System.arraycopy(data, 5, avroData, 0, avroData.length); try { GenericDatumReader<GenericRecord> reader = new GenericDatumReader<>(schema); return reader.read(null, DecoderFactory.get().binaryDecoder(avroData, null)); } catch (IOException e) { throw new RuntimeException("反序列化Avro消息失败", e); } } @Override public void close() {} }
在消费者配置中指定该反序列化器:
key.deserializer=org.apache.kafka.common.serialization.StringDeserializer value.deserializer=com.你的包名.LocalSchemaAvroDeserializer avro.schema.path=/本地Schema文件路径/your-schema.avsc
3. 恢复Schema Registry数据库备份
如果Schema Registry使用外部数据库(如PostgreSQL)或有定期备份:
- 停止Schema Registry服务
- 恢复数据库备份到对应存储位置
- 重启Schema Registry服务,即可找回被删除的Schema ID 100869
4. 提取原始数据尝试推断Schema
如果没有Schema备份,可提取Kafka消息的原始Avro数据,尝试推断Schema:
- 导出原始消息:
kafka-console-consumer.sh --bootstrap-server 你的Kafka地址:9092 --topic 目标主题 --from-beginning --property print.value=false --property print.key=false --property value.deserializer=org.apache.kafka.common.serialization.ByteArrayDeserializer > raw-messages.bin - 提取每个消息前5字节后的二进制内容,使用Avro工具尝试解析或推断Schema:
avro-tools tojson --schema-file inferred-schema.avsc raw-data.bin注:这种方法可能无法完全准确恢复字段名和复杂类型,优先找消息生成方获取原始Schema。
内容的提问来源于stack exchange,提问作者Muhammed Rahiz
相关产品推荐
相关产品推荐

