使用BeamRunPythonPipelineOperator时Composer安装setup.py报错
问题描述
在使用BeamRunPythonPipelineOperator传递setup_file参数时,部分DAG运行失败,随机出现文件不存在或IO错误,核心报错日志如下:
[2022-11-16, 05:03:19 UTC] {beam.py:127} WARNING - error: [Errno 2] No such file or directory: 'csv_converter-0.0.1/csv_converter.egg-info/PKG-INFO' [2022-11-16, 05:03:20 UTC] {beam.py:127} WARNING - subprocess.CalledProcessError: Command '['/usr/bin/python3', 'setup.py', 'sdist', '--dist-dir', '/tmp/tmpifl6ty8k']' returned non-zero exit status 1.
使用的算子代码:
BeamRunPythonPipelineOperator( task_id='xxxx', runner="DataflowRunner", py_file=f'/home/airflow/gcs/data/csv_converter/main.py', pipeline_options={ 'project_id': project_id, 'input_path': input_path, 'output_path': output_path, 'schema_path': schema_path, 'service_account': service_account, 'no_use_public_ips': True, 'subnetwork': subnetwork, 'staging_location': staging_location, 'temp_location': temp_location, "setup_file": f'/home/airflow/gcs/data/csv_converter/setup.py', "machine_type": "n1-standard-4", "num_workers": 5, "max_num_workers": 10, }, py_options=[], py_interpreter='python3', py_system_site_packages=False, dataflow_config=DataflowConfiguration( job_name='{{task.task_id}}', location=gce_region, wait_until_finished=False, gcp_conn_id="dataflow_conn" ), )
背景说明:CSV文件存入GCS时触发DAG,已将调度器从2个2vCPU扩容到4个4vCPU,问题依然存在。环境版本:
- Composer 2.0.31
- Airflow 2.3.3
- apache-airflow-providers-google 8.1.0
- apache-beam 2.41.0
解决方案
1. 改用GCS路径传递setup_file
当前使用的/home/airflow/gcs/data/...是Airflow worker挂载的本地路径,并发场景下容易出现文件访问冲突或临时文件被清理的问题。直接指定GCS路径让Dataflow自行拉取,避免本地文件竞争:
"setup_file": f'gs://{你的存储桶名}/data/csv_converter/setup.py',
2. 隔离临时构建目录
Beam构建sdist包时默认使用系统临时目录,多任务并发时会出现目录复用冲突。可以给每个任务分配独立临时目录:
在算子中添加环境变量配置:
BeamRunPythonPipelineOperator( # 其他参数不变 env_vars={"TMPDIR": "/tmp/airflow_dataflow_build_{{task_instance_id}}"}, )
或者在setup.py开头添加清理逻辑,避免残留旧构建文件:
import shutil import os # 清理旧构建目录 for dir_name in ['csv_converter.egg-info', 'dist']: if os.path.exists(dir_name): shutil.rmtree(dir_name)
3. 提前构建sdist包,跳过动态构建
跳过Beam自动执行setup.py sdist的步骤,本地提前构建好包上传到GCS:
python setup.py sdist
将生成的dist/csv_converter-0.0.1.tar.gz上传到GCS,然后修改pipeline参数,用extra_package指定包路径:
pipeline_options={ # 替换setup_file为extra_package "extra_package": f'gs://{你的存储桶名}/data/csv_converter/dist/csv_converter-0.0.1.tar.gz', # 其他参数保留 }
4. 检查setup.py路径依赖
确保setup.py中引用的所有文件(模块、资源文件)使用相对路径,避免硬编码本地绝对路径。例如使用package_data时,路径要相对于setup.py所在目录:
setup( # 其他配置 package_data={ 'csv_converter': ['schemas/*.json'], }, )
5. 限制DAG并发数
减少单DAG同时运行的任务数,降低文件系统竞争压力:
在DAG定义中添加并发限制:
default_args={ # 其他默认参数 'max_active_tasks': 3, 'max_active_runs': 3, } dag = DAG( 'your_dag_id', default_args=default_args, # 其他DAG配置 )
内容的提问来源于stack exchange,提问作者nano
相关产品推荐
相关产品推荐

