PySpark如何实现RDD元素值与所属分区号相加
PySpark 分区内元素叠加分区编号实现
核心思路
要给每个元素加上所属分区的编号,需要使用PySpark RDD提供的mapPartitionsWithIndex算子——这个算子在处理分区数据时,会同时传入当前分区的索引编号,以及分区内的元素迭代器,刚好匹配需求。
你创建的4分区RDD默认会均匀切分0-99的数值:
- 分区0存储024,元素加0后结果为024
- 分区1存储2549,元素加1后结果为2650
- 分区2存储5074,元素加2后结果为5276
- 分区3存储7599,元素加3后结果为78102
切分计算后的结果和你给出的期望输出完全一致。
补全后可运行代码
A = sc.parallelize(range(100), 4) B = A.mapPartitionsWithIndex(lambda partition_id, elements: (ele + partition_id for ele in elements)) print(B.collect())
代码说明
- 传入
mapPartitionsWithIndex的处理函数接收两个参数:第一个是当前处理的分区编号,第二个是当前分区所有元素组成的迭代器 - 用生成器表达式遍历分区内元素做计算,不会把整个分区的数据一次性加载到内存,执行效率更高
内容的提问来源于stack exchange,提问作者Brandon Ong
相关产品推荐
相关产品推荐

