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

Airflow执行Python脚本失败,subprocess配置求助

问题:Airflow DAG无法执行Python脚本

我尝试在Airflow中运行一个DAG,用于执行环境中的Python脚本。在脚本所在目录测试相关命令时逻辑正常,但在Airflow中脚本完全无法执行。我在第一个任务中添加了打印语句,得到如下日志:

[2023-11-17, 07:00:07 UTC] {logging_mixin.py:137} INFO - 脚本/opt/***/src/etl/extract.py不存在或无法执行
[2023-11-17, 07:00:07 UTC] {python.py:177} INFO - Done. Returned value was: None
[2023-11-17, 07:00:07 UTC] {taskinstance.py:1323} INFO - Marking task as SUCCESS. dag_id=run_scripts_daily, task_id=run_extract, execution_date=20231117T070006, start_date=20231117T070007, end_date=20231117T070007
[2023-11-17, 07:00:07 UTC] {local_task_job.py:208} INFO - Task exited with return code 0
[2023-11-17, 07:00:07 UTC] {taskinstance.py:2578} INFO - 1 downstream tasks scheduled from follow-on schedule check

我已尝试解决该问题数小时,认为subprocess.run是可行方案,但无法正确配置,希望获得帮助。我的DAG代码如下:

from airflow.decorators import dag
from airflow.operators.python_operator import PythonOperator
from datetime import datetime, timedelta
import subprocess
import os

default_args = {
    "owner": "airflow",
    "depends_on_past": False,
    "email_on_failure": False,
    "email_on_retry": False,
    "retries": 1,
    "retry_delay": timedelta(minutes=5),
}

def run_extract():
    airflow_home = '/opt/airflow/'
    script_path = os.path.join(airflow_home, 'src/etl/extract.py')
    if os.path.exists(script_path) and os.access(script_path, os.X_OK):
        subprocess.run(['python', script_path], shell=True)
    else:
        print(f'O script {script_path} não existe ou não pode ser executado')

def run_pre_validate():
    #subprocess.run(['python', '../src/validators/pre_validate.py'])
    pass

def run_transform():
    # subprocess.run(['python', '../src/etl/transform.py'], shell=True)
    pass

def run_pos_validate():
    # subprocess.run(['python', '../src/validators/pos_validate.py'])
    pass

def run_load():
    # subprocess.run(['python', '/src/etl/load.py'])
    pass

@dag(
    "run_scripts_daily",
    start_date=datetime(2021, 12, 1),
    max_active_runs=1,
    schedule="@daily",
    default_args=default_args,
    catchup=False,
)
def run_scripts_daily():
    opr_run_extract = PythonOperator(
        task_id="run_extract",
        python_callable=run_extract
    )

    opr_run_pre_validate = PythonOperator(
        task_id="run_pre_validate",
        python_callable=run_pre_validate
    )

    opr_run_transform = PythonOperator(
        task_id="run_transform",
        python_callable=run_transform
    )

    opr_run_pos_validate = PythonOperator(
        task_id="run_pos_validate",
        python_callable=run_pos_validate
    )

    opr_run_load = PythonOperator(
        task_id="run_load",
        python_callable=run_load
    )

    opr_run_extract >> opr_run_pre_validate >> opr_run_transform >> opr_run_pos_validate >> opr_run_load

run_scripts_daily_dag = run_scripts_daily()

解决方案

1. 排查路径与权限问题

  • 确认Airflow Worker执行路径:在run_extract函数中添加print(f"当前工作目录: {os.getcwd()}"),查看Airflow任务的运行路径是否与本地测试一致,避免相对路径错误。
  • 修复文件及目录权限:
    • 给脚本添加执行权限:chmod +x /opt/airflow/src/etl/extract.py,确保Airflow运行用户(通常为airflow)有读和执行权限。
    • 确保父目录权限允许访问:比如/opt/airflow/src/etl权限至少设置为755。

2. 修正subprocess.run配置

  • 参数格式错误:shell=True搭配列表参数会导致执行异常,二选一调整:
    # 方式1:去掉shell=True,用列表参数(推荐,更安全)
    subprocess.run(['python', script_path])
    # 方式2:保留shell=True,改用字符串参数
    subprocess.run(f'python {script_path}', shell=True)
    
  • 添加错误捕获:启用check=True让脚本执行失败时抛出异常,Airflow会标记任务失败,便于排查问题:
    def run_extract():
        airflow_home = '/opt/airflow/'
        script_path = os.path.join(airflow_home, 'src/etl/extract.py')
        if os.path.exists(script_path) and os.access(script_path, os.X_OK):
            try:
                # 捕获执行输出并检查结果
                result = subprocess.run(['python', script_path], check=True, capture_output=True, text=True)
                print(f"脚本输出: {result.stdout}")
            except subprocess.CalledProcessError as e:
                print(f"脚本执行失败: {e.stderr}")
                raise  # 抛出异常让Airflow标记任务失败
        else:
            print(f'脚本 {script_path} 不存在或无法执行')
    

3. 更优的Airflow执行方式

无需使用subprocess,直接导入脚本中的函数执行,更符合Airflow设计逻辑,还能共享Airflow的Python环境:

# 假设extract.py中有main函数
from src.etl.extract import main

def run_extract():
    main()

内容的提问来源于stack exchange,提问作者samurai-py

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 02:35:04