基于PySpark的加权标签得分计算及目标周标签选择问题
PySpark 百万级DataFrame标签得分计算与最优标签选择
原始DataFrame
+-----+-------+---+ | name|ts_week|tag| +-----+-------+---+ | Bob| week1| a| | Bob| week1| b| | Bob| week1| c| | Bob| week2| a| | Bob| week2| b| | Bob| week2| d| | Bob| week3| c| | Bob| week3| d| | Bob| week4| a| | Bob| week4| d| |Allen| week1| a| |Allen| week2| c| |Allen| week3| a| |Allen| week3| b| |Allen| week4| | +-----+-------+---+
评分规则
- 统计范围限定为week1-week4(一个月时间窗口)
- 各周标签对应分值:
- week1:0.1分
- week2:0.2分
- week3:0.3分
- week4:0.1分
预期标签得分计算结果
Bob: a= 0.7 b= 0.3 c= 0.4 d= 0.7 Allen: a= 0.4 b= 0.3 c= 0.2
需求说明
需要为每个用户的目标周(week4)选择得分最高的标签,且DataFrame规模达百万行级别,必须避免使用PySpark pandas(防止出现OOM内存溢出问题),需采用纯PySpark分布式处理方案。
原始数据生成代码
data_ls = [('Bob', 'week1', 'a'), ('Bob', 'week1', 'b'), ('Bob', 'week1', 'c'), ('Bob', 'week2', 'a'), ('Bob', 'week2', 'b'), ('Bob', 'week2', 'd'), ('Bob', 'week3', 'c'), ('Bob', 'week3', 'd'), ('Bob', 'week4', 'a'), ('Bob', 'week4', 'd'), ('Allen', 'week1', 'a'), ('Allen', 'week2', 'c'), ('Allen', 'week3', 'a'), ('Allen', 'week3', 'b'), ('Allen', 'week4', '')] data_sdf = spark.sparkContext.parallelize(data_ls).toDF(['name', 'ts_week', 'tag'])
纯PySpark解决方案
步骤1:映射周数到对应分值并过滤空标签
先过滤无效的空标签,再通过条件映射将周数转换为对应的分值:
from pyspark.sql import functions as F from pyspark.sql.window import Window score_sdf = data_sdf.filter(F.col("tag") != "") \ .withColumn("score", F.when(F.col("ts_week") == "week1", 0.1) .when(F.col("ts_week") == "week2", 0.2) .when(F.col("ts_week") == "week3", 0.3) .when(F.col("ts_week") == "week4", 0.1) .otherwise(0) )
步骤2:按用户和标签聚合计算总得分
通过分组聚合,计算每个用户每个标签的累计得分:
tag_total_score_sdf = score_sdf.groupBy("name", "tag") \ .agg(F.sum("score").alias("total_score"))
步骤3:筛选每个用户得分最高的标签
使用窗口函数按用户分组,对标签得分降序排序,取排名第一的标签(若存在并列得分,可根据业务需求调整排序逻辑):
# 定义窗口:按用户分组,按总得分降序排列 user_window = Window.partitionBy("name").orderBy(F.desc("total_score")) # 添加排名列并筛选Top1标签 result_sdf = tag_total_score_sdf.withColumn("rank", F.row_number().over(user_window)) \ .filter(F.col("rank") == 1) \ .select("name", "tag", "total_score")
查看最终结果
result_sdf.show()
输出示例:
+-----+---+-----------+ | name|tag|total_score| +-----+---+-----------+ |Allen| a| 0.4| | Bob| a| 0.7| | Bob| d| 0.7| +-----+---+-----------+
- 若需保留所有并列最高得分的标签,可将
row_number()替换为rank()或dense_rank(),此时会输出所有并列第一的标签。
性能说明
该方案完全基于PySpark分布式API实现,所有计算都在集群节点分布式执行,不会将大量数据拉取到Driver端,能够高效处理百万行级别的DataFrame,从根源上避免OOM问题。
内容的提问来源于stack exchange,提问作者mark
相关产品推荐
相关产品推荐

