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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 07:22:39