Airflow中BashOperator未执行Python脚本致CSV未生成求助
Airflow任务显示成功但未执行Python脚本的排查方案
我是Airflow新手,尝试通过BashOperator运行DAG执行ETL Python脚本,脚本会更新pandas DataFrame并输出.csv文件。但Airflow Web UI显示任务成功,却没生成目标文件,推测Python脚本没被实际执行。以下是我的DAG脚本和日志信息:
DAG脚本代码
from airflow.operators.bash import BashOperator from airflow.models import DAG from airflow.operators.bash_operator import BashOperator from datetime import datetime with DAG('tester', start_date=datetime(2022, 9, 27), schedule_interval='*/10 * * * *', catchup=False) as dag: task1 = BashOperator( task_id='task1', bash_command='echo python3 /G:/xxx/xxxxx/xx/xxxx/t3.py' ) task2 = BashOperator( task_id='task2', bash_command='echo python3 /C:/airflow_docker/scripts/t1.py', ) task3 = BashOperator( task_id = 'task3', bash_command='echo python3 /G:/xxx/xxxxx/xx/xxxx/t2.py' )
任务日志信息
*** Reading local file: /opt/airflow/logs/dag_id=tester/run_id=manual__2022-09-28T10:15:38.095133+00:00/task_id=empresas/attempt=1.log [2022-09-28, 10:15:39 UTC] {taskinstance.py:1171} INFO - Dependencies all met for <TaskInstance: tester.empresas manual__2022-09-28T10:15:38.095133+00:00 [queued]> [2022-09-28, 10:15:39 UTC] {taskinstance.py:1171} INFO - Dependencies all met for <TaskInstance: tester.empresas manual__2022-09-28T10:15:38.095133+00:00 [queued]> [2022-09-28, 10:15:39 UTC] {taskinstance.py:1368} INFO - -------------------------------------------------------------------------------- [2022-09-28, 10:15:39 UTC] {taskinstance.py:1369} INFO - Starting attempt 1 of 1 [2022-09-28, 10:15:39 UTC] {taskinstance.py:1370} INFO - -------------------------------------------------------------------------------- [2022-09-28, 10:15:39 UTC] {taskinstance.py:1389} INFO - Executing <Task(BashOperator): empresas> on 2022-09-28 10:15:38.095133+00:00 [2022-09-28, 10:15:39 UTC] {standard_task_runner.py:52} INFO - Started process 9879 to run task [2022-09-28, 10:15:39 UTC] {standard_task_runner.py:79} INFO - Running: ['***', 'tasks', 'run', 'tester', 'empresas', 'manual__2022-09-28T10:15:38.095133+00:00', '--job-id', '1381', '--raw', '--subdir', 'DAGS_FOLDER/another.py', '--cfg-path', '/tmp/tmptz45sf6g', '--error-file', '/tmp/tmp57jeddaf'] [2022-09-28, 10:15:39 UTC] {standard_task_runner.py:80} INFO - Job 1381: Subtask empresas [2022-09-28, 10:15:39 UTC] {task_command.py:371} INFO - Running <TaskInstance: tester.empresas manual__2022-09-28T10:15:38.095133+00:00 [running]> on host 620a4d8bf7f5 [2022-09-28, 10:15:39 UTC] {taskinstance.py:1583} INFO - Exporting the following env vars: AIRFLOW_CTX_DAG_OWNER=*** AIRFLOW_CTX_DAG_ID=tester AIRFLOW_CTX_TASK_ID=empresas AIRFLOW_CTX_EXECUTION_DATE=2022-09-28T10:15:38.095133+00:00 AIRFLOW_CTX_TRY_NUMBER=1 AIRFLOW_CTX_DAG_RUN_ID=manual__2022-09-28T10:15:38.095133+00:00 [2022-09-28, 10:15:39 UTC] {subprocess.py:62} INFO - Tmp dir root location: /tmp [2022-09-28, 10:15:39 UTC] {subprocess.py:74} INFO - Running command: ['/bin/bash', '-c', 'echo /C:/***_docker/scripts/empresas.py'] [2022-09-28, 10:15:39 UTC] {subprocess.py:85} INFO - Output: [2022-09-28, 10:15:39 UTC] {subprocess.py:92} INFO - /C:/***_docker/scripts/empresas.py [2022-09-28, 10:15:39 UTC] {subprocess.py:96} INFO - Command exited with return code 0 [2022-09-28, 10:15:39 UTC] {taskinstance.py:1412} INFO - Marking task as SUCCESS. dag_id=tester, task_id=empresas, execution_date=20220928T101538, start_date=20220928T101539, end_date=20220928T101539 [2022-09-28, 10:15:39 UTC] {local_task_job.py:156} INFO - Task exited with return code 0 [2022-09-28, 10:15:39 UTC] {local_task_job.py:279} INFO - 0 downstream tasks scheduled from follow-on schedule check
问题原因及修复方案
1. 多余的echo导致脚本未执行
你的所有BashOperator的bash_command都带了echo前缀,这只会打印Python命令字符串,不会实际执行脚本。从日志也能看到,实际执行的命令是echo /C:/xxx/scripts/empresas.py,输出的就是路径,没有运行脚本。
修复: 去掉每个bash_command里的echo,直接写执行命令:
task1 = BashOperator( task_id='task1', bash_command='python3 /G:/xxx/xxxxx/xx/xxxx/t3.py' ) task2 = BashOperator( task_id='task2', bash_command='python3 /C:/airflow_docker/scripts/t1.py', ) task3 = BashOperator( task_id = 'task3', bash_command='python3 /G:/xxx/xxxxx/xx/xxxx/t2.py' )
2. Docker环境下的路径兼容性问题
从日志主机标识620a4d8bf7f5可以看出,Airflow运行在Docker容器中,而你使用的Windows本地路径(G:/、C:/)无法被容器直接访问。
修复:
- 在
docker-compose.yml中添加本地目录到容器的挂载配置:volumes: - ./airflow_docker/scripts:/opt/airflow/scripts - ./your_G_drive_folder:/opt/airflow/data - 将BashOperator中的路径改为容器内部路径:
task2 = BashOperator( task_id='task2', bash_command='python3 /opt/airflow/scripts/t1.py', )
额外优化建议
- 去掉重复的
BashOperator导入,保留from airflow.operators.bash import BashOperator即可。 - 可以改用
PythonOperator直接执行Python函数,避免路径和环境变量问题:from airflow.operators.python import PythonOperator def run_etl(): import pandas as pd # 这里写你的ETL逻辑 df = pd.read_csv('/opt/airflow/data/source.csv') # 数据处理操作 df.to_csv('/opt/airflow/output/result.csv', index=False) etl_task = PythonOperator( task_id='run_etl', python_callable=run_etl )
内容的提问来源于stack exchange,提问作者Paulo Hader
相关产品推荐
相关产品推荐

