Flink消费Kafka Avro数据时报Expecting type to be a PojoTypeInfo错误求助
问题现象
通过maven-avro插件生成了SpecificRecord类型的Customer类,在Java 8环境下运行Flink消费Kafka Avro格式数据的程序时触发报错,错误信息如下:
Exception in thread "main" java.lang.IllegalStateException: Expecting type to be a PojoTypeInfo [main] INFO org.apache.flink.api.java.typeutils.TypeExtractor - class com.example.Customer does not contain a setter for field first_name [main] INFO org.apache.flink.api.java.typeutils.TypeExtractor - Class class com.example.Customer cannot be used as a POJO type because not all fields are valid POJO fields, and must be processed as GenericType. Please read the Flink documentation on "Data Types & Serialization" for details of the effect on performance.
问题根因
- 你当前avro-maven插件配置中
<createSetters>false</createSetters>,生成的Customer类没有字段对应的setter方法,Flink默认的类型提取器无法将其识别为合法POJO类型 - Flink Scala API的自动类型提取没有正确识别到Customer是Avro SpecificRecord类型,没有匹配到对应的专用序列化器,而是尝试按POJO规则校验,最终抛出异常
- 额外注意:你的POM中maven-compiler-plugin配置里
<source>1.8</source>后面多了无效字符dd,会导致编译异常,需要先删掉这个多余字符
解决方案
方案1(推荐,无需修改生成类结构)
显式指定Flink使用Avro SpecificRecord对应的TypeInformation,避免自动类型提取出错:
- 先导入Avro类型信息类
import org.apache.flink.formats.avro.typeutils.AvroTypeInfo
- 在创建Source前显式声明Customer类的隐式类型信息
// 显式创建Customer类对应的Avro TypeInformation,覆盖默认的类型提取逻辑 implicit val customerTypeInfo: AvroTypeInfo[Customer] = new AvroTypeInfo[Customer](classOf[Customer]) val userKafkaReaderResult = env.addSource(new FlinkKafkaConsumer[Customer]("customer-avro", ConfluentRegistryAvroDeserializationSchema.forSpecific(classOf[Customer],schemaRegistryUrl), properties).setStartFromEarliest())
修改后Flink会直接使用专门的Avro序列化器处理Customer类,不再做POJO规则校验。
方案2(修改生成类配置)
修改POM中avro-maven插件的配置,开启setter生成,重新生成Customer类即可通过Flink的POJO校验:
<plugin> <groupId>org.apache.avro</groupId> <artifactId>avro-maven-plugin</artifactId> <version>${avro.version}</version> <executions> <execution> <phase>generate-sources</phase> <goals> <goal>schema</goal> <goal>protocol</goal> <goal>idl-protocol</goal> </goals> <configuration> <sourceDirectory>${project.basedir}/src/main/resources/avro</sourceDirectory> <stringType>String</stringType> <!-- 把false修改为true,生成setter方法 --> <createSetters>true</createSetters> <enableDecimalLogicalType>true</enableDecimalLogicalType> <fieldVisibility>private</fieldVisibility> </configuration> </execution> </executions> </plugin>
修改后重新执行mvn clean compile生成带setter的Customer类,即可解决POJO校验报错。
内容的提问来源于stack exchange,提问作者codebye
相关产品推荐
相关产品推荐

