如何创建基于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
- 批量执行:将所有插入语句拼接为单字符串一次性执行,适合中小规模表批量操作;若表数量过大,可拆分任务分批执行
注意事项
- 确保Airflow已配置正确的Snowflake连接(
snowflake_default) - 目标Schema下的所有表必须包含
COL1字段,否则执行会抛出字段不存在错误 - 批量插入大量数据时,需关注Snowflake的事务大小限制,必要时添加提交逻辑
内容的提问来源于stack exchange,提问作者My80
相关产品推荐
相关产品推荐

