如何在Snowpark(Python)1.0.0中不使用explode()实现行展开功能
针对Snowpark Python 1.0.0版本实现类似explode的拆分方案
由于Snowpark Python 1.0.0版本未内置explode()函数,可借助Snowflake原生SQL表函数STRTOK_SPLIT_TO_TABLE实现单列表格拆分,步骤如下:
实现步骤
1. 修正原DataFrame创建(注意二维数组格式)
你的示例DataFrame创建需调整为二维数组,否则字段映射会出错:
df = session.create_dataframe( [["rest_of_the_row", "A|B|C"]], schema=["record", "product_code"] )
2. 添加临时唯一行号
为确保拆分后的行能正确关联回原数据,需给原DataFrame添加临时行标识:
from snowflake.snowpark.functions import row_number, lit from snowflake.snowpark.window import Window df_with_rowid = df.withColumn("ROW_ID", row_number().over(Window.order_by(lit(1))))
3. 调用Snowflake表函数拆分字符串
通过STRTOK_SPLIT_TO_TABLE按指定分隔符拆分product_code列:
# 创建临时视图供SQL查询使用 df_with_rowid.create_or_replace_temp_view("temp_split_source") # 执行拆分SQL并转为Snowpark DataFrame split_df = session.sql(""" SELECT t.ROW_ID, s.VALUE AS product_code FROM temp_split_source t, TABLE(STRTOK_SPLIT_TO_TABLE(t.product_code, '|')) s """)
4. 关联原数据并清理临时字段
将拆分结果与原DataFrame关联,去掉临时行号后得到目标输出:
final_df = df_with_rowid.join(split_df, on="ROW_ID").drop("ROW_ID") # 查看结果 final_df.show()
后续聚合支持
拆分后的final_df可直接进行product_code相关聚合操作,例如统计各代码出现次数:
from snowflake.snowpark.functions import count agg_result = final_df.groupBy("product_code").agg(count("record").alias("count")) agg_result.show()
内容的提问来源于stack exchange,提问作者Deepa KP
相关产品推荐
相关产品推荐

