【问题标题】:Scala program using futures is not terminating使用期货的 Scala 程序没有终止
【发布时间】:2020-01-31 17:56:39
【问题描述】:

我正在尝试学习 Scala 中的并发性,并使用 Scala 期货生成带有随机字符串的数据集。我想创建一个应用程序,它应该生成一个包含任意数量记录的文件并且它应该是可扩展的。

代码:

import java.util.concurrent.{ExecutorService, Executors}
import scala.util.{Failure, Random, Success}
import scala.concurrent.duration._

  object datacreator {

    implicit val ec: ExecutionContext = new ExecutionContext {

    val threadPool: ExecutorService = Executors.newFixedThreadPool(4)

    def execute(runnable: Runnable) {
      threadPool.submit(runnable)
    }

    def reportFailure(t: Throwable) {}
  }

  def getRecord : String = {
    "Random string"
  }

  def main(args: Array[String]): Unit = {

    val filename = args(0)
    val number_of_records = args(1)
    val file_Object = new FileWriter(filename, true)

    val data: Future[Iterable[String]] = Future {
      for (i <- 1 to number_of_records.toInt)
        yield getRecord
    }

    val result = data.map{
      result => result.foreach(record => file_Object.write(record))
    }

    result.onComplete{
          case Success(value) => {
            println("Success")
            file_Object.close()
          }
          case Failure(e) => e.printStackTrace()
    }
  }
}

我面临以下问题:

  1. 当我使用 SBT 运行程序时,它会将结果写入文件,但不会因进入无限模式而终止。
[info] Loading project definition from /Users/cw0155/PersonalProjects/datagen/project
[info] Loading settings for project datagen from build.sbt ...
[info] Set current project to datagenerator (in build file:/Users/cw0155/PersonalProjects/datagen/)
[info] running com.generator.DataGenerator xyz.csv 100
Success
  | => datagen / Compile / runMain 255s
  1. 当我使用 Jar 运行程序时:

scala -cp target/scala-2.13/datagenerator_2.13-0.1.jar com.generator.DataGenerator "pqr.csv" "1000" 它正在等待无限时间并且不写入文件。

非常感谢任何帮助:)

【问题讨论】:

  • 第一条线索:result 是什么类型? (即Future 什么?)
  • 是的,它是Future[Unit]

标签: scala concurrency


【解决方案1】:

试试这个版本

bar.scala

import scala.concurrent.{Await, Future, ExecutionContext}
import scala.concurrent.duration._
import scala.util.{Success, Failure}
import ExecutionContext.Implicits.global
import java.io.FileWriter
object bar {
  def getRecord: String = "Random string\n"
  def main(args: Array[String]): Unit = {
    val filename = args(0)
    val number_of_records = args(1)
    val data: Future[Iterable[String]] = Future {
      for (i <- 1 to number_of_records.toInt)
        yield getRecord
    }
    val file_Object = new FileWriter(filename, true)
    val result      = data.map( r => r.foreach(record => file_Object.write(record)) )
    result.onComplete {
      case Success(value) =>
        println("Success")
        file_Object.close()
      case Failure(e) =>
        e.printStackTrace()
    }
    Await.result( result, 10.second )
  }
}

当我这样运行时,您的原始版本给了我预期的输出

bash-3.2$ scala bar.scala /dev/fd/1 10
Success
Random string
Random string
Random string
Random string
Random string
Random string
Random string
Random string
Random string
Random string

但是,如果没有 Await.result,您的程序可以在未来完成之前退出。

【讨论】:

  • 感谢您的回答 :) 如何知道超时的正确值,因为记录的数量可能会有所不同,假设它可能是数十亿。没有等待怎么办?我还看到您在创建全局执行上下文时导入了全局执行上下文,这就是问题的原因。
  • @Ajit K'sagar 不幸的是,没有很好的答案什么是正确的超时,这取决于您的问题的上下文。作为选项之一,您可以使用scala.concurrent.duration.Duration.Inf - 无限持续时间,这意味着您的应用程序只有在未来完成时才会退出,根本没有任何时间限制。或者使用您可以接受的最大超时时间 - 比如说 1 小时。
猜你喜欢
  • 2020-11-03
  • 1970-01-01
  • 2021-09-12
  • 1970-01-01
  • 2012-05-20
  • 2022-06-17
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多