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

删除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命令注册:
    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 Registry地址、subject名称和Schema内容。

    注:如果原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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 11:12:23