Scala/Spark如何提取RDD每行嵌套元组的第一个元素
Scala 实现对等PySpark逻辑的方案
你给出的PySpark逻辑是对键值对RDD的value部分(每个value为二元组列表),提取所有二元组的第一个元素组成新的列表,Scala Spark的实现代码如下:
// 构造和PySpark完全一致的输入RDD val rdd = sc.parallelize(Seq( (0, Seq((0, "a"), (1, "b"), (2, "c"))), (1, Seq((3, "x"), (5, "y"), (6, "z"))) )) // 核心转换逻辑,等价于PySpark的mapValues+itemgetter(0) val mapped = rdd.mapValues(_.map(_._1))
执行mapped.collect()得到的输出和PySpark运行结果完全一致:
Array((0,List(0, 1, 2)), (1,List(3, 5, 6)))
逻辑对应说明
- PySpark中
itemgetter(0)取元组第一个元素的操作,对应Scala中访问元组第一个元素的语法_._1 - Scala和PySpark的
mapValues算子功能完全一致,仅处理RDD每个键值对的value部分,不会修改原有key
内容的提问来源于stack exchange,提问作者Ged
相关产品推荐
相关产品推荐

