【发布时间】: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