Apache Beam Python运行报Can't pickle generator objects错误咨询
错误产生原因
- 核心根因是你使用
beam.Map调用了带yield关键字的assign_columns函数:yield会把普通函数变为生成器函数,调用后返回的是生成器对象,而Apache Beam依赖的pickle序列化机制不支持序列化生成器对象,直接触发了报错。 - 附加说明:
beam.Map设计用来接收单输入对应单输出的处理函数,会直接把函数返回值作为输出元素;如果你的处理逻辑是单输入对应多输出,应该使用beam.FlatMap或者beam.ParDo,这两个算子会自动展开生成器、列表这类可迭代对象的内部元素,不会把可迭代对象本身作为输出。
修复方案
你当前的assign_columns逻辑是单row输入、单row输出,不需要生成器特性,直接把函数末尾的yield row改为return row即可,修改后的代码如下:
def assign_columns(row, mapping): main_account = row['filepath_platform_id'] main_report = row['filepath_report_name'] for _mapping in mapping: _side_account = _mapping['account_id'] _side_report = _mapping['report_name'] if main_account == _side_account and main_report == _side_report: _cols = json.loads(_mapping['mapping']) row = renames_columns(row, _cols) return row
如果你希望保留yield写法,也可以不修改函数,只将'Assign Side Columns'步骤的beam.Map替换为beam.FlatMap即可。
内容的提问来源于stack exchange,提问作者Thomas GUILLAUX
相关产品推荐
相关产品推荐

