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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 03:18:11