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

Airflow中DAG运行时如何获取参数并动态选择数据表?

问题:Airflow DAG通过dag_run参数选择表时触发TypeError错误

场景说明

  • 现有一个Airflow DAG,用于在指定表上执行多段SQL脚本,支持两种运行模式:
    • 操作生产表(production tables)
    • 操作归档冻结表(frozen archived tables)
  • 期望通过dag_run的参数值切换目标表,但运行代码时出现TypeError: unhashable type: 'PlainXComArg'错误

原始代码

from datetime import datetime, timedelta, date
from airflow import DAG
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator
from airflow.models import DagRun
from airflow.decorators import task

data_dict = {
    'prod':{'tab_1': 'table_name_1','tab_2': 'table_name_2'},
    'arch':{'tab_1': 'table_name_1_arch', 'tab_2': 'table_name_2_arch'}
}


with DAG(
        'sgk_test_2',
        description='sgk_test_2',
        tags=["sgk_test"],
        schedule_interval=None,
        start_date=datetime(2025, 7, 1),
        default_args={
            'retries': 0,
            'retry_delay': timedelta(minutes=1),
            'conn_id': 'sgk_gp_tau_pvr'
        },
        params={
            'tab_type':'',
        }
) as dag:

    @task(task_id='task_0')
    def get_type(**context):
        params = context.get('params', {})
        tab_type = params.get('tab_type')
        return tab_type

    tab_type = get_type()

    task_1 = SQLExecuteQueryOperator(
        task_id='task_1',
        sql=f"select * from {data_dict[tab_type]['tab_1']}"
    )
    
    task_2 = SQLExecuteQueryOperator(
        task_id='task_2',
        sql=f"select * from {data_dict[tab_type]['tab_2']}"
    )

    task_0 >> task_1 >> task_2

错误信息

sql=f"select * from {data_dict[tab_type]['tab_1']}"
                         ~~~~~~~~~^^^^^^^^^^
TypeError: unhashable type: 'PlainXComArg'

问题根源

tab_type是Airflow的PlainXComArg对象(任务返回值的封装容器),并非实际字符串值。Airflow解析DAG时会提前执行模板渲染,此时tab_type尚未获取到运行时的参数值,直接用它作为字典键访问data_dict会触发类型错误。

解决方案

方案一:直接利用Jinja模板访问参数(推荐)

无需额外任务获取参数,将表映射字典注入DAG的自定义宏,通过Jinja模板直接读取dag_run参数并匹配表名:

from datetime import datetime, timedelta, date
from airflow import DAG
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator

data_dict = {
    'prod':{'tab_1': 'table_name_1','tab_2': 'table_name_2'},
    'arch':{'tab_1': 'table_name_1_arch', 'tab_2': 'table_name_2_arch'}
}

with DAG(
        'sgk_test_2',
        description='sgk_test_2',
        tags=["sgk_test"],
        schedule_interval=None,
        start_date=datetime(2025, 7, 1),
        default_args={
            'retries': 0,
            'retry_delay': timedelta(minutes=1),
            'conn_id': 'sgk_gp_tau_pvr'
        },
        params={
            'tab_type':'prod',  # 设置默认值避免空参数报错
        },
        user_defined_macros={'table_mapping': data_dict}  # 注入表映射字典到模板
) as dag:

    task_1 = SQLExecuteQueryOperator(
        task_id='task_1',
        sql="select * from {{ table_mapping[params.tab_type]['tab_1'] }}"
    )
    
    task_2 = SQLExecuteQueryOperator(
        task_id='task_2',
        sql="select * from {{ table_mapping[params.tab_type]['tab_2'] }}"
    )

    task_1 >> task_2

方案二:保留参数获取任务,通过XCom传递值

如果需要对参数做额外处理(如校验、转换),可保留参数获取任务,通过Jinja模板读取XCom值:

from datetime import datetime, timedelta, date
from airflow import DAG
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator
from airflow.decorators import task

data_dict = {
    'prod':{'tab_1': 'table_name_1','tab_2': 'table_name_2'},
    'arch':{'tab_1': 'table_name_1_arch', 'tab_2': 'table_name_2_arch'}
}


with DAG(
        'sgk_test_2',
        description='sgk_test_2',
        tags=["sgk_test"],
        schedule_interval=None,
        start_date=datetime(2025, 7, 1),
        default_args={
            'retries': 0,
            'retry_delay': timedelta(minutes=1),
            'conn_id': 'sgk_gp_tau_pvr'
        },
        params={
            'tab_type':'',
        },
        render_template_as_native_obj=True  # 开启原生对象渲染支持
) as dag:

    @task(task_id='get_type')
    def get_type(**context):
        params = context.get('params', {})
        tab_type = params.get('tab_type')
        # 可在此添加参数校验逻辑
        return tab_type

    tab_type_task = get_type()

    task_1 = SQLExecuteQueryOperator(
        task_id='task_1',
        sql="select * from {{ data_dict[ti.xcom_pull(task_ids='get_type')]['tab_1'] }}",
        templates_dict={'data_dict': data_dict}  # 传递表映射字典到模板
    )
    
    task_2 = SQLExecuteQueryOperator(
        task_id='task_2',
        sql="select * from {{ data_dict[ti.xcom_pull(task_ids='get_type')]['tab_2'] }}",
        templates_dict={'data_dict': data_dict}
    )

    tab_type_task >> task_1 >> task_2

说明

  • 方案一更简洁高效,直接利用Airflow的参数传递和模板能力,无需额外任务
  • 方案二适合需要对参数做预处理的场景,比如参数合法性校验、格式转换等

内容的提问来源于stack exchange,提问作者Serge Kashlik

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 08:43:20