Snowflake Python存储过程使用filter时出现_udf_code报错
问题分析与解决方案
核心问题
SHOW STREAMS结果非结构化:SHOW STREAMS属于Snowflake元数据命令,返回的是半结构化输出,直接用Snowpark DataFrame的filter操作会因列名识别、数据结构不兼容触发错误。- 返回值类型不匹配:存储过程声明返回
STRING,但collect()返回的是Row对象列表,无法直接转换为字符串,导致类型错误。
修正后的存储过程代码
CREATE OR REPLACE PROCEDURE utility.procedure.RECREATE_STALE_STREAM_PYTHON() RETURNS STRING LANGUAGE PYTHON RUNTIME_VERSION = '3.8' PACKAGES = ('snowflake-snowpark-python') HANDLER = 'run' AS $$ from snowflake.snowpark.functions import col def run(session): # 执行SHOW命令后,通过RESULT_SCAN将输出转为结构化表 session.sql("SHOW STREAMS IN ACCOUNT;").collect() streams_df = session.sql("SELECT * FROM TABLE(RESULT_SCAN(LAST_QUERY_ID()))") # 过滤STALE状态为TRUE的流(元数据列默认大写) stale_streams = streams_df.filter(col('STALE') == 'TRUE').collect() # 将Row列表转换为字符串格式返回,符合存储过程返回类型要求 return "\n".join([str(row) for row in stale_streams]) $$;
关键修正说明
- 转换元数据输出为结构化表:使用
TABLE(RESULT_SCAN(LAST_QUERY_ID()))将SHOW命令的半结构化输出转换为标准可查询表,确保Snowpark能正确识别列名和数据结构。 - 匹配列名大小写:Snowflake元数据命令返回的列名默认大写(如
STALE),必须使用大写列名进行过滤,避免列不存在的错误。 - 统一返回值类型:将
collect()得到的Row列表转换为字符串格式,满足存储过程RETURNS STRING的类型要求。
内容的提问来源于stack exchange,提问作者Brandon Coleman
相关产品推荐
相关产品推荐

