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
相关产品推荐
相关产品推荐

