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

Airflow DAG中Snowflake批量插入Python函数去重与日志处理问题

Snowflake批量INSERT生成及错误处理优化(适配Airflow)

问题背景

  • 刚接触Airflow,需编写Python函数生成Snowflake批量INSERT语句,用于SnowflakeOperator
  • 当前函数生成重复INSERT语句,需修复
  • 需添加snowflake.connector.errors.ProgrammingError错误日志捕获
  • 函数需输出字符串,适配SnowflakeOperator的要求

测试用DataFrame代码:

import pandas as pd

data = [['test', 'table01'], ['test', 'table02'], ['my_schema', 'table03'], ['schemaxxx', 'table04']]
df_new = pd.DataFrame(data, columns=['schema', 'table'])

原函数及问题:

my_schema = 'my_schema'
def my_insert_fnc(df):
    my_dict = dict()
    for i in df['schema'].unique().tolist():
        df_x = df[df['schema'] == i]
        my_dict[i] = df_x['table'].tolist()
        sql_list = []
        for schema, tables in my_dict.items():
            for table in tables:
                if schema == my_schema:
                    sql_list.append(f"INSERT INTO {schema}.{table} SELECT * FROM {schema}.{table} where col2 = 1;")
                    print(sql_list)

执行后重复输出:

['INSERT INTO my_schema.table03 SELECT * FROM my_schema.table03 where col2 = 1;']
['INSERT INTO my_schema.table03 SELECT * FROM my_schema.table03 where col2 = 1;']

问题原因

原函数将sql_list初始化、内层遍历my_dict的逻辑嵌套在外层for i in df['schema'].unique()循环内部。每遍历一个新schema,my_dict会新增键值对,内层循环会重新遍历整个my_dict,导致my_schema下的表被重复生成INSERT语句。

优化后的实现

方案1:生成无重复SQL字符串(适配SnowflakeOperator)

仅生成正确的批量SQL语句,直接供SnowflakeOperator执行:

my_schema = 'my_schema'

def generate_insert_sql(df):
    sql_list = []
    # 直接筛选目标schema的表,减少多余遍历
    target_df = df[df['schema'] == my_schema]
    for _, row in target_df.iterrows():
        schema_name = row['schema']
        table_name = row['table']
        sql = f"INSERT INTO {schema_name}.{table_name} SELECT * FROM {schema_name}.{table_name} WHERE col2 = 1;"
        sql_list.append(sql)
    # 拼接为单个字符串,SnowflakeOperator支持执行多条语句
    return ' '.join(sql_list)

方案2:带错误捕获的执行逻辑(适配PythonOperator)

如果需要自定义错误日志捕获,建议用PythonOperator直接执行插入逻辑,灵活处理异常:

import logging
import pandas as pd
from snowflake.connector import connect
from snowflake.connector.errors import ProgrammingError
from airflow.operators.python import PythonOperator

# Snowflake连接配置(建议从Airflow变量/密钥管理器读取)
SNOWFLAKE_CONFIG = {
    "account": "你的账号",
    "user": "你的用户名",
    "password": "你的密码",
    "warehouse": "你的仓库",
    "database": "你的数据库"
}

my_schema = 'my_schema'

def execute_snowflake_inserts(df):
    conn = None
    cur = None
    try:
        conn = connect(**SNOWFLAKE_CONFIG)
        cur = conn.cursor()
        target_df = df[df['schema'] == my_schema]
        for _, row in target_df.iterrows():
            schema_name = row['schema']
            table_name = row['table']
            sql = f"INSERT INTO {schema_name}.{table_name} SELECT * FROM {schema_name}.{table_name} WHERE col2 = 1;"
            try:
                cur.execute(sql)
                logging.info(f"成功执行{schema_name}.{table_name}的插入操作")
            except ProgrammingError as e:
                logging.error(f"{schema_name}.{table_name}插入失败: {str(e)}")
                # 可选:根据业务需求决定是否继续执行其他语句
                continue
        conn.commit()
    except Exception as e:
        logging.error(f"Snowflake操作发生意外错误: {str(e)}")
        if conn:
            conn.rollback()
    finally:
        # 确保资源释放
        if cur:
            cur.close()
        if conn:
            conn.close()

# 在Airflow DAG中调用
insert_task = PythonOperator(
    task_id="snowflake_batch_insert",
    python_callable=execute_snowflake_inserts,
    op_kwargs={"df": df_new}
)

代码说明

  1. 重复语句修复:

    • 直接筛选目标my_schema的DataFrame,避免无效遍历
    • 移除嵌套的字典遍历逻辑,改为逐行处理目标表,确保每个表只生成一次INSERT语句
  2. 错误日志处理:

    • 精准捕获ProgrammingError,输出包含具体表名的错误日志
    • 添加全局异常捕获,处理其他意外错误
    • 完善连接关闭逻辑,避免资源泄漏
  3. Airflow适配:

    • 方案1生成的字符串可直接传入SnowflakeOperator的sql参数
    • 方案2通过PythonOperator实现自定义执行逻辑,更灵活控制错误处理流程

内容的提问来源于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 03:40:33