如何在不修改DoFn的process方法下,将Apache Beam输入与输出组合成元组?
解决Apache Beam中不修改现有DoFn实现输出元组的问题
你可以通过以下两种方式,在不修改LolString类代码的前提下,实现输出包含原字符串和处理后字符串的元组:
方案一:包装新的DoFn类
创建一个新的DoFn,内部调用原LolString的process方法,同时返回原元素与处理结果组成的元组:
class LolStringWithOriginal(apache_beam.DoFn): def process(self, element: str): processed_str = LolString().process(element) yield (element, processed_str)
在管道中替换为这个新的DoFn:
strings | "Lol string with original" >> apache_beam.ParDo(LolStringWithOriginal())
方案二:使用Map转换
直接通过Map函数,在lambda表达式中调用原LolString的处理逻辑,生成元组:
strings | "Generate lol tuple" >> apache_beam.Map(lambda x: (x, LolString().process(x)))
两种方式都能让管道输出类似("Stack", "Stack_lol")的元组结果,且完全不需要改动原LolString类的代码。
内容的提问来源于stack exchange,提问作者Rafaél
相关产品推荐
相关产品推荐

