PySpark如何从DataFrame的JSON列表列中提取score生成独立列
PySpark 从JSON对象列表列提取score字段生成多列方案
适用场景:DataFrame中某列为JSON对象组成的列表,需要提取列表每个元素的
score字段,生成独立列存储不同位置的score值。
实现步骤
- 第一步:导入依赖并构造测试DataFrame
from pyspark.sql import SparkSession import pyspark.sql.functions as F # 初始化SparkSession spark = SparkSession.builder.appName("extract_score_from_json_array").getOrCreate() # 测试数据 raw_data = [{"user_id" : 1234, "col" : [{"id":14577120145280,"score":64.71,"Elastic_position":0},{"id":14568530280240,"score":88.53,"Elastic_position":1},{"id":14568530119661,"score":63.75,"Elastic_position":2},{"id":14568530205858,"score":62.79,"Elastic_position":3},{"id":14568530414899,"score":60.88,"Elastic_position":4}]}] df = spark.createDataFrame(raw_data)
- 第二步:(可选)JSON字符串列转结构化数组
如果col列存储的是JSON字符串而非原生结构化数组,需要先解析为数组类型:
# 定义数组元素schema json_schema = "array<struct<id:long, score:double, Elastic_position:int>>" df = df.withColumn("col", F.from_json(F.col("col"), json_schema))
- 第三步:动态提取各位置score生成独立列
# 获取数组最大长度,适配长度不固定的场景 max_array_len = df.agg(F.max(F.size("col"))).head()[0] # 批量构造score列 score_cols = [F.col("col")[i].score.alias(f"score_{i}") for i in range(max_array_len)] # 生成结果表 result_df = df.select("user_id", *score_cols)
- 最终输出效果
执行result_df.show()即可得到预期结果:
+-------+-------+-------+-------+-------+-------+ |user_id|score_0|score_1|score_2|score_3|score_4| +-------+-------+-------+-------+-------+-------+ | 1234| 64.71| 88.53| 63.75| 62.79| 60.88| +-------+-------+-------+-------+-------+-------+
内容的提问来源于stack exchange,提问作者Teresa
相关产品推荐
相关产品推荐

