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

Airflow脚本改造:实现Python函数独立运行并监控任务状态

解决Airflow任务拆分与状态跟踪问题

你的思路完全正确——把单函数逻辑拆分成独立任务,就能清晰跟踪每个步骤的运行状态,但目前的实现有几个关键问题需要修正,尤其是不能用XCom传递数据库连接对象,因为数据库连接属于不可序列化的资源,XCom没办法存储这类对象。下面我会给出两种可行的改造方案:

问题分析(你的版本2存在的问题)

  • insert_data函数里return (v1)之后的代码永远不会执行,导致插入逻辑根本无法运行
  • 尝试用XCom传递db_con连接对象,这是行不通的,XCom仅支持传递可序列化的数据(比如字符串、数字、字典等)
  • insert_data里没有从kwargs中获取task_instance,会直接报变量未定义错误

方案一:每个任务独立创建数据库连接(简单直接)

这种方案下,每个任务自己负责创建和关闭数据库连接,虽然会多一次连接操作,但胜在简单清晰,适合小型任务场景:

## 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")

#######################
## 复用的数据库连接函数
def get_db_connection():
    try:
        db_con = psycopg2.connect(
            "dbname='name' user='user' password='pass' host='host' port='port' sslmode='require'"
        )
        print('Connected successfully')
        return db_con
    except Exception as e:
        print(f"Connection Failed: {str(e)}")
        raise  # 抛出异常让Airflow标记任务失败

## 任务1:仅验证数据库连接(单独跟踪连接状态)
def db_connect_task(**kwargs):
    conn = get_db_connection()
    conn.close()  # 验证完成后关闭连接

## 任务2:执行数据插入操作
def insert_data_task(**kwargs):
    conn = get_db_connection()
    try:
        cur = conn.cursor()
        cur.execute("""insert into tbl_1 select id,bill_no,status from tbl_2 limit 2;""")
        conn.commit()  # 务必提交事务
        print("Data inserted successfully")
    except Exception as e:
        conn.rollback()  # 出错时回滚事务
        print(f"Insert failed: {str(e)}")
        raise
    finally:
        cur.close()
        conn.close()  # 确保连接最终关闭

##########################################
t1 = PythonOperator(
    task_id='DB_Connect_Check',
    python_callable=db_connect_task,
    provide_context=True,
    dag=dag
)

t2 = PythonOperator(
    task_id='Insert_Data_To_Table',
    python_callable=insert_data_task,
    provide_context=True,
    dag=dag
)

t1 >> t2

方案一的优点:

  • 每个任务完全独立,连接失败会直接标记DB_Connect_Check任务失败,插入失败标记Insert_Data_To_Table失败,状态清晰
  • 彻底避免了XCom传递不可序列化对象的问题
  • 每个任务都正确处理了连接关闭和事务的提交/回滚,避免资源泄漏

方案二:使用Airflow Connections管理(推荐,更规范)

Airflow自带Connections管理功能,你可以在Airflow UI的Admin -> Connections里配置数据库连接,这样代码里不用硬编码账号密码,更安全也更易维护:

步骤1:在Airflow UI配置连接

  • 进入Admin -> Connections,点击"Create"
  • Conn Id: 自定义为my_postgres_conn(代码里会用到这个ID)
  • Conn Type: 选择Postgres
  • 填写Host、Schema(对应dbname)、Login、Password、Port等信息,SSL Mode选择require

步骤2:改造后的代码

## Third party Library Imports
from airflow import DAG
from airflow.operators.python_operator import PythonOperator
from airflow.hooks.postgres_hook import PostgresHook  # 使用Airflow自带的Postgres Hook
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")

#######################
## 任务1:验证数据库连接
def db_connect_task(**kwargs):
    hook = PostgresHook(postgres_conn_id='my_postgres_conn')
    # 尝试获取连接来验证可用性
    conn = hook.get_conn()
    conn.close()
    print("Database connection verified successfully")

## 任务2:执行数据插入操作
def insert_data_task(**kwargs):
    hook = PostgresHook(postgres_conn_id='my_postgres_conn')
    try:
        # 使用hook的run方法执行SQL,自动处理连接和事务
        hook.run("""insert into tbl_1 select id,bill_no,status from tbl_2 limit 2;""", autocommit=True)
        print("Data inserted successfully")
    except Exception as e:
        print(f"Insert failed: {str(e)}")
        raise

##########################################
t1 = PythonOperator(
    task_id='DB_Connect_Check',
    python_callable=db_connect_task,
    provide_context=True,
    dag=dag
)

t2 = PythonOperator(
    task_id='Insert_Data_To_Table',
    python_callable=insert_data_task,
    provide_context=True,
    dag=dag
)

t1 >> t2

方案二的优点:

  • 不用硬编码数据库账号密码,安全性更高
  • Airflow的Hook已经封装了连接的创建、关闭、事务处理,减少重复代码
  • 同样实现了任务拆分,每个任务的运行状态都能单独跟踪

关键注意点

  • 永远不要用XCom传递数据库连接、文件句柄这类不可序列化的对象
  • 数据库操作一定要处理事务:成功提交,失败回滚,避免数据不一致
  • 任务失败时要抛出异常,这样Airflow才会正确标记任务为失败状态,而不是仅仅打印错误信息

内容的提问来源于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 07:41:34