PySpark RDD去重需求:基于前两列匹配保留第三列最大值
解决PySpark RDD按前两列分组取第三列最大值的问题
没问题,我来帮你搞定这个需求!你需要的是将RDD中前两列相同的行归为一组,然后保留每组里第三列最大的那个值,具体实现步骤如下:
步骤1:准备原始数据并创建RDD
首先我们把你提供的原始数据转换成PySpark RDD的元素,注意要把第三列转换成浮点型,这样才能正确比较大小:
from pyspark import SparkContext sc = SparkContext("local", "MaxValueByGroup") # 原始数据,按行拆分 raw_data = [ "2,10482422,0.18", "2,10482422,0.4", "2,10482423,0.15", "2,10482423,0.43", "2,10482424,0.18", "2,10482424,0.49", "2,10482425,0.21", "2,10482425,0.52", "2,10482426,0.27", "2,10482426,0.64", "2,10482427,0.73" ] # 创建RDD并转换为( (col1, col2), col3 )的结构,col3转成float rdd = sc.parallelize(raw_data).map(lambda x: ((x.split(',')[0], x.split(',')[1]), float(x.split(',')[2])))
步骤2:按前两列分组并取第三列最大值
这里推荐使用reduceByKey操作,它比groupByKey更高效(会在节点本地先做聚合,减少shuffle的数据量),对每个键(前两列的组合)对应的所有值取最大值:
# 对每个键对应的第三列值取最大值 max_rdd = rdd.reduceByKey(lambda a, b: a if a > b else b)
步骤3:转换回目标格式并输出
最后把键值对结构转换成你需要的三列格式,再转成字符串输出:
# 转换为( col1, col2, max_col3 )的格式,再转成字符串 result_rdd = max_rdd.map(lambda x: f"{x[0][0]},{x[0][1]},{x[1]}") # 打印结果 for line in result_rdd.collect(): print(line)
运行结果
执行上面的代码后,你会得到期望的输出:
2,10482422,0.4 2,10482423,0.43 2,10482424,0.49 2,10482425,0.52 2,10482426,0.64 2,10482427,0.73
如果你的原始RDD已经是拆分好的元组(不是字符串),那可以跳过第一步的字符串拆分,直接用前两个元素作为键就行啦~
内容的提问来源于stack exchange,提问作者Sami
相关产品推荐
相关产品推荐

