Airflow:跨DAG XCom自定义优先级权重插件触发任务报错排查
问题:跨DAG优先级策略引发SQLAlchemy事务关闭错误
我需要根据DAG A中get_username任务的XCom值,为DAG A的TriggerDagRunOperator触发的DAG B内run_test_task设置优先级。基于Airflow内置的_DownstreamPriorityWeightStrategy类实现了自定义权重规则,插件代码如下:
class CustomPriorityStrategy(_DownstreamPriorityWeightStrategy): @provide_session def get_weight(self, ti: TaskInstance, session: Session) -> int: weight = super().get_weight(ti) username = ti.xcom_pull(task_ids='get_username', session=session) if username == "john": weight += 1000 return weight
在run_test_task中配置使用该策略:
run_test_task = RunTestOperator( ... weight_rule="custom_priority_weight_strategy.CustomPriorityStrategy", ... )
触发任务时出现以下错误:
Traceback (most recent call last): File "/home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/orm/session.py", line 3901, in _bulk_save_mappings persistence._bulk_insert( File "/home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/orm/persistence.py", line 74, in _bulk_insert connection = session_transaction.connection(base_mapper) File "/home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/orm/session.py", line 627, in connection self._assert_active() File "/home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/orm/session.py", line 620, in _assert_active raise sa_exc.ResourceClosedError(closed_msg) sqlalchemy.exc.ResourceClosedError: This transaction is closed
解决方案
- 替换会话获取方式,避免复用已关闭的事务:
@provide_session传入的session在Airflow批量处理场景下可能已被关闭,改用create_session上下文管理器获取活跃会话:
from airflow.utils.session import create_session class CustomPriorityStrategy(_DownstreamPriorityWeightStrategy): def get_weight(self, ti: TaskInstance) -> int: weight = super().get_weight(ti) with create_session() as session: # 指定源DAG ID,跨DAG获取XCom username = ti.xcom_pull( task_ids='get_username', dag_id='dag_a_id', # 替换为实际的DAG A ID session=session ) if username == "john": weight += 1000 return weight
- 修正跨DAG XCom获取逻辑:必须在
xcom_pull中指定dag_id参数,否则默认只会查询当前任务(DAG B)所属DAG的XCom数据,无法拿到DAG A中get_username的结果。
错误原因说明
Airflow在批量处理任务实例时,会复用数据库会话但可能提前终止事务,此时直接使用@provide_session传入的session就会触发事务已关闭的错误。使用create_session可以确保每次获取到的是处于活跃状态的会话。同时跨DAG获取XCom必须明确指定源DAG的ID,否则逻辑无法达到预期效果。
内容的提问来源于stack exchange,提问作者Ivan Agapitov
相关产品推荐
相关产品推荐

