Apache Beam多PCollection输出时MultiOutputReceiver类型不兼容错误排查
问题原因及解决方案
1. 类型不兼容报错的核心诱因
你遇到的logSchema cannot be converted to capture#1 of ?报错,本质是泛型类型擦除导致的:
TupleTagList存储的是无具体类型的TupleTag<?>实例,你通过tags.get(index)获取到的tag携带通配符类型,MultiOutputReceiver.get(tag)拿到的输出接收器也无法识别具体接收的元素类型,编译器无法确认你输出的logSchema类型是否匹配,就会抛出类型不兼容错误。- 你后续的强制转换警告,是因为你定义
logObjects时没有指定泛型,写了原生类型PCollection logObjects而非PCollection<logSchema> logObjects,泛型擦除后编译器无法确认apply返回的结果是PCollectionTuple,才要求你手动强转。
2. 具体修复步骤
步骤1:修改DoFn传参逻辑,不要传TupleTagList
直接将具体的TupleTag<logSchema>实例作为branching类的成员变量,不要用索引从TupleTagList取tag,保证类型明确:
static class branching extends DoFn<logSchema, logSchema> { private final TupleTag<logSchema> noticesTag; private final TupleTag<logSchema> errorsTag; private final TupleTag<logSchema> warningsTag; private final TupleTag<logSchema> soutTag; public branching(TupleTag<logSchema> noticesTag, TupleTag<logSchema> errorsTag, TupleTag<logSchema> warningsTag, TupleTag<logSchema> soutTag) { this.noticesTag = noticesTag; this.errorsTag = errorsTag; this.warningsTag = warningsTag; this.soutTag = soutTag; } @ProcessElement public void processElement(@Element logSchema log, MultiOutputReceiver out ) { if (log.getType().equals("[notice]")) out.get(noticesTag).output(log); else if (log.getType().equals("[error]")) out.get(errorsTag).output(log); else if (log.getType().equals("[warn]")) out.get(warningsTag).output(log); else if (log.getType().equals("[sout]") ) out.get(soutTag).output(log); } }
步骤2:补上PCollection的泛型声明,去掉不必要的强转
修改run方法中logObjects的定义,指定泛型为logSchema:
// 补全泛型,不要写原生PCollection PCollection<logSchema> logObjects = input .apply("Conform", ParDo.of(new conformToSchema())); // 无需强制转换,编译器可识别返回类型为PCollectionTuple PCollectionTuple multipleOutputs = logObjects.apply("Branch", ParDo.of(new branching(noticesTag, errorsTag, warningsTag, soutTag)) .withOutputTags(all, tags));
额外优化建议
你的主输出all这个TupleTag实际上没有用到,因为你所有的输出都走了额外输出tag,如果不需要这个主输出,可以把branching类的父类定义改成DoFn<logSchema, Void>,主输出tag声明为TupleTag<Void>即可,避免冗余。
内容的提问来源于stack exchange,提问作者Pelleri
相关产品推荐
相关产品推荐

