Apache Beam Dataflow运行模板时出现dict_level1未定义错误
问题:Dataflow模板运行时触发NameError,自定义函数未定义
我在Google Dataflow上基于Apache Beam开发数据处理任务,编写了dict_level1、unnest_dict、dict_level0三个自定义函数,在Pipeline中通过以下代码调用:
| "Unnest 1" >> beam.Map(lambda record: dict_level1(record)) | "Unnest 2" >> beam.Map(lambda record: unnest_dict(record)) | "Unnest 3" >> beam.Map(lambda record: dict_level0(record))
本地执行代码成功生成存储在GCS中的Dataflow模板,但上传模板并启动任务时触发NameError,提示名称dict_level1未定义,错误信息如下:
File "/Users/dario/Repo-c3tech/c3t-tango/./composer/dags/gcp_to_bq_table.py", line 76, in NameError: name 'dict_level1' is not defined
完整代码实现:
import apache_beam as beam import os from apache_beam.options.pipeline_options import PipelineOptions # 生成输出和模板 pipeline_options = { 'project': 'c3t-tango-dev', 'runner': 'DataflowRunner', 'region': 'us-central1', # 请指定正确的区域 'staging_location': 'gs://dario-dev-gcs/dataflow-course/staging', 'template_location': 'gs://dario-dev-gcs/dataflow-course/templates/batch_job_df_gcs_flights4' } pipeline_options = PipelineOptions.from_dictionary(pipeline_options) table_schema = 'airport:STRING, list_delayed_num:INTEGER, list_delayed_time:INTEGER' table = 'c3t-tango-dev:dataflow.flights_aggr' class Filter(beam.DoFn): def process(self, record): if int(record[8]) > 0: return [record] def dict_level1(record): dict_ = {} dict_['airport'] = record[0] dict_['list'] = record[1] return (dict_) def unnest_dict(record): def expand(key, value): if isinstance(value, dict): return [(key + '_' + k, v) for k, v in unnest_dict(value).items()] else: return [(key, value)] items = [item for k, v in record.items() for item in expand(k, v)] return dict(items) def dict_level0(record): #print("Record in dict_level0:", record) dict_ = {} dict_['airport'] = record['airport'] dict_['list_Delayed_num'] = record['list_Delayed_num'][0] dict_['list_Delayed_time'] = record['list_Delayed_time'][0] return (dict_) with beam.Pipeline(options=pipeline_options) as p1: serviceAccount = "./composer/dags/c3t-tango-dev-591728f351ee.json" os.environ["GOOGLE_APPLICATION_CREDENTIALS"] = serviceAccount Delayed_time = ( p1 | "Import Data time" >> beam.io.ReadFromText("gs://dario-dev-gcs/dataflow-course/input/voos_sample.csv", skip_header_lines=1) | "Split by comma time" >> beam.Map(lambda record: record.split(',')) | "Filter Delays time" >> beam.ParDo(Filter()) | "Create a key-value time" >> beam.Map(lambda record: (record[4], int(record[8]))) | "Sum by key time" >> beam.CombinePerKey(sum) ) Delayed_num = ( p1 | "Import Data" >> beam.io.ReadFromText("gs://dario-dev-gcs/dataflow-course/input/voos_sample.csv", skip_header_lines=1) | "Split by comma" >> beam.Map(lambda record: record.split(',')) | "Filter Delays" >> beam.ParDo(Filter()) | "Create a key-value" >> beam.Map(lambda record: (record[4], int(record[8]))) | "Count by key" >> beam.combiners.Count.PerKey() ) Delay_table = ( {'Delayed_num': Delayed_num, 'Delayed_time': Delayed_time} | "Group By" >> beam.CoGroupByKey() | "Unnest 1" >> beam.Map(lambda record: dict_level1(record)) | "Unnest 2" >> beam.Map(lambda record: unnest_dict(record)) | "Unnest 3" >> beam.Map(lambda record: dict_level0(record)) #| beam.Map(print) | "Write to BQ" >> beam.io.WriteToBigQuery( table, schema=table_schema, write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, custom_gcs_temp_location="gs://dario-dev-gcs/dataflow-course/staging") ) p1.run()
原因分析
使用Dataflow Runner生成模板并远程运行时,lambda表达式中的自定义函数无法被完整序列化到模板中。Dataflow序列化Pipeline时,lambda的上下文不会包含这些顶层函数的完整定义,导致Worker节点执行时找不到对应的函数名。
解决方案
方法1:直接传递函数引用给beam.Map
移除lambda包裹,直接将自定义函数本身传给beam.Map,这样Dataflow能正确识别并序列化函数定义:
| "Unnest 1" >> beam.Map(dict_level1) | "Unnest 2" >> beam.Map(unnest_dict) | "Unnest 3" >> beam.Map(dict_level0)
方法2:将自定义函数封装为DoFn类
如果函数逻辑复杂,推荐将函数改写为beam.DoFn的子类,序列化可靠性更高:
class DictLevel1(beam.DoFn): def process(self, record): dict_ = {} dict_['airport'] = record[0] dict_['list'] = record[1] return [dict_] # 调用时使用ParDo | "Unnest 1" >> beam.ParDo(DictLevel1())
额外注意事项
- 确保所有自定义函数/类都定义在主脚本中,或被正确打包到Dataflow的依赖包中(若使用外部模块)。
- 生成模板时,确认
staging_location配置正确,Dataflow会将依赖文件上传至该GCS路径供Worker节点调用。
内容的提问来源于stack exchange,提问作者Dario Villalon
相关产品推荐
相关产品推荐

