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

