为何在Apache Beam ParDo函数中需显式返回列表?
Apache Beam ParDo中显式返回列表的原因解析
问题场景
我编写了一段处理CSV文件的Apache Beam代码:
class SplitRow(beam.DoFn): def process(self, element): return [element.split(',')] class FilterCardioPatients(beam.DoFn): def process(self, element): if element[3] == 'cardio': return [element] class PairPatients(beam.DoFn): def process(self, element): return [(element[1], 1)] class Counting(beam.DoFn): def process(self, element): (key, values) = element return [(key, sum(values))] p1 = beam.Pipeline() visit_count = ( p1 |beam.io.ReadFromText(inputs_pattern) |beam.ParDo(SplitRow()) |beam.Map(print) # |beam.ParDo(FilterCardioPatients()) # |beam.ParDo(PairPatients()) # |beam.GroupByKey() # | beam.ParDo(Counting()) # |beam.io.WriteToText('parddo_output.txt') ) p1.run()
两种返回方式的差异
- 当
SplitRow的process方法返回[element.split(',')]时,输出是每行对应的完整列表:['2984641', 'Emily', '35', 'cardio', '1/9/21'] ['9454384', 'Riikka', '86', 'ortho', '21-07-2021'] - 当改为直接返回
element.split(',')时,输出变为每个元素单独一行:2984641 Emily 35 cardio 1/9/21 9454384 Riikka 86 ortho 21-07-2021 9266396
既然Python的split函数本身就返回列表,为何在Apache Beam的ParDo函数中仍需显式返回列表?
原因解析
这是由Apache Beam中ParDo的process方法的设计逻辑决定的:
process方法的返回值是一个可迭代对象,Beam会遍历这个对象,把其中的每一个元素都作为单独的输出元素传递给下游处理。- 当你直接返回
element.split(',')时,这个split得到的列表就是可迭代对象,Beam会逐个取出列表里的字符串,每个字符串都成为一个独立的输出元素,所以下游的print会逐个打印这些字符串。 - 当你返回
[element.split(',')]时,外层的列表是可迭代对象,里面只有一个元素——就是split得到的完整列表,Beam会把这个完整列表作为一个单独的输出元素传递给下游,所以print会直接输出整个列表。
简单来说:如果想把某个对象(比如split后的完整列表)作为一个独立的输出元素,就必须把它包裹在一个可迭代对象(比如列表)里,告诉Beam“这是一个输出单元,不要拆分它”;如果直接返回拆分后的列表,Beam会把列表里的每个元素都当成独立的输出单元处理。
内容的提问来源于stack exchange,提问作者Amarjeet Kushwaha
相关产品推荐
相关产品推荐

