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
相关产品推荐
相关产品推荐

