Java中使用SpecificAvroSerde出现找不到类的问题求助
解决Kafka Streams中SpecificAvroSerde找不到类的问题
咱们一步步来排查你遇到的SerializationException问题,从最明显的点开始:
1. 首先注意错误日志里的拼写错误!
你看到的错误是:
Caused by: org.apache.kafka.common.errors.SerializationException: Could not find class mybject specified in writer's schema whilst finding reader's schema for a SpecificRecord.
这里的类名是mybject(注意少了一个字母c),但你代码里用的是myObject,生成的类路径里是myobject——这明显是拼写/大小写不匹配问题!
Java是大小写敏感的,而且Avro会严格按照writer schema里的类名去匹配本地生成的类。你需要:
- 检查JDBC Connector使用的Avro Schema(可以通过Schema Registry UI查看对应的subject的schema),确认schema里的
name字段是myObject还是mybject - 打开你生成的Java类文件
target/generated-sources/avro/com/myproject/myclasses/myobject.java,看类定义是public class myobject(全小写)还是public class MyObject(首字母大写),确保代码里的类名和这个完全一致
2. 确认生成的类在运行时类路径中
你提到类是通过Maven插件生成在target/generated-sources/avro下,但这个目录需要被IDE和Maven识别为源码目录,否则编译时不会把这些类打包到最终的jar里,导致运行时找不到:
- 在IDE(比如IntelliJ)中,右键
target/generated-sources/avro目录,选择「Mark Directory as -> Sources Root」 - 检查你的
avro-maven-plugin配置,确保生成的类包路径正确,和你导入的com.myproject.myclasses一致 - 手动执行
mvn clean compile,然后查看target/classes/com/myproject/myclasses下是否存在对应的.class文件(比如myobject.class或MyObject.class)
3. 验证SpecificAvroSerde的配置细节
虽然你的配置看起来没问题,但有几个小细节可以确认:
- 确保
SpecificAvroSerde的schema.registry.url和JDBC Connector使用的是同一个Schema Registry,避免schema不匹配 - 确认你生成的类确实实现了
SpecificRecord接口(Avro Maven插件生成的类默认会实现,可打开类文件确认) - 如果你给全局配置了
DEFAULT_VALUE_SERDE_CLASS_CONFIG为SpecificAvroSerde.class,那在Consumed.with()里其实可以不用单独指定,但单独指定也没问题,只要配置一致就行
修正后的代码示例(假设类名是MyObject)
import com.myproject.myclasses.MyObject; // ... 配置部分 ... Properties config = new Properties(); config.put(StreamsConfig.APPLICATION_ID_CONFIG,this.applicaionId); config.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG,svrConfig.getBootstrapServers()); config.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); config.put(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, svrConfig.getCacheMaxBytesBufferingConfig()); config.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE); config.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.StringSerde.class); config.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, SpecificAvroSerde.class); config.put(KafkaAvroSerializerConfig.SCHEMA_REGISTRY_URL_CONFIG, "http://localhost:8081"); // 流构建部分 final Map<String, String> serdeConfig = Collections.singletonMap("schema.registry.url","http://localhost:8081"); Serdes.StringSerde keySerde = new Serdes.StringSerde(); final Serde<MyObject> valueSpecificAvroSerde = new SpecificAvroSerde<>(); valueSpecificAvroSerde.configure(serdeConfig, false); // false表示这是value serde StreamsBuilder builder = new StreamsBuilder(); KStream<String, MyObject> purchaseInvoices = builder.stream(this.sourceTopic, Consumed.with(keySerde, valueSpecificAvroSerde));
总结
最可能的原因是拼写/大小写不匹配或者生成的类不在运行时类路径中。先排查这两个点,应该就能解决问题。如果还是不行,可以对比Schema Registry里的writer schema和你本地生成类对应的schema,确保两者完全兼容。
内容的提问来源于stack exchange,提问作者Benny Chan
相关产品推荐
相关产品推荐

