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

PySpark中如何按Key分组RDD并统计唯一字符串频次,获取用户最受欢迎流派?

嘿,我来帮你搞定这个PySpark RDD的统计需求!你要的是先按用户ID统计各流派的出现次数,再找出每个用户出现次数最多的流派对吧?咱们一步步来实现:

首先,先把你的初始数据转换成RDD,方便后续操作:

from pyspark import SparkContext
sc = SparkContext.getOrCreate()

# 你的原始数据
data = [(1, "Western"), (1, "Western"), (1, "Drama"), (2, "Western"), (2, "Romance"), (2, "Romance")]
rdd = sc.parallelize(data)

第一步:统计每个用户各流派的出现次数

你之前尝试的sortByKey().countByValue()不太对,因为countByValue是统计整个RDD里每个完整元素的出现次数,没法帮你按用户分组统计流派。正确的做法是先把每个元素转换成((用户ID, 流派), 1)的格式,再用reduceByKey求和,这样就能得到每个用户每个流派的次数:

# 转换键值对,然后按(user, genre)聚合计数
genre_count_rdd = rdd.map(lambda x: ((x[0], x[1]), 1)).reduceByKey(lambda a, b: a + b)
# 此时RDD的内容是:((1, 'Western'), 2), ((1, 'Drama'), 1), ((2, 'Western'), 1), ((2, 'Romance'), 2)

如果想要得到你说的1: {"Western":2, "Drama":1}这种字典格式,可以把结构再转换一下,按用户分组后转成字典:

# 转换为(用户ID, (流派, 次数))的格式,再分组转字典
user_genre_dict_rdd = genre_count_rdd.map(lambda x: (x[0][0], (x[0][1], x[1])))\
                                    .groupByKey()\
                                    .mapValues(lambda genre_counts: dict(genre_counts))

# 查看结果
print(user_genre_dict_rdd.collect())
# 输出:[(1, {'Western': 2, 'Drama': 1}), (2, {'Western': 1, 'Romance': 2})]

第二步:筛选每个用户出现次数最多的流派

接下来,我们要从每个用户的流派计数里找出次数最高的那个。还是基于刚才的genre_count_rdd,转换结构后按用户分组,然后在每个组里用max函数找出次数最大的流派:

# 转换为(用户ID, (流派, 次数)),分组后取次数最多的流派
user_top_genre_rdd = genre_count_rdd.map(lambda x: (x[0][0], (x[0][1], x[1])))\
                                   .groupByKey()\
                                   .mapValues(lambda genre_list: max(genre_list, key=lambda item: item[1])[0])

# 查看结果
print(user_top_genre_rdd.collect())
# 输出:[(1, 'Western'), (2, 'Romance')]

核心思路是先按(用户ID, 流派)作为唯一key统计次数,再转换结构进行分组处理——而直接用countByValue更适合统计全局元素的出现频率,没法满足你分组统计的需求。

内容的提问来源于stack exchange,提问作者michael green

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 06:59:47