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
相关产品推荐
相关产品推荐

