Apache Beam中GenericRecord结合窗口时AvroRuntimeException的解决方法
解决方案
问题出在AvroCoder默认使用SpecificData处理序列化,而GenericRecord是接口,SpecificData无法为接口生成Schema,必须显式指定使用GenericData来处理GenericRecord的序列化。
修正后的代码
将原代码中的AvroCoder.of(GenericRecord.class, schema)替换为显式指定GenericData的版本:
PCollection<KV<KV<GenericRecord, GenericRecord>, Boolean>> windowedOutput = inputRecords .apply(Window.<KV<KV<GenericRecord, GenericRecord>, Boolean>>into(FixedWindows.of(Duration.standardMinutes(1)))) .setCoder(KvCoder.of( KvCoder.of( // 显式指定使用GenericData处理GenericRecord AvroCoder.of(GenericRecord.class, keySchema, GenericData.get()), AvroCoder.of(GenericRecord.class, valueSchema, GenericData.get()) ), BooleanCoder.of() ));
错误原因分析
从堆栈信息可以看到,错误源于SpecificData.createSchema方法——它试图为GenericRecord接口生成Schema,但SpecificData仅支持Avro Specific类(即继承自SpecificRecord的实体类),不支持GenericRecord接口。
未使用窗口的代码能正常运行,是因为此时编码上下文自动适配了GenericData;而窗口操作改变了PCollection的内部处理路径,导致Beam默认使用SpecificData进行序列化,从而触发异常。
内容的提问来源于stack exchange,提问作者Ayush Sharma
相关产品推荐
相关产品推荐

