如何将DataFrame中分数字符串列转换为struct列表或动态扩展列
可行实现方案(基于PySpark,Spark SQL逻辑可直接复用)
方案一:转换为struct数组(推荐,最适合后续聚合计算)
该方案不需要处理动态列兼容问题,后续统计平均值、最高分、最低分等操作均可直接通过Spark内置函数完成,性能和易用性远优于动态宽表方案。
实现代码如下:
from pyspark.sql import functions as F from pyspark.sql.types import StructType, StructField, StringType, FloatType # 定义成绩struct的结构 score_schema = StructType([ StructField("student", StringType(), nullable=True), StructField("score", FloatType(), nullable=True) ]) # 转换为struct数组 df_struct = df \ # 按冒号拆分得到单个成绩项数组 .withColumn("score_list", F.split(F.col("scores"), ":")) \ # 遍历每个成绩项转换为struct .withColumn("scores", F.transform( F.col("score_list"), lambda item: F.struct( F.split(item, ",")[0].alias("student"), F.split(item, ",")[1].cast(FloatType()).alias("score") ) )) \ .drop("score_list") # 测试:直接计算ID为0000000003的平均分数 target_avg = df_struct \ .filter(F.col("id") == "0000000003") \ .select(F.expr("aggregate(scores, 0.0, (acc, x) -> acc + x.score, acc -> acc / size(scores))").alias("avg_score")) \ .first()["avg_score"]
执行后df_struct即为你要求的struct列表格式输出。
方案二:转换为动态扩展列(仅当业务强制要求宽表时使用)
该方案通过行转列实现动态多列,但如果不同ID对应的成绩项数量差异较大,会产生大量空值列,占用额外存储,后续聚合计算也更复杂。
实现代码如下:
from pyspark.sql import functions as F from pyspark.sql import Window df_wide = df \ .withColumn("score_list", F.split(F.col("scores"), ":")) \ # 拆分为单行单成绩项 .withColumn("score_item", F.explode(F.col("score_list"))) \ # 给同一ID下的成绩项编序号 .withColumn("item_index", F.row_number().over(Window.partitionBy("id").orderBy(F.lit(1)))) \ # 拆分姓名和分数 .withColumn("student", F.split(F.col("score_item"), ",")[0]) \ .withColumn("score", F.split(F.col("score_item"), ",")[1].cast(FloatType())) \ # 按序号行转列 .groupBy("id") \ .pivot("item_index") \ .agg( F.first("student").alias("student"), F.first("score").alias("score") )
执行后df_wide的列会自动生成为1_student、1_score、2_student、2_score格式,和你要求的动态列输出一致。
内容的提问来源于stack exchange,提问作者219CID
相关产品推荐
相关产品推荐

