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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 07:58:17