【发布时间】:2015-08-29 05:35:31
【问题描述】:
我正在编写使用执行器服务来并行化任务的代码(想想机器学习计算一遍又一遍地在小数据集上完成)。 我的目标是尽可能快地多次执行一些代码并将结果存储在某处(总执行量将至少为 100M 次)。
逻辑看起来像这样(它是一个简化的例子):
dbconn = new dbconn() //This is reused by all threads
for a in listOfSize1000:
for b in listofSize10:
for c in listOfSize2:
taskcompletionexecutorservice.submit(new runner(a, b, c, dbconn))
最后,taskcompletionexecutorservice.take() 被调用,我将“未来”的结果存储在数据库中。 但这种方法并没有真正在一个点之后进行缩放。
这就是我现在在 spark 中所做的事情(这是一个残酷的 hack,但我正在寻找有关如何最好地构建它的建议):
sparkContext.parallelize(listOfSize1000).filter(a -> {
dbconn = new dbconn() //Cannot init it outsize parallelize since its not serializable
for b in listofSize10:
for c in listOfSize2:
Result r = new runner(a, b, c. dbconn))
dbconn.store(r)
return true //It serves no purpose.
}).count();
这种方法在我看来效率低下,因为它并没有真正在最小的工作单元上进行并行化,尽管这项工作工作正常。另外 count 并没有真正为我做任何事情,我添加它来触发执行。它的灵感来自计算此处的 pi 示例:http://spark.apache.org/examples.html
那么对于如何更好地构建我的 spark runner 以便我可以有效地使用 spark executor 有什么建议吗?
【问题讨论】:
标签: java apache-spark