如何在PySpark中实现两个RDD的按行一一映射?
解决RDD按行一一配对的问题
看起来你用join操作走偏啦,join是基于键值对的key匹配来关联数据的,而你想要的是按行位置一一对应,这时候应该用zip操作才对!
问题出在哪?
你之前写的rdd2.join(rdd1.map(lambda x: (x[0], x[0:])))有两个核心问题:
rdd2是纯文本RDD,不是键值对类型(没有key),而join要求两个输入RDD都必须是(key, value)格式;- 就算你把rdd1转成键值对,
join是找相同key的元素关联,不是按位置配对,所以自然找不到匹配项,返回空RDD。
正确的解决方案
要实现“第1行对应第1行”的配对,用zip就可以直接搞定,前提是两个RDD的元素数量相同、分区数一致:
- 先确认两个RDD的元素数量是否匹配(这是zip的必要条件):
print("rdd1元素数:", rdd1.count()) print("rdd2元素数:", rdd2.count())
- 直接用
zip配对:
paired_rdd = rdd1.zip(rdd2)
执行后,paired_rdd里的每个元素就是(rdd1的行元素, rdd2的对应行元素),比如你想要的示例输出:
(0, "i hate painting i have white paint all over my hands.") (0, "Bawww I need a haircut No1 could fit me in before work tonight. Sigh.") (4, "I had a great day") ...
- (可选)如果需要转换成键值对格式方便后续操作:
key_value_rdd = paired_rdd.map(lambda x: (x[0], x[1]))
注意事项
如果执行zip时报错说分区数不一致,先调整两个RDD的分区数再配对:
# 把rdd1的分区数调整成和rdd2一致 rdd1 = rdd1.coalesce(rdd2.getNumPartitions()) paired_rdd = rdd1.zip(rdd2)
内容的提问来源于stack exchange,提问作者Rahul Anand
相关产品推荐
相关产品推荐

