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

