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

