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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 02:23:18