【问题标题】:Porting a multi-threaded compute intensive job to spark移植多线程计算密集型作业以激发
【发布时间】: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


    【解决方案1】:

    所以我们可以做一些事情来让这段代码更像 Spark。首先是您使用了filtercount,但实际上使用了两者的结果。函数foreach 可能更接近你想要的。

    话虽如此,您正在创建一个数据库连接来存储结果,我们可以通过几种方式来实现这一点。一个是:如果数据库真的是您想要用于存储的东西,您可以使用foreachPartitionmapPartitionsWithIndex 为每个分区只创建一个连接,然后执行count()(我知道这有点难看,但@ 987654328@ 自 1.0.0 起已弃用)。您也可以只做一个简单的map,然后将您的结果保存为许多支持的输出格式之一(例如 saveAsSequenceFile)。

    【讨论】:

    • 谢谢霍尔顿!我也将dbconn 传递给runner,其中我的实际逻辑命中 db 并将一些数据存储在内存中(在 jvm 中)。它有一些遗留代码,我需要做一些实质性的重构来删除这种依赖关系并创建一个 RDD,而不是从 runner 中访问 db。
    • 好的,那么执行mapPartitionsWithIndex 允许您对每个分区中的所有元素重复使用 dbconn 可能是最好的方法。
    • 是的,刚刚试过。 mapPartitionsWithIndex 说它缺少一个返回语句(需要一个 RDD 迭代器)。但是我真的没有什么可以回报的,我该如何解决呢?
    • 啊,是的,只需返回一个空的迭代器。
    • 不确定我是否错过了这里的基本内容:markdownshare.com/view/c52a7ffb-27ef-44aa-978f-f4822ec876a3
    【解决方案2】:

    您可以尝试另一种方法来更好地并行化它,但要付出代价。代码在 scala 中,但 python 有 cartesian 方法。为简单起见,我的列表包含整数。

    val rdd1000 = sc.parallelize(list1000)
    val rdd10 = sc.parallelize(list10)
    val rdd2 = sc.parallelize(list2)
    
    rdd1000.cartesian(rdd10).cartesian(rdd2)
        .foreachPartition((tuples: Iterator[Tuple2[Tuple2[Int, Int], Int]]) => {
            dbconn =...
            for (tuple <- tuples) {
                val a = tuple._1._1
                val b = tuple._1._2
                val c = tuple._2
    
                val r = new Result(a, b, c, dbconn)
                dbconn.store(r)
            }
        })
    

    Filter 在您的情况下是一个 transformation,它是惰性的 - spark 不会在调用时对其进行评估。该过程仅在调用操作时开始。 Collect 是一个 action,它在您的示例中开始实际处理。 ForeachPartition 也是一个动作,spark 会立即启动它。 ForeachPartition 在这里是必需的,因为它允许为数据的整个分区打开一次连接。

    笛卡尔的坏处可能是它可能意味着对集群进行洗牌,所以如果你有复杂的对象,那可能会影响性能。如果您要从外部源读取数据,则可能会发生这种情况。如果您要使用并行化,那很好。

    还有一点需要注意的是,根据集群 Spark 的大小,可能会对您使用的数据库造成相当大的压力。

    【讨论】:

    • 这种情况,每次运行都会发生dbconn init 吗?
    • 一个缺点是笛卡尔是一项昂贵的操作(它会导致洗牌)。
    • @Holden,是的,它可能。我提到过。另一方面,我在结果 rdd 上打印了 debugString ,它显示没有洗牌。我认为这是因为一切都是本地的,尽管我将 spark master 指定为本地 [8]。所以,我会说这种方法很好,因为它使用标准 API,并且总是可以查看解释。
    • @zengr, dbconn init 在整个迭代器中只发生一次。
    • 好吧,它可以工作,但我现在很好奇,为所有列表创建 RDD 与循环 list10、list2 等有何不同? spark 是否更有效地处理任务分配?这个 impl 在性能方面与 @Hoden 的实现有何不同?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-05-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-03-12
    • 1970-01-01
    • 2015-05-19
    相关资源
    最近更新 更多