【发布时间】:2016-02-06 08:06:34
【问题描述】:
你好,有人可以向我解释为什么mapPartitions 对这两个函数的反应不同吗? (我看过这个this thread,但我认为我的问题不在于我创建它时我的可迭代对象是TraversableOnce。
L=range(10)
J=range(5,15)
K=range(8,18)
data=J+K+L
def function_1(iter_listoflist):
final_iterator=[]
for sublist in iter_listoflist:
final_iterator.append([x for x in sublist if x%9!=0])
return iter(final_iterator)
def function_2(iter_listoflist):
final_iterator=[]
listoflist=list(iter_listoflist)
for i in range(len(listoflist)):
for j in range(i+1,len(listoflist)):
sublist=listoflist[i]+listoflist[j]
final_iterator.append([x for x in sublist if x%9!=0])
pass
pass
return iter(final_iterator)
sc.parallelize(data,3).glom().mapPartitions(function_1).collect()
返回它应该同时
sc.parallelize(data,3).glom().mapPartitions(function_2).collect()
返回一个空数组,我通过在末尾返回一个列表来检查代码,它会执行我想要的操作。
感谢您的帮助
菲利普·C
【问题讨论】:
标签: python apache-spark pyspark partition