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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 21:45:42