Apache Beam Python 传递上一步变量报PCollection不可下标错误
错误产生原因
报错TypeError: 'PCollection' object is not subscriptable的核心原因有两个:
- 你在
FetchFileNameDoFn中使用了TaggedOutput输出多分支PCollection,但调用ParDo时没有通过.with_outputs()方法显式声明标签输出,此时ParDo返回的结果是普通单PCollection对象,本身不支持下标访问,也不会自动绑定标签对应的分支PCollection。 - 代码中实例化
FetchFileName时传入了'start'参数,但DoFn类没有实现对应的构造方法接收该参数,运行时也会触发参数不匹配的错误。
修正方案
要正确实现跨步骤的侧输入传递,按如下步骤调整即可:
- 调用ParDo时追加
.with_outputs()声明所有标签输出,指定主输出别名,此时返回的结果对象可以通过.标签名的方式获取对应分支的PCollection - 移除DoFn实例化时多余的
'start'传参,或补全DoFn的__init__方法接收对应参数 - 分别引用主输出PCollection和标签输出PCollection,将标签输出封装为
AsSingleton侧输入传入后续ParDo步骤
修正后完整代码
import json import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions from apache_beam import pvalue # 自定义DoFn class FetchFileName(beam.DoFn): def process(self, element): element = json.loads(element) gdata = element['data'] # 主输出:data数组下的每个元素 for elm in gdata: yield elm # 标签输出:JSON根层级的id值 yield pvalue.TaggedOutput('sdata', element['id']) # 示例AdjustObject实现,用于验证侧输入传递 class AdjustObject(beam.DoFn): def process(self, element, service_id): element['service_id'] = service_id yield element if __name__ == '__main__': input_files = 'message.json' options = PipelineOptions() with beam.Pipeline(options=options) as p: # 读取JSON并拆分多输出 multi_output_result = ( p | 'Read JSON' >> beam.io.ReadFromText(input_files) | 'Fetch Data and File Name' >> beam.ParDo(FetchFileName()).with_outputs( 'sdata', # 声明标签输出 main='emp_data' # 指定主输出的别名 ) ) # 将标签输出sdata封装为单例侧输入 sdata_side_input = beam.pvalue.AsSingleton(multi_output_result.sdata) # 基于主输出做后续转换,传入侧输入 final_result = ( multi_output_result.emp_data | 'Parallel Transform of Data' >> beam.ParDo(AdjustObject(), service_id=sdata_side_input) | 'PRINT' >> beam.Map(print) )
注意事项
如果你的输入路径包含多个JSON文件、或单个文件内有多行JSON记录,会生成多个
sdata值,此时AsSingleton会抛出“PCollection包含多于1个元素”的校验错误。如果业务上所有记录的id一致,可以先对sdata分支做去重后再转Singleton;如果id不唯一,需要根据业务逻辑改用AsList等其他侧输入类型。
内容的提问来源于stack exchange,提问作者Sharvil Popli
相关产品推荐
相关产品推荐

