如何移除PySpark RDD笛卡尔积结果中两个元素相同的元组对
解决方案
核心实现逻辑
直接对笛卡尔积返回的RDD调用filter()算子完成过滤即可,不需要把数据拉取到本地,整个操作完全分布式执行,适配十亿级以上数据量场景。
另外你当前代码里先collect()再sc.parallelize()的操作会把全量数据拉到Driver端,大数据量下极易触发OOM,直接删除这部分冗余操作即可。
完整可运行代码
def compute_cartesian(rdd): # 直接生成笛卡尔积,无需拉取到本地 cartesian_rdd = rdd.cartesian(rdd) # 过滤掉两个元素完全相等的元组对 result_rdd = cartesian_rdd.filter(lambda x: x[0] != x[1]) # 以下为调试输出,生产环境可删除 print(type(result_rdd)) print(result_rdd.collect()) return result_rdd # 测试调用 test = sc.parallelize([(1,0), (2,0), (3,0)]) compute_cartesian(test)
输出结果
<class 'pyspark.rdd.PipelinedRDD'> [((1, 0), (2, 0)), ((1, 0), (3, 0)), ((2, 0), (1, 0)), ((2, 0), (3, 0)), ((3, 0), (1, 0)), ((3, 0), (2, 0))]
补充说明
filter()是PySpark RDD的转换算子,会分布式遍历RDD的每一条元素,只保留满足判断条件的记录,全程不会把数据同步到本地Driver节点。- 你之前尝试的
distinct()作用是去除完全重复的整条记录,和当前过滤元组内部两个元素相等的场景不匹配,所以结果不符合预期;dropDuplicates()是DataFrame专属API,RDD本身没有该方法,所以调用会报错。
内容的提问来源于stack exchange,提问作者MarkS
相关产品推荐
相关产品推荐

