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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 21:03:30