PySpark MapReduce:统计人员同日期共同出现次数的实现问题
解决Spark中人员共同出现日期次数统计问题
原始数据格式
每行是人员ID和日期,以空格分隔:
A 2022-08-13 B 2022-08-14 B 2022-08-13 A 2022-05-04 B 2022-05-04 C 2022-08-14 ...
现有代码问题
你当前的reduceByKey逻辑完全偏离目标:当前键是(人员, 日期),这种维度的reduce操作根本无法关联不同人员的共同日期。要实现统计两人共同出现的日期次数,核心思路应该是先按日期分组,再生成该日期下的人员对,最后统计每对的出现次数。
正确实现步骤
- 修正数据分割逻辑:原始数据用空格分隔,你之前的
split(',')是错误的,需改为空格分割。 - 按日期聚合人员:把同一天的所有人员归为一组。
- 生成无序人员对:对每个日期下的人员列表,生成所有不重复的无序对(比如
(A,B)和(B,A)视为同一对,避免重复统计)。 - 统计人员对出现次数:对生成的人员对做计数累加,得到最终结果。
完整代码实现
# 1. 读取原始数据并转换为(日期, 人员)键值对 worker_shifts = sc.textFile("你的数据文件路径") date_worker = worker_shifts.map(lambda x: (x.split()[1], x.split()[0])) # 2. 按日期分组,得到(日期, 当日人员列表) date_workers = date_worker.groupByKey() # 3. 生成每个日期下的无序人员对 def generate_pairs(workers): worker_list = list(workers) pairs = [] # 双重循环生成i<j的无序对,避免重复 for i in range(len(worker_list)): for j in range(i+1, len(worker_list)): # 按字典序排序,确保(A,B)和(B,A)对应同一个键 sorted_pair = tuple(sorted((worker_list[i], worker_list[j]))) pairs.append((sorted_pair, 1)) return pairs worker_pairs = date_workers.flatMap(lambda x: generate_pairs(x[1])) # 4. 统计每对人员的共同日期次数 common_days_count = worker_pairs.reduceByKey(lambda x, y: x + y) # 转换为期望的输出格式:(人员A, 人员B, 共同日期次数) final_result = common_days_count.map(lambda x: (x[0][0], x[0][1], x[1])) # 查看结果 final_result.collect()
关键细节说明
- 替换
split(',')为split():匹配原始数据的空格分隔规则,避免分割错误。 - 无序对生成:通过
sorted统一人员对的顺序,确保同一组人员不管顺序如何都被统计为同一个键。 - 按日期分组:这是实现需求的核心前提,只有把同一天的人员聚在一起,才能找到共同出现的人员组合。
内容的提问来源于stack exchange,提问作者sadmansh
相关产品推荐
相关产品推荐

