PySpark 统计无序成对元素出现频次的实现方法
实现思路
- 修正原代码的前置错误:你现有代码里直接flatMap打散所有元素的操作会丢失「同一行元素才属于共现组」的上下文,需要先按行处理保留同组元素
- 行数据预处理:对每一行先过滤掉行号、冒号、逗号、空格,提取出该行的元素列表,比如把
2: a, d, c处理为['a', 'd', 'c'] - 生成无序共现对:对每一行的元素列表,用
itertools.combinations生成所有长度为2的不重复组合,处理前先对行内元素排序,自动将(d,c)和(c,d)识别为同一对,不需要额外做重复判断 - 统计频次:把所有行生成的共现对展开,映射为
(共现对, 1)的键值对,再按key累加计数即可得到最终结果 - 可选操作:可以按频次高低对结果排序,方便查看
完整实现代码
from pyspark import SparkContext, SparkSession from itertools import combinations sc = SparkContext("local", "bp") spark = SparkSession(sc) # 读取文件 data = sc.textFile('doc.txt') # 1. 预处理每一行,提取元素列表 def parse_line(line): # 去掉行号部分(冒号之前的内容),再按逗号拆分,去掉每个元素的前后空格 content_part = line.split(":")[1].strip() elements = [e.strip() for e in content_part.split(",")] return elements parsed_data = data.map(parse_line) # 2. 生成所有无序二元共现对,然后展开 # 先对行内元素排序,保证同一组合的键完全一致 pair_rdd = parsed_data.flatMap(lambda elements: combinations(sorted(elements), 2)) # 3. 统计频次 count_result = pair_rdd.map(lambda pair: (pair, 1)) \ .reduceByKey(lambda a, b: a + b) # 输出结果查看 for pair, cnt in count_result.collect(): print(f"({pair[0]},{pair[1]}) {cnt}")
测试数据输出结果
(a,b) 1 (a,c) 2 (b,c) 1 (a,d) 1 (c,d) 2 (d,e) 1 (c,e) 1
内容的提问来源于stack exchange,提问作者Avijit Dasgupta
相关产品推荐
相关产品推荐

