【问题标题】:How to use mapPartitions in pyspark如何在 pyspark 中使用 mapPartitions
【发布时间】:2017-03-23 15:54:23
【问题描述】:

遵循 Apache Spark 文档后,我尝试使用 mapPartition 模块。在下面的代码中,我希望看到函数 myfunc 中的初始 RDD,我只是在打印值后返回迭代器。但是当我在 RDD 上执行 collect 时,它是空的。

from pyspark import SparkConf
from pyspark import SparkContext                                                          

def myfunc(it):
    print(it.next())
    return it

def fun1(sc):
    n = 5
    rdd = sc.parallelize([x for x in range(n+1)], n)
    print(rdd.mapPartitions(myfunc).collect())


if __name__ == "__main__":                                                                
    conf = SparkConf().setMaster("local[*]")                                              
    conf = conf.setAppName("TEST2")                                                       
    sc = SparkContext(conf = conf)                                                        
    fun1(sc)

【问题讨论】:

    标签: python apache-spark pyspark rdd itertools


    【解决方案1】:

    mapPartitions 在这里不相关。迭代器(这里是itertools.chain)是有状态的,只能被遍历一次。当您调用 it.next() 时,您会读取并丢弃第一个元素,返回的是序列的尾部。

    当分区只有一项时(应该是除一项之外的所有项目),您实际上丢弃了整个分区。

    一些注意事项:

    • 将任何内容放入任务中的标准输出通常是没有用的。
    • 您使用 next 的方式不可移植,不能在 Python 3 中使用。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2014-12-31
      • 1970-01-01
      • 2016-02-06
      • 1970-01-01
      • 2019-07-10
      • 1970-01-01
      • 1970-01-01
      • 2017-07-14
      相关资源
      最近更新 更多