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

Airflow中Dataflow Pipeline序列化(Pickle)失败问题求助

问题

运行的Dataflow Pipeline代码如下:

def dataflow_job_run(files_list, bucket):
    import apache_beam as beam
    from apache_beam.io import ReadAllFromText
    from apache_beam.io import filesystems
    from apache_beam.options.pipeline_options import PipelineOptions
    from <custom_module> import func1

    class SpecialOutputClass(beam.DoFn):
        def process(self, element):
            path = bucket+element[0]
            writer = filesystems.FileSystems.create(f'{path}')
            writer.write(bytes(element[1], 'utf-8'))
            writer.close()

    dataflow_options = <dataflow_options_list>
    options = PipelineOptions(dataflow_options)

    p = beam.Pipeline(options=options)
    transform  = (p | beam.Create(files_list) | ReadAllFromText(with_filename=True) | beam.Map(func1) | beam.GroupByKey() | beam.ParDo(SpecialOutputClass()))

    result = p.run()

(注:原代码中writer.write语句缺少闭合括号,已修正)

该代码在本地运行完全正常,但部署到Airflow(v1.10)(使用CeleryExecutor)时出现错误,问题集中在最后一个ParDo步骤,报错信息如下:

File "/usr/local/lib/python3.8/site-packages/apache_beam/internal/dill_pickler.py", line 285, in loads
return dill.loads(s)
File "/usr/local/lib/python3.8/site-packages/dill/_dill.py", line 275, in loads
return load(file, ignore, **kwds)
File "/usr/local/lib/python3.8/site-packages/dill/_dill.py", line 270, in load
return Unpickler(file, ignore=ignore, **kwds).load()
File "/usr/local/lib/python3.8/site-packages/dill/_dill.py", line 472, in load
obj = StockUnpickler.load(self)
File "/usr/local/lib/python3.8/site-packages/dill/_dill.py", line 826, in _import_module
return __import__(import_name)
ModuleNotFoundError: No module named 'unusual_prefix_*' 

其中*代表DAG名称。推测是Airflow的DAG序列化方式导致ParDo类无法正确序列化。若将SpecialOutputClass改为函数并使用beam.Map()则可正常运行,但希望保留原类结构让代码正常工作,且未设置donot_pickle = True,请问有什么解决办法?

解决办法

方案1:将SpecialOutputClass移至模块顶层

Airflow序列化DAG时,函数内部定义的类会被dill序列化并带上Airflow生成的临时模块前缀(即报错中的unusual_prefix_*),导致Worker节点无法找到对应模块。把类移到模块级别(函数外部)即可避免该问题,同时通过构造函数传入依赖的bucket参数:

import apache_beam as beam
from apache_beam.io import filesystems
from <custom_module> import func1

class SpecialOutputClass(beam.DoFn):
    def __init__(self, bucket):
        self.bucket = bucket

    def process(self, element):
        path = self.bucket + element[0]
        writer = filesystems.FileSystems.create(f'{path}')
        writer.write(bytes(element[1], 'utf-8'))
        writer.close()

def dataflow_job_run(files_list, bucket):
    from apache_beam.io import ReadAllFromText
    from apache_beam.options.pipeline_options import PipelineOptions

    dataflow_options = <dataflow_options_list>
    options = PipelineOptions(dataflow_options)

    p = beam.Pipeline(options=options)
    transform = (
        p 
        | beam.Create(files_list) 
        | ReadAllFromText(with_filename=True) 
        | beam.Map(func1) 
        | beam.GroupByKey() 
        | beam.ParDo(SpecialOutputClass(bucket))
    )

    result = p.run()

方案2:使用@beam.DoFn装饰器定义处理逻辑

该方式可以保留DoFn的特性,同时避免函数内部类的序列化问题(部分Beam版本适用):

def dataflow_job_run(files_list, bucket):
    import apache_beam as beam
    from apache_beam.io import ReadAllFromText
    from apache_beam.io import filesystems
    from apache_beam.options.pipeline_options import PipelineOptions
    from <custom_module> import func1

    @beam.DoFn
    def SpecialOutputFn(element):
        path = bucket + element[0]
        writer = filesystems.FileSystems.create(f'{path}')
        writer.write(bytes(element[1], 'utf-8'))
        writer.close()

    dataflow_options = <dataflow_options_list>
    options = PipelineOptions(dataflow_options)

    p = beam.Pipeline(options=options)
    transform = (
        p 
        | beam.Create(files_list) 
        | ReadAllFromText(with_filename=True) 
        | beam.Map(func1) 
        | beam.GroupByKey() 
        | beam.ParDo(SpecialOutputFn())
    )

    result = p.run()

注:此方式在部分Beam版本中可能仍存在序列化隐患,优先推荐方案1。

方案3:调整Airflow序列化配置(全局变更)

修改Airflow的airflow.cfg配置文件,调整DAG序列化参数:

  • 设置dag_serialization=False,禁用DAG序列化(但可能降低CeleryExecutor的性能)
  • 或设置pickle_protocol=4,提升序列化兼容性

此方案属于集群全局配置变更,需要权衡稳定性与兼容性。

内容的提问来源于stack exchange,提问作者Aaron Gonzalez

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 18:54:24