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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 03:10:42