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

PySpark如何合并相同userId、movieId行的playtime值生成DataFrame

PySpark 双主键重复行合并playtime实现方案

核心逻辑:以userId、movieId两个字段作为分组键,对同分组下的playtime做聚合计算,即可得到去重合并后的结果DataFrame,常规播放时长场景默认使用求和聚合。

完整实现代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import sum

# 初始化Spark会话
spark = SparkSession.builder \
    .appName("merge_duplicate_play_record") \
    .getOrCreate()

# --------------------------
# 这里替换成你自己的 DataFrame 读取逻辑
# 示例为构造和参考样例结构一致的测试数据
raw_data = [
    (1, 101, 12),
    (1, 101, 8),
    (1, 102, 15),
    (2, 101, 20),
    (2, 102, 5),
    (2, 102, 9)
]
raw_df = spark.createDataFrame(raw_data, schema=["userId", "movieId", "playtime"])
# --------------------------

# 分组聚合合并重复行
result_df = raw_df.groupBy("userId", "movieId").agg(
    sum("playtime").alias("playtime")
)

# 输出查看处理结果
result_df.show()

运行结果

处理后输出的DataFrame内容如下,相同userId+movieId的行已经完成playtime合并:

+------+-------+--------+
|userId|movieId|playtime|
+------+-------+--------+
|     1|    101|      20|
|     1|    102|      15|
|     2|    101|      20|
|     2|    102|      14|
+------+-------+--------+

注意事项

  • 如果你的playtime合并规则不是累加,替换agg内的聚合函数即可:取同组最大播放时长用max("playtime"),需要拼接所有playtime记录用concat_ws(",", collect_list("playtime"))
  • 聚合操作默认仅保留分组字段和聚合字段,如果原始表存在其他字段,需要提前明确对应字段的聚合规则,否则相关字段不会出现在结果表中

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 07:34:23