Snowpark中用循环动态扁平化JSON:匹配同层级对象
解决Snowpark DataFrame动态扁平化JSON列的问题
你的核心问题是循环中每次select都会覆盖结果,而且没有实现每行对应列表中元素的动态匹配——因为list_1和list_2是和DataFrame行数等长的,需要让每行根据自身对应的transport子对象名称和要提取的字段名去解析JSON,而不是全局循环取同一个字段。
修正思路
- 先把两个列表作为新列添加到DataFrame中,让每行都能拿到对应的
transport子对象名称和要提取的字段名。 - 利用Snowpark支持用列作为JSON索引的特性,动态提取每行对应的JSON值。
修正后的代码
from snowflake.snowpark.functions import col, lit, array_construct, row_number # 原始DataFrame list_1 = ['plane', 'ship'] list_2 = ['flight_class', 'ship_cabin'] # 先添加行号列,用于绑定列表元素与DataFrame行 dataframe = dataframe.with_column("ROW_ID", row_number().over(order_by=lit(1))) # 将列表转为Snowpark数组列,通过行号匹配每行对应的元素 dataframe = dataframe.with_columns( [ ("TRANSPORT_TYPE", array_construct(*[lit(item) for item in list_1])[col("ROW_ID") - 1]), ("TARGET_FIELD", array_construct(*[lit(item) for item in list_2])[col("ROW_ID") - 1]) ] ) # 动态提取JSON中对应字段,生成扁平化结果 dataframe_return = dataframe.select( # 保留原表需要的列(可根据需求调整) col("*"), # 按行匹配提取JSON字段 col("STREAMING_DATA")["transport"][col("TRANSPORT_TYPE")][col("TARGET_FIELD")].alias("FLATTENED_RESULT") ) # 可选:移除临时添加的行号和辅助列 dataframe_return = dataframe_return.drop("ROW_ID", "TRANSPORT_TYPE", "TARGET_FIELD")
代码说明
row_number().over(order_by=lit(1))生成行号,实现Python列表元素与DataFrame行的一一绑定。array_construct(*[lit(item) for item in list_1])把Python列表转为Snowpark数组,再通过行号索引取出每行对应的元素。col("STREAMING_DATA")["transport"][col("TRANSPORT_TYPE")]用列值作为JSON键,动态获取每行对应的transport子对象,再通过TARGET_FIELD提取目标字段,完成扁平化。
内容的提问来源于stack exchange,提问作者CodingNewbie
相关产品推荐
相关产品推荐

