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

Snowflake查询ProgrammingError:Airflow DAG插入数据跳过错误记录并导出GCS

问题根因

你当前的代码是整批执行INSERT SELECT语句,只要存在一条违反目标表约束的记录(比如字段类型不匹配、长度超限、非空约束冲突等),整个SQL事务就会整体失败。外层的try-except只能捕获到整个语句执行失败的异常,无法实现逐行跳过错误记录的效果。

解决方案

核心思路是在SQL层提前拆分合法、非法记录,避免整批插入失败,再通过Snowflake内置的导出能力直接把错误记录写入GCS,不需要在Python侧拉取数据再上传,性能更高:

  • 第一步:按照目标表的约束条件筛选合法记录,直接插入目标表
  • 第二步:筛选出不符合约束的错误记录,通过Snowflake的COPY INTO语句直接导出到提前配置好的GCS外部Stage
修改后代码
from airflow import DAG
from airflow.operators import python_operator
from datetime import datetime
import logging
from my_folder import param_file
from snowflake.connector.errors import ProgrammingError, OperationalError

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

args = {
    "owner": "test_dag",
    "start_date": datetime(2021, 3, 2, 10, 1),
    'depends_on_past': False,
    'email': ['mymailid@example.com'],
    'email_on_failure': True
}

# 注意:提前在Snowflake中创建指向目标GCS桶的外部Stage,替换下面的YOUR_GCS_STAGE名称
GCS_STAGE_NAME = "YOUR_GCS_STAGE"
# 错误记录导出到GCS的文件前缀
ERROR_FILE_PREFIX = "employees_error_records/"

# 合法记录插入SQL:按照你目标表的实际约束调整筛选条件
insert_valid_query = """
insert into employees(first_name, last_name, workphone, city, postal_code)
select 
    contractor_first,
    contractor_last,
    worknum,
    null,
    zip_code 
from contractors
where 
    -- 按需调整校验规则,比如字段非空、长度限制、类型合法等
    contractor_first is not null
    and contractor_last is not null
    and length(worknum) <= 20
    and try_cast(zip_code as varchar(10)) is not null
"""

# 错误记录查询SQL:筛选不符合上面约束的记录
get_error_query = """
copy into @{stage_name}/{file_prefix}
from (
    select 
        contractor_first,
        contractor_last,
        worknum,
        zip_code,
        current_timestamp() as error_time
    from contractors
    where 
        contractor_first is null
        or contractor_last is null
        or length(worknum) > 20
        or try_cast(zip_code as varchar(10)) is null
)
file_format = (type = 'CSV' field_delimiter = ',' skip_header = 0)
overwrite = false
single = false
max_file_size = 1073741824
""".format(stage_name=GCS_STAGE_NAME, file_prefix=ERROR_FILE_PREFIX)

def run_sql(**context):
    # 不要用全局游标,避免连接超时失效,每次任务运行新建连接
    conn = None
    cur = None
    try:
        conn = param_file.get_snowflake_conn() # 替换成你自己的获取Snowflake连接的方法
        cur = conn.cursor()
        # 插入合法记录
        cur.execute(insert_valid_query)
        logger.info(f"成功插入{cur.rowcount}条合法记录")
        # 导出错误记录到GCS
        cur.execute(get_error_query)
        logger.info("错误记录已导出到GCS存储桶")
        conn.commit()
    except ProgrammingError as db_ex:
        logger.error(f"SQL执行错误: {db_ex}")
        if conn:
            conn.rollback()
        raise
    except OperationalError as conn_ex:
        logger.error(f"Snowflake连接错误: {conn_ex}")
        raise
    finally:
        if cur:
            cur.close()
        if conn:
            conn.close()

with DAG(
    dag_id="EXCEPTION_testing", 
    schedule_interval=None, 
    max_active_runs=1,
    catchup=False,
    default_args=args
) as dag:

    task1= python_operator.PythonOperator(
        task_id='snowflake_query_execution',
        python_callable=run_sql,
        provide_context=True,
        dag=dag
    )
注意事项
  • 需要提前在Snowflake中创建指向目标GCS存储桶的外部Stage,配置好GCS服务账号的写入权限,保证Snowflake有权限写入你的存储桶
  • 上述代码中的校验规则仅为示例,需要按照你实际的employees表约束调整,比如枚举值范围、特殊字符限制等
  • 全局游标写法存在连接超时失效的问题,建议每次任务运行时新建Snowflake连接,避免偶发的连接错误

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 09:15:02