Airflow触发后脚本未执行排查:Redshift数据插入任务异常
问题排查与修复方案
从你的描述和提供的代码、运行日志来看,核心问题出在PythonOperator的配置错误,导致Airflow在解析DAG文件的阶段就提前执行了任务逻辑,而非将任务正确注册到调度系统中,这就是Graph视图始终显示no status的根本原因。
关键错误分析
1. PythonOperator的python_callable参数误用
在你更新后的代码里,这一行是致命错误:
t1 = PythonOperator( task_id='DWH_Connect', python_callable=job_run(), # 错误:直接调用了函数,而非传递函数引用 # bash_command='python3 ~/airflow/dags/sample.py', dag=dag)
当你写job_run()时,Airflow在加载解析DAG文件的瞬间就会立即执行这个函数(这也是日志里提前出现Task Complete: Insert success的原因),而不是把它作为待调度的任务逻辑保存下来。正确的写法应该是传递函数名job_run(不带括号)。
2. 初始脚本的冗余参数干扰
初始脚本里给PythonOperator添加了bash_command参数,这是BashOperator的专属配置项,PythonOperator不需要该参数,它会干扰任务的正常注册流程。
修复后的完整代码
## Third party Library Imports import psycopg2 from airflow import DAG from airflow.operators.python_operator import PythonOperator from datetime import datetime, timedelta # Following are defaults which can be overridden later on default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime(2018, 1, 23, 12), 'email': ['airflow@airflow.com'], 'email_on_failure': False, 'email_on_retry': False, 'retries': 1, 'retry_delay': timedelta(minutes=5), } dag = DAG('sample_dag', default_args=default_args, catchup=False, schedule_interval="@once") ####################### ## Login to DB def data_warehouse_login(): try: dwh_connection = psycopg2.connect( "dbname='name' user='user' password='pass' host='host' port='port' sslmode='require'" ) print('Connected successfully') return dwh_connection except Exception as e: print(f"Connection Failed: {str(e)}") raise # 抛出异常让Airflow捕获,标记任务失败 def insert_data(dwh_connection): cur = dwh_connection.cursor() try: cur.execute("""insert into tbl_1 select id,bill_no,status from tbl_2 limit 2;""") dwh_connection.commit() print('Task Complete: Insert success') except Exception as e: dwh_connection.rollback() print(f"Insert failed: {str(e)}") raise # 抛出异常触发任务失败状态 finally: cur.close() dwh_connection.close() def job_run(): conn = data_warehouse_login() insert_data(conn) ########################################## t1 = PythonOperator( task_id='DWH_Connect', python_callable=job_run, # 正确:传递函数引用,不带括号 dag=dag)
额外优化建议
- 避免全局变量:原代码使用
global dwh_connection,在Airflow分布式环境中,任务可能在不同Worker节点执行,全局变量无法跨节点共享,改用函数返回值传递连接更可靠。 - 完善异常处理:原代码的
except块仅打印错误但不抛出,Airflow会误以为任务执行成功,添加raise可以让Airflow正确标记任务失败状态。 - 显式关闭连接:任务结束后主动关闭数据库连接,避免资源泄漏。
验证步骤
- 将修复后的代码替换到你的DAG文件中
- 重启Airflow Scheduler和Webserver(或等待DAG自动重新解析)
- 触发DAG后,查看Graph视图,任务会正常进入
running状态,执行完成后根据结果显示success或failed
内容的提问来源于stack exchange,提问作者dark horse
相关产品推荐
相关产品推荐

