Snowpark中记录无法插入Notification Table问题求助
问题排查与修复:Snowpark脚本无法插入记录到Notification Table
核心问题分析
你的脚本能打印符合条件的记录但无法插入,主要是以下几个错误导致:
1. INSERT语句语法错误
脚本中写了两次VALUES子句,第一次是无效的占位符写法,第二次才是实际值,这会导致SQL语法错误,数据库无法执行该语句。
2. 未触发SQL执行
Snowpark中session.sql()仅创建查询对象,不会自动执行,必须调用.collect()或.execute()才能触发实际的数据库操作。
3. 字符串字段未加引号
STATUS字段是字符串类型,格式化时未添加单引号,会导致SQL语法错误。
4. 列数不匹配
NOTIFICATION_TABLE包含FILE_NAME字段,但你的INSERT语句未包含该列,也未提供默认值,插入时会因列数不匹配失败。
5. 未执行建表语句
create_table_query仅定义但未执行,若NOTIFICATION_TABLE不存在,插入操作会直接失败。
6. 时间类型格式化风险
直接将Python datetime对象格式化到SQL中可能出现格式不兼容问题,同时存在SQL注入风险。
修正后的完整脚本
import snowflake.snowpark as snowpark from datetime import datetime schemas = ['AW'] database = 'DEV_DB' create_table_query = """ CREATE TABLE IF NOT EXISTS AW.NOTIFICATION_TABLE ( TABLE_NAME VARCHAR, FILE_NAME VARCHAR, LATEST_TIMESTAMP TIMESTAMP, SCAN_TIMESTAMP TIMESTAMP, MAX_AGE INT, STATUS VARCHAR, MESSAGE VARCHAR ) """ def main(session: snowpark.Session): # 先执行建表语句,确保表存在 session.sql(create_table_query).collect() table_notification_intervals = { 'Table1': 1 } # 处理AW schema下的表 for schema in schemas: print("Processing schema: {0}".format(schema)) # 获取指定schema下的表和视图列表 query = "SHOW OBJECTS IN {0}.{1}".format(database, schema) objects = session.sql(query).collect() for object_row in objects: object_name = object_row['name'] object_type = object_row['kind'] if object_type == 'TABLE' and object_name in table_notification_intervals: print("Processing table: {0}".format(object_name)) # 检查表是否包含FILE_NAME列 check_column_query = """ SELECT COLUMN_NAME FROM INFORMATION_SCHEMA.COLUMNS WHERE TABLE_SCHEMA = '{0}' AND TABLE_NAME = '{1}' AND COLUMN_NAME = 'FILE_NAME' """.format(schema, object_name) column_result = session.sql(check_column_query).collect() if len(column_result) == 0: # 表不包含FILE_NAME列,查询表的最后修改时间 subquery = """ SELECT LAST_ALTERED AS TIME_STAMP FROM INFORMATION_SCHEMA.TABLES WHERE TABLE_SCHEMA = '{0}' AND TABLE_NAME = '{1}' """.format(schema, object_name) result = session.sql(subquery).collect() if len(result) > 0: latest_timestamp = result[0]['TIME_STAMP'] days_since_update = (datetime.utcnow().date() - latest_timestamp.date()).days notification_interval = table_notification_intervals[object_name] if days_since_update > notification_interval: # 准备插入通知表的数据,补充FILE_NAME字段(无数据则设为NULL) table_data = [ object_name, None, latest_timestamp, datetime.utcnow(), notification_interval, 'Not Updated', "The table '{0}' has not been updated for {1} days. Please take action!".format(object_name, days_since_update) ] print(table_data) # 使用参数绑定构建INSERT语句,避免语法错误和SQL注入 insert_query = """ INSERT INTO AW.NOTIFICATION_TABLE ( TABLE_NAME, FILE_NAME, LATEST_TIMESTAMP, SCAN_TIMESTAMP, MAX_AGE, STATUS, MESSAGE ) VALUES (?, ?, ?, ?, ?, ?, ?) """ # 执行插入操作 session.sql(insert_query, table_data).collect()
关键修复点说明
- 执行建表语句:在main函数开头添加
session.sql(create_table_query).collect(),确保通知表存在。 - 修正INSERT语法:移除重复的
VALUES子句,添加FILE_NAME列,使用参数绑定(?占位符)传递值。 - 触发SQL执行:所有DDL/DML语句后添加
.collect(),确保数据库执行相应操作。 - 处理FILE_NAME字段:补充该字段的值(这里设为NULL,可根据业务逻辑调整),保证列数匹配。
- 参数绑定:避免直接字符串格式化带来的语法错误和SQL注入风险,同时正确处理时间类型。
内容的提问来源于stack exchange,提问作者Karolina
相关产品推荐
相关产品推荐

