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

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正确标记任务失败状态。
  • 显式关闭连接:任务结束后主动关闭数据库连接,避免资源泄漏。

验证步骤

  1. 将修复后的代码替换到你的DAG文件中
  2. 重启Airflow Scheduler和Webserver(或等待DAG自动重新解析)
  3. 触发DAG后,查看Graph视图,任务会正常进入running状态,执行完成后根据结果显示success或failed

内容的提问来源于stack exchange,提问作者dark horse

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 06:54:53