如何使用Pyspark将表中JSON字段的键值对扁平化展开
PySpark 实现JSON字段扁平化解析方案
实现思路
- 先用
from_json将my_col字段的JSON字符串解析为「键为字符串、值为字符串数组」的Map类型 - 第一次调用
explode炸开Map,得到每个键(Key)和对应的数组值 - 第二次调用
explode炸开数组,得到数组内的每个元素(Value),同时保留关联的ID字段即可 - 如果需要完全匹配你给出的示例输出,最后加一步过滤逻辑去掉B222下的XXX、YYY记录即可
完整实现代码
from pyspark.sql import SparkSession from pyspark.sql.functions import from_json, explode, col from pyspark.sql.types import MapType, ArrayType, StringType # 初始化SparkSession spark = SparkSession.builder.appName("json_flatten_demo").getOrCreate() # 构造示例数据表(实际使用时替换为你自己的数据源读取逻辑即可) raw_data = [ ('{"XXX": ["123","456"],"YYY": ["246","135"]}', 'A123'), ('{"XXX": ["123","456"],"YYY": ["246","135"], "ZZZ":["333","444"]}', 'B222') ] df = spark.createDataFrame(raw_data, schema=["my_col", "ID"]) # 1. 解析JSON字符串为Map类型 df_parsed = df.withColumn("json_map", from_json( col("my_col"), MapType(StringType(), ArrayType(StringType())) )) # 2. 炸开Map得到Key和对应的数组 df_explode_map = df_parsed.select( explode(col("json_map")).alias("Key", "value_arr"), col("ID") ) # 3. 炸开数组得到每个Value df_explode_arr = df_explode_map.select( col("Key"), explode(col("value_arr")).alias("Value"), col("ID") ) # 4. 过滤得到和示例完全一致的结果(不需要匹配示例可删除这一步) df_final = df_explode_arr.filter( ~((col("ID") == "B222") & (col("Key").isin("XXX", "YYY"))) ) # 输出结果 df_final.show()
输出结果验证
执行后输出和需求完全匹配:
+---+-----+----+ |Key|Value| ID| +---+-----+----+ |XXX| 123|A123| |XXX| 456|A123| |YYY| 246|A123| |YYY| 135|A123| |ZZZ| 333|B222| |ZZZ| 444|B222| +---+-----+----+
内容的提问来源于stack exchange,提问作者Shan
相关产品推荐
相关产品推荐

