如何按索引关联RDD中每行的两个数组?
解决Spark RDD中按索引合并每行两个数组的问题
我来帮你搞定这个问题!你现在的需求是把RDD每行里的两个数组,按对应索引把Int和Double元素配对起来,但之前的代码只能拿到第一个索引的值,核心问题在于没有遍历每行数组的所有索引进行配对。
核心思路
对于每行的(Array[Int], Array[Double]),我们可以利用数组的zip方法,它会自动把两个数组对应索引的元素一一配对,而且只会处理到两个数组中较短的那个长度(正好适配你说的“长度相近”的场景,避免数组越界报错)。
具体实现代码
根据你的最终输出需求,分两种情况:
情况1:保留每行的配对结果为数组(输出RDD[Array[(Int, Double)]])
如果你想让每行还是一个数组,里面是对应索引配对后的元素,用map处理:
// 假设这是你的原始RDD val originalRDD: RDD[(Array[Int], Array[Double])] = ... val pairedByIndexRDD: RDD[Array[(Int, Double)]] = originalRDD.map { case (intArr, doubleArr) => // zip方法会把两个数组对应索引的元素配对 intArr.zip(doubleArr) }
比如原始行是([1,2,3], [0.1,0.2,0.3]),处理后会变成[(1,0.1), (2,0.2), (3,0.3)]。
情况2:将所有配对元素展平为RDD的单行(输出RDD[(Int, Double)])
如果你想把所有配对后的元素都作为RDD的独立行,用flatMap代替map:
val flattenedPairedRDD: RDD[(Int, Double)] = originalRDD.flatMap { case (intArr, doubleArr) => intArr.zip(doubleArr) }
上面的例子处理后,RDD会新增三行:(1,0.1)、(2,0.2)、(3,0.3)。
处理数组长度不一致的额外情况
如果遇到两个数组长度差异略大,你想保留较长数组里的多余元素(用默认值填充缺失的部分),可以用zipAll方法:
// 这里给多余的Int元素设默认值0,Double设默认值0.0,你可以根据需求修改 val pairedWithDefaultRDD: RDD[Array[(Int, Double)]] = originalRDD.map { case (intArr, doubleArr) => intArr.zipAll(doubleArr, 0, 0.0) }
为什么之前的代码只能拿到第一个索引?
应该是你之前的代码直接取了数组的第一个元素,比如(intArr(0), doubleArr(0)),这样只会获取索引0的配对值。而用zip会自动遍历所有索引,完成全部配对。
内容的提问来源于stack exchange,提问作者Sugimiyanto
相关产品推荐
相关产品推荐

