如何在PySpark RDD中找出不考虑顺序的双向关联对?
问题分析与代码修正
原代码存在的问题
- 输入处理逻辑错误:通过替换逗号和空格将所有字符合并,虽然对于单字符标识符看似能生成正确对,但这种方式不健壮(若标识符为多字符则完全失效),且逻辑上存在隐患。
- 核心逻辑错误:
findpairs函数设计为接收成对列表,但代码中通过filter将单个元组传入该函数,导致函数内部循环无法找到反向对,最终过滤掉所有元素,输出为空。 - 冗余导入:导入了
LinearRegression、SparkSession等无关模块,增加不必要的依赖。
修正后的代码实现
from pyspark import SparkContext sc = SparkContext("local", "DoubleRDD") def parse_line(line): # 跳过无箭头的无效行 if '->' not in line: return [] # 拆分源节点和目标节点部分 source, targets_str = line.split('->', 1) source = source.strip() # 拆分并清理每个目标节点 targets = [t.strip() for t in targets_str.split(',') if t.strip()] # 返回所有(源节点, 目标节点)元组 return [(source, target) for target in targets] # 读取文本文件 text_rdd = sc.textFile("path to the .txt") # 解析所有节点对 pairs_rdd = text_rdd.flatMap(parse_line) # 过滤掉自引用的对(如A->A) valid_pairs = pairs_rdd.filter(lambda x: x[0] != x[1]) # 将每个对转换为排序后的元组,便于统计双向关系 sorted_pairs = valid_pairs.map(lambda x: tuple(sorted(x))) # 统计每个排序后对的出现次数 pair_counts = sorted_pairs.countByKey() # 提取出现次数为2的对(即存在双向关系) mutual_pairs = [pair for pair, count in pair_counts.items() if count == 2] # 按期望格式输出结果 print("Output:") for pair in mutual_pairs: print(f"{pair[0]}, {pair[1]} //(this means {pair[0]} has supplied goods to {pair[1]} and {pair[1]} has also supplied some good to {pair[0]})") # 停止Spark上下文 sc.stop()
代码说明
- 输入解析:通过
parse_line函数正确拆分每行的源节点和目标节点,生成标准的(源, 目标)元组,确保处理逻辑健壮。 - 双向关系检测:将每个对转换为排序后的元组(如(K,M)和(M,K)都转为(K,M)),统计每个排序后对的出现次数。若次数为2,说明存在双向关系。
- 结果输出:按题目要求的格式打印所有双向关系对。
内容的提问来源于stack exchange,提问作者Salman Rasheed
相关产品推荐
相关产品推荐

