我是Spark的新手,我正在尝试用马尔科夫模型表示的质心实现一些迭代算法进行聚类(期望最大化).所以我需要进行迭代和连接.
我遇到的一个问题是每次迭代时间呈指数增长.
经过一些实验后,我发现在进行迭代时需要保留将在下一次迭代中重用的RDD,否则每个迭代spark都会创建执行计划,从开始重新计算RDD,从而增加计算时间.
init = sc.parallelize(xrange(10000000), 3)
init.cache()
for i in range(6):
print i
start = datetime.datetime.now()
init2 = init.map(lambda n: (n, n*3))
init = init2.map(lambda n: n[0])
# init.cache()
print init.count()
print str(datetime.datetime.now() - start)
Run Code Online (Sandbox Code Playgroud)
结果是:
0
10000000
0:00:04.283652
1
10000000
0:00:05.998830
2
10000000
0:00:08.771984
3
10000000
0:00:11.399581
4
10000000
0:00:14.206069
5
10000000
0:00:16.856993
Run Code Online (Sandbox Code Playgroud)
因此添加cache()会有所帮助,迭代时间也会变得不变.
init = sc.parallelize(xrange(10000000), 3)
init.cache()
for i in range(6):
print i
start = datetime.datetime.now()
init2 = init.map(lambda n: (n, n*3))
init = init2.map(lambda n: …Run Code Online (Sandbox Code Playgroud)