如何合并PCollections并保留TaggedOutput?Apache Beam技术问询
Apache Beam(Dataflow)多标签PCollection合并与标签管理问题
我正在为公司的流水线和业务场景开发专属工具包,以简化后续Apache Beam(Dataflow)流水线的开发,核心需求之一是实现失败捕获与重定向。
我有多个使用TaggedOutput的PTransform,有时需要合并这些PTransform的输出,但仍希望保留TaggedOutput。当只有一个TaggedOutput时可实现合并,但合并含多标签的内容时会遇到问题。
单标签场景的可行处理方式
在某PTransform中,我能正常管理标签,代码如下:
fails = [] for timestamp_field in self.timestamp_fields: output, fail = output | f'Convert "{timestamp_field}"' >> ParDo( ConvertTimestamp(self.fail_tag), timestamp_field ).with_outputs(self.fail_tag, main=self.success_tag) # type: ignore fails.append(fail) fails = ( fails | "Flatten all timestamp fails" >> Flatten() | Map(lambda elem: pvalue.TaggedOutput(self.fail_tag, elem)) )
此场景下,我可展平失败输出列表后重新添加原有标签。
多标签合并的问题
但后续我希望该PTransform返回不同类型标签,同时返回已用标签区分的成功与失败输出,尝试合并时遇到了标签丢失的问题:
result = ( ( output, fails, ) | "Flatten output and errors" >> Flatten() ) return result
此场景中,Flatten()会移除所有标签,result丢失标签,且因需应用不同标签,无法复用上述单标签场景的方法。
问题
- 是否存在合并PCollections且不丢失标签的方法?
- 是否可以为PTransform添加标签?(若能在PTransform上使用类似
with_outputs的方法,可重构代码管理标签) - 是否有办法返回结果并区分失败与成功输出?
我当前使用输出标签,若标签方案不可行,也可采用其他方法。
我全程使用Python开发,但若您更擅长Java并能提供清晰解释,我可将Java方案转换为Python实现。
内容的提问来源于stack exchange,提问作者Mike Williamson
相关产品推荐
相关产品推荐

