如何使用PySpark统计各键对应列表长度及查询任意好友对共同好友
PySpark 好友关系统计实现
前置初始化
首先构造Spark上下文与输入数据RDD:
from pyspark import SparkContext # 初始化本地运行的SparkContext sc = SparkContext("local", "FriendStatsDemo") # 输入原始数据:键为用户ID,值为该用户的好友列表 raw_records = [('a1', ['b1', 'c1', 'd1', 'e1']), ('a2', ['b1', 'c2', 'd2', 'e1']), ('a3', ['b1', 'c2', 'd1', 'e2'])] friend_rdd = sc.parallelize(raw_records)
需求1:统计每个用户的好友总数量
直接对每个用户的好友列表计算长度即可:
user_friend_count = friend_rdd.map(lambda x: (x[0], len(x[1]))) # 输出结果 print(user_friend_count.collect())
运行输出:
[('a1', 4), ('a2', 4), ('a3', 4)]
需求2:查询任意一对用户的共同好友列表
实现逻辑:为每个用户的好友列表生成所有两两无序好友对,将当前用户作为该好友对的共同好友候选,最后分组聚合得到每对好友的共同好友列表:
def build_friend_pair(item): current_user = item[0] friends = item[1] result = [] # 生成所有两两好友组合,排序后作为key避免(A,B)和(B,A)被识别为不同对 for i in range(len(friends)): for j in range(i+1, len(friends)): sorted_pair = tuple(sorted([friends[i], friends[j]])) result.append((sorted_pair, current_user)) return result common_friends = friend_rdd.flatMap(build_friend_pair) \ .groupByKey() \ .mapValues(list) # 输出结果 print(common_friends.collect())
部分运行输出示例:
[(('b1', 'e1'), ['a1', 'a2']), (('c2', 'd1'), ['a3']), (('b1', 'c1'), ['a1'])]
内容的提问来源于stack exchange,提问作者behnaz.sheikhi
相关产品推荐
相关产品推荐

