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

