【问题标题】:Scala batch jobs in parallelScala 批处理作业并行
【发布时间】:2014-09-10 00:27:21
【问题描述】:

我在 Scala 中有一个可以进行批处理的函数。它被设计为自同步,因此您可以在一台机器或一个集群上运行它的 1 个或 1000 个实例,并且它将使用外部中间件进行同步。

为了提高性能,我希望一个 JVM 在多个线程中运行该函数。 (我希望在同一个 JVM 中完成以节省 RAM)。理想情况下,代码应如下所示:

execInParallel(9, myBatchFunction) // Starts 9 threads and invokes myBatchFunction() in each one

有什么简单的方法可以做到这一点?

【问题讨论】:

    标签: scala concurrency parallel-processing


    【解决方案1】:

    非阻塞版本:

    import scala.concurrent._
    import java.util.concurrent._
    import collection.JavaConverters._
    
    def execInParallel[T](numberThreads: Int, body: => T): util.List[Future[T]] = {
      val javaExecutor = Executors.newFixedThreadPool(numberThreads) // fixed thread pool
      val collections = Seq.fill(numberThreads) {
        new Callable[T]() {
          def call = body
        }
      }
      val futures = javaExecutor.invokeAll(collections.asJavaCollection) // run pool, first convert Seq to java.util.Collection
      // Here you have to be sure, that all task run
      import ExecutionContext.Implicits.global
      concurrent.Future(javaExecutor.shutdown()) // shutdown in new thread
      futures // return java futures !!!
    }
    
    val futures = execInParallel(9, Thread.currentThread.getName)
    println("Hurray")
    futures.asScala.foreach(x => println(x.get))
    

    【讨论】:

    • 代码看起来很棒,但我的目标是学习,而不仅仅是获取代码。你能解释一下你做了什么以及它是如何工作的吗?
    • 检查 java.util.concurrent 包,或 RxJava(RxScala 是包装器)——在那里你会找到你需要的一切。
    • 是错字吗? “asJavaColRxScalalection”应该是“asJavaCollection”
    • 这是一个 java 解决方案,而不是 scala 解决方案。
    • Scala 包装了 java.util.concurrent。 Scala 在 jvm 上运行。如果你想使用低级 api - 你必须使用 java api。
    【解决方案2】:

    这可能有效:

     (1 to 10).toList.par.foreach(myBatchFunction)
    

    它将调用 myBatchFunction 10 次,允许 Scala 确定何时并行执行它们。

    【讨论】:

    • 不确定是否会创建10个线程,可能与核心数有关
    猜你喜欢
    • 2012-11-12
    • 2014-09-28
    • 2013-10-21
    • 1970-01-01
    • 1970-01-01
    • 2023-04-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多