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

为何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 20:22:41