如何将Snowpark中df_output写入Snowflake分区表并追加数据?
问题
我参考了一段用ARIMA实现多序列时间序列预测的代码(测试用第197行至末尾的代码),想把代码里第136行的df_output写入Snowflake分区表,实现每个序列的预测结果追加到表中。
试过两种方法都踩坑了:
- 用Snowpark的
write_pandas方法,但这个方法需要session对象,在UDTF的运行环境里拿不到; - 尝试通过
session.add_import导入sqlalchemy的连接配置,结果报了错:
类型错误:无法序列化'sqlalchemy.cprocessors.UnicodeResultProcessor'对象:你可能需要先将无法序列化的对象保存到本地环境,通过session.add_import()将其添加到UDF,然后从UDF中读取。
解决方案
推荐方案:UDTF输出结果,外部统一写入分区表
这是最符合Snowpark设计逻辑的做法,别在UDTF内部搞写入:
- 保持原UDTF的逻辑不变,把
df_output作为UDTF的返回结果输出就行(原代码已经定义了UDTF的返回结构,直接return即可); - 在调用UDTF的主程序里(也就是你用来测试的第197行及以后的代码),拿到UDTF的输出结果,再用Snowpark的写入方法把数据追加到分区表:
# 调用UDTF获取所有序列的预测结果 predictions_df = session.table("你的输入表名").select( col("SERIES_ID"), col("DATE"), col("VALUE"), # 调用自定义表函数,按序列分区计算 your_arima_udtf(col("SERIES_ID"), col("DATE"), col("VALUE")).over( partition_by="SERIES_ID" ).alias("预测结果") ).select("SERIES_ID", "DATE", "VALUE", "预测结果.*") # 追加写入分区表(假设分区键是SERIES_ID) predictions_df.write.mode("append").save_as_table("你的分区表名", partition_by=["SERIES_ID"])
备选方案:UDTF内部动态创建连接写入(不推荐)
如果一定要在UDTF里写表,别直接传连接对象,传连接参数字符串:
- 把Snowflake的连接参数(账户、仓库、数据库、Schema、用户名、密码)作为字符串参数传给UDTF;
- 在UDTF的
__call__方法里,每次需要写表时动态创建sqlalchemy连接:
# 在UDTF的__call__方法内部 from sqlalchemy import create_engine import pandas as pd # 从UDTF的参数中拿到连接信息(提前把这些参数作为UDTF的输入参数) engine = create_engine( f"snowflake://{用户名}:{密码}@{账户}/{数据库}/{Schema}?warehouse={仓库}" ) # 把df_output追加写入分区表 df_output.to_sql( name="你的分区表名", con=engine, if_exists="append", index=False )
⚠️ 注意:这种方法会给每个序列的计算实例都创建一个数据库连接,数据量大的时候会导致连接数爆炸,性能拉胯,只适合小批量场景。
为什么之前的方法报错?
sqlalchemy的连接对象包含底层的系统资源(比如网络连接句柄),这些东西没法被pickle序列化,所以不能通过session.add_import传递。你得传递能序列化的字符串参数,在UDTF内部自己创建连接。
内容的提问来源于stack exchange,提问作者janicebaratheon
相关产品推荐
相关产品推荐

