Snowflake/Snowpark中动态展平JSON并转置为列的技术问题
问题描述
我需要在Snowflake中展平名为EVENT_PARAMS_JSON的JSON类型列。执行查询select EVENT_PARAMS_JSON from GA4_EVENT_DETAILS limit 1;得到的JSON数组结构如下(每个元素包含"key"字段和带有多类型值的"value"字段,且仅一个值非空):
[ {"key": "param1", "value": {"string_value": "abc"}}, {"key": "param2", "value": {"int_value": 123}}, {"key": "param3", "value": {"float_value": 4.56}} ]
我希望最终输出表将这些key作为列名,对应的非空value作为行值。但用Snowpark实现时出现报错:
100357 (P0000): Python Interpreter Error: raise error_class( snowflake.snowpark.exceptions.SnowparkSQLException: (1304): 01b0d218-0001-087c-0000-24952743d6f2: 001007 (22023): SQL compilation error: invalid type [VARCHAR(16777216)] for parameter '1'
使用的Snowpark代码如下:
import snowflake.snowpark as snowpark from snowflake.snowpark.functions import col def main(session: snowpark.Session): tableName = 'GA4_EVENT_DETAILS' df = session.table(tableName).filter( (col("app_info_id") == 'au.com.caltex.flagship') & (col("event_name") == 'screen_view') & (col("event_date") == '20230602') & (col("EVENT_PARAMS_FIREBASE_SCREEN_CLASS")=='View') ) df.show() df = df.join_table_function("flatten", df["EVENT_PARAMS_JSON"]).drop(["SEQ", "PATH", "INDEX", "THIS"])
解决方案
1. 修复Flatten调用错误
报错原因是join_table_function调用flatten时参数格式不正确,Snowpark中需通过session.table_function传递flatten的输入参数,同时要提取JSON中的key和对应非空值:
2. 动态行转列实现需求
要把key转为列名,需先提取所有唯一的key,再通过pivot完成转置。完整实现代码如下:
import snowflake.snowpark as snowpark from snowflake.snowpark.functions import col, coalesce def main(session: snowpark.Session): tableName = 'GA4_EVENT_DETAILS' df = session.table(tableName).filter( (col("app_info_id") == 'au.com.caltex.flagship') & (col("event_name") == 'screen_view') & (col("event_date") == '20230602') & (col("EVENT_PARAMS_FIREBASE_SCREEN_CLASS")=='View') ) # 展开JSON数组,提取key和对应非空值 flattened_df = df.join_table_function( session.table_function("flatten", input=col("EVENT_PARAMS_JSON")) ).with_column("param_key", col("VALUE")["key"]) \ .with_column("param_value", coalesce( col("VALUE")["value"]["string_value"], col("VALUE")["value"]["int_value"].cast("STRING"), col("VALUE")["value"]["float_value"].cast("STRING") )) \ .drop("SEQ", "PATH", "INDEX", "THIS", "VALUE") # 获取所有唯一参数key,用于转置列 param_keys = [row["PARAM_KEY"] for row in flattened_df.select("param_key").distinct().collect()] # 转置为宽表,key作为列名,对应value作为行值 pivot_df = flattened_df.pivot("param_key", param_keys).agg(col("param_value").alias("")) pivot_df.show() return pivot_df
关键说明
- 使用
coalesce从多类型的value字段中提取非空值,统一转为字符串类型保证列类型一致 - 通过
pivot实现行转列,需先获取所有唯一的key值作为目标列名 - 调用
flatten时通过session.table_function传递参数,解决原代码中的类型不匹配报错
内容的提问来源于stack exchange,提问作者Meg
相关产品推荐
相关产品推荐

