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

如何创建基于Schema名称执行SQL插入的Airflow DAG

实现Airflow动态Schema批量插入DAG

核心逻辑

  • 从LOG_SCHEMA.LOG001中筛选指定Schema下的所有表名
  • 动态生成目标表的插入语句:INSERT INTO [Schema].[表名] SELECT * FROM [Schema].[表名] WHERE COL1=2;
  • 支持通过Airflow变量指定目标Schema,批量执行插入操作

完整DAG代码

from airflow import DAG
from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator
from airflow.operators.python import PythonOperator
from airflow.models import Variable
from datetime import datetime

def generate_insert_queries(**context):
    # 读取Airflow变量指定的目标Schema,默认值为TEST
    target_schema = Variable.get("target_schema", default_var="TEST")
    # 从上游任务拉取表列表数据
    ti = context['ti']
    table_list = ti.xcom_pull(task_ids='get_target_tables')
    
    # 批量生成插入SQL
    insert_queries = []
    for table in table_list:
        if table['SCHEMA_NAME'] == target_schema:
            table_name = table['TABLE_NAME']
            query = f"INSERT INTO {target_schema}.{table_name} SELECT * FROM {target_schema}.{table_name} WHERE COL1=2;"
            insert_queries.append(query)
    
    # 将生成的SQL推送给下游任务
    ti.xcom_push(key='insert_queries', value='\n'.join(insert_queries))

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2024, 1, 1),
    'retries': 1
}

with DAG(
    'snowflake_dynamic_batch_insert',
    default_args=default_args,
    schedule_interval='@daily',
    catchup=False
) as dag:

    # 任务1:获取目标Schema下的所有表名
    get_target_tables = SnowflakeOperator(
        task_id='get_target_tables',
        snowflake_conn_id='snowflake_default',
        sql="""
            SELECT SCHEMA_NAME, TABLE_NAME 
            FROM LOG_SCHEMA.LOG001 
            WHERE SCHEMA_NAME = '{{ var.value.target_schema | default('TEST') }}';
        """,
        do_xcom_push=True,
        return_last=False
    )

    # 任务2:生成批量插入SQL语句
    generate_queries = PythonOperator(
        task_id='generate_insert_queries',
        python_callable=generate_insert_queries,
        provide_context=True
    )

    # 任务3:执行批量插入操作
    execute_batch_insert = SnowflakeOperator(
        task_id='execute_batch_insert',
        snowflake_conn_id='snowflake_default',
        sql="{{ ti.xcom_pull(key='insert_queries', task_ids='generate_insert_queries') }}"
    )

    # 任务依赖关系
    get_target_tables >> generate_queries >> execute_batch_insert

关键细节说明

  • 动态Schema配置:可通过Airflow UI的「变量」功能设置target_schema值,运行时自动生效
  • 数据传递:通过XCom完成任务间的表列表、SQL语句传递,无需额外存储
  • 过滤逻辑:仅处理目标Schema下的表,避免生成无关SQL
  • 批量执行:将所有插入语句拼接为单字符串一次性执行,适合中小规模表批量操作;若表数量过大,可拆分任务分批执行

注意事项

  1. 确保Airflow已配置正确的Snowflake连接(snowflake_default)
  2. 目标Schema下的所有表必须包含COL1字段,否则执行会抛出字段不存在错误
  3. 批量插入大量数据时,需关注Snowflake的事务大小限制,必要时添加提交逻辑

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 19:50:21