You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark字符串列转数值、分组计算均值及新增等级列实现

PySpark 无效值过滤、分组聚合、等级打标实现

首先导入需要依赖的PySpark函数:

from pyspark.sql.types import IntegerType, StringType
from pyspark.sql import functions as F

假设原始输入DataFrame变量名为df,按需求分三步实现如下:

  • 第一步:过滤subject_score为"NA"的无效记录,将有效分数转为整数类型,生成不覆盖原列的新列converted_score,仅保留ID和新分数列即可匹配Output1
df_output1 = df.filter(F.col("subject_score") != "NA") \
    .withColumn("converted_score", F.col("subject_score").cast(IntegerType())) \
    .select("ID", "converted_score")
  • 第二步:按ID字段分组,计算有效分数的平均值,结果取整为整数类型,命名为avg_subject_score即可匹配Output2
df_output2 = df_output1.groupBy("ID") \
    .agg(F.round(F.avg("converted_score")).cast(IntegerType()).alias("avg_subject_score"))
  • 第三步:基于分组聚合结果,按规则新增StringType类型的grade列,规则匹配后即可得到Output3
df_output3 = df_output2.withColumn(
    "grade",
    F.when(F.col("avg_subject_score") < 50, "Fail")
     .when(F.col("avg_subject_score").between(50, 99), "Good")
     .otherwise("Very Good")
     .cast(StringType())
).withColumnRenamed("ID", "id")

如果不需要保留中间结果,可以把三段逻辑合并为链式调用,一次执行得到最终结果:

df_final = df.filter(F.col("subject_score") != "NA") \
    .withColumn("converted_score", F.col("subject_score").cast(IntegerType())) \
    .groupBy("ID") \
    .agg(F.round(F.avg("converted_score")).cast(IntegerType()).alias("avg_subject_score")) \
    .withColumn(
        "grade",
        F.when(F.col("avg_subject_score") < 50, "Fail")
         .when(F.col("avg_subject_score").between(50, 99), "Good")
         .otherwise("Very Good")
         .cast(StringType())
    ).withColumnRenamed("ID", "id")

验证说明:执行对应DataFrame的.show()方法,即可得到对应阶段的输出结果,所有字段类型均符合要求:converted_score/avg_subject_score为整数类型,grade为字符串类型,原始输入列不会被修改。

内容的提问来源于stack exchange,提问作者AndyD

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.03 10:31:01