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

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操作根本无法关联不同人员的共同日期。要实现统计两人共同出现的日期次数,核心思路应该是先按日期分组,再生成该日期下的人员对,最后统计每对的出现次数。

正确实现步骤

  1. 修正数据分割逻辑:原始数据用空格分隔,你之前的split(',')是错误的,需改为空格分割。
  2. 按日期聚合人员:把同一天的所有人员归为一组。
  3. 生成无序人员对:对每个日期下的人员列表,生成所有不重复的无序对(比如(A,B)和(B,A)视为同一对,避免重复统计)。
  4. 统计人员对出现次数:对生成的人员对做计数累加,得到最终结果。

完整代码实现

# 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 16:10:41