【问题标题】:Scala Futures with DB带有 DB 的 Scala 期货
【发布时间】:2015-11-13 20:35:17
【问题描述】:

我正在 scala/play 中使用 anorm/postgres 编写代码,以根据用户配置文件生成匹配。以下代码有效,但我已注释掉导致问题的部分,即while 循环。我在运行它时注意到前 3 个 Future 似乎是同步工作的,但是当我在第四步中检索表中的行数时问题就来了。

第四步返回上述插入实际发生之前的计数。据我所知,步骤 1-3 正在同步排队等待 postgres,但检索计数的调用似乎在前 3 个步骤完成之前返回,这对我来说毫无意义。如果前 3 步以正确的顺序排队,为什么第四步不等到插入发生后才返回计数?

当我取消注释 while 循环时,将调用匹配生成和插入函数,直到内存用完,因为返回的计数一直低于所需的阈值。

我知道格式本身是不合格的,但我的问题不是关于如何编写最优雅的 scala 代码,而只是如何让它现在工作。

def matchGeneration(email:String,itNum:Int) = {
  var currentIterationNumber = itNum
  var numberOfMatches = MatchData.numberOfCurrentMatches(email)
    while(numberOfMatches < 150){
          Thread.sleep(25000)//delay while loop execution time
          generateUsers(email) onComplete {
            case(s) => {
               print(s">>>>>>>>>>>>>>>>>>>>>>>>>>>STEP 1")
               Thread.sleep(5000)//Time for initial user generation to take place
              genDemoMatches(email, currentIterationNumber) onComplete {
                case (s) => {
                  print(s">>>>>>>>>>>>>>>>>>>>>>>>>>>STEP 2")
                  genIntMatches(email,currentIterationNumber) onComplete {
                    case(s) => {
                      print(s">>>>>>>>>>>>>>>>>>>>>>>>>>>STEP 3")
                      genSchoolWorkMatches(email,currentIterationNumber) onComplete {
                        case(s) => {
                          Thread.sleep(10000)
                          print(s">>>>>>>>>>>>>>>>>>>>>>>>>>>STEP 4")
                          incrementNumberOfMatches(email) onComplete {
                            case(s) => {
                              currentIterationNumber+=1
                              println(s"current number of matches: $numberOfMatches")
                              println(s"current Iteration: $currentIterationNumber")
                            }
                          }
                        }
                      }
                    }
                  }
                }
              }
            }
          }
        //}
  }

匹配函数被定义为future,例如:

def genSchoolWorkMatches(email:String,currentIterationNumber:Int):Future[Unit]=
Future(genUsersFromSchoolWorkData(email, currentIterationNumber))

genUsersFromSchoolWorkData(email:String) 遵循与其他两个相同的形式。这是一个函数,最初获取用户在其个人资料 (SELECT major FROM school_work where email='$email') 中填写的所有学校/工作字段,并生成一个 dummyUser,其中包含与 email:String 的该用户相同的字段之一。打印这个函数大约需要 30-40 行代码,所以如果需要,我可以进一步解释。

我已经编辑了我的代码,到目前为止我发现让它工作的唯一方法是使用Thread.sleep() 破解它。我认为问题可能在于异常 正如我的 Future 逻辑构造确实按我预期的那样工作,但问题在于写入发生的时间与读取返回的时间不一致。 numberOfCurrentMatches(email:String) 函数返回匹配的数量,因为它是一个简单的 SELECT count(email) from table where email='$email'。问题是有时在插入 23 个匹配项后,计数返回为 0,然后在第二次迭代后它将返回 46。我假设 onComplete() 将绑定到用DB.withConnection() 定义的底层异常函数,但显然它也可能是远未实现这一目标。在这一点上,我不确定要研究什么或进一步查找以尝试解决此问题,而不是编写一种单独的主管函数以返回接近 150 的值。

更新

感谢用户的建议,并尝试通过以下链接了解 Scala 的文档:Scala Futures and Promises

我已经更新了我的代码,使其更具可读性和 Scala 风格:

  def genMatchOfTypes(email:String,iterationNumber:Int) = {
    genDemoMatches(email,iterationNumber)
    genIntMatches(email,iterationNumber)
    genSchoolWorkMatches(email,iterationNumber)
  }
  def matchGeneration(email:String) = {
  var currentIterationNumber = 0
  var numberOfMatches = MatchData.numberOfCurrentMatches(email)
  while (numberOfMatches < 150) {
    println(s"current number of matches: $numberOfMatches")
    Thread.sleep(30000)
    generateUsers(email)
      .flatMap(users =>  genMatchOfTypes(email,currentIterationNumber))
      .flatMap(matches => incrementNumberOfMatches(email))
      .map{
          result =>
            currentIterationNumber += 1
            println(s"current Iteration2: $currentIterationNumber")
            numberOfMatches = MatchData.numberOfCurrentMatches(email)
            println(s"current number of matches2: $numberOfMatches")
        }
    }

  }

我仍然严重依赖 Thread.sleep(30000) 来提供足够的时间来运行 while 循环,然后再尝试再次循环。它仍然是一个笨拙的黑客。当我取消注释 Thread.sleep()

我在 bash 中的输出如下所示:

users for match generation createdcurrent number of matches: 0
[error] c.MatchDataController - here is the list: jnkj
[error] c.MatchDataController - here is the list: hbhjbjjnkjn
current number of matches: 0
current number of matches: 0
current number of matches: 0
current number of matches: 0
current number of matches: 0

这当然是截断的输出。它一遍又一遍地运行,直到我收到关于打开文件过多的错误并且 JVM/play 服务器完全崩溃。

【问题讨论】:

  • is now about 应该读 is not about ?
  • @Odomontois 是的,抱歉,打字很快会编辑
  • 重要问题。 currentIterationNumber 是否会在未来成功的分支中一直增加或可能存在某些条件?
  • 您可以通过foldandThen 链接期货。
  • 您能否在问题中包含函数genSchoolWorkMatchesincrementNumberOfMatches 的定义?

标签: postgresql scala future


【解决方案1】:

一种解决方案是使用Future.traverse 获取已知的迭代次数

暗示

object MatchData {
  def numberOfCurrentMatches(email: String) = ???
}

def generateUsers(email: String): Future[Unit] = ???
def incrementNumberOfMatches(email: String): Future[Int] = ???
def genDemoMatches(email: String, it: Int): Future[Unit] = ???
def genIntMatches(email: String, it: Int): Future[Unit] = ???
def genSchoolWorkMatches(email: String, it: Int): Future[Unit] = ???

你可以写这样的代码

def matchGeneration(email: String, itNum: Int) = {
  val numberOfMatches = MatchData.numberOfCurrentMatches(email)
  Future.traverse(Stream.range(itNum, 150 - numberOfMatches + itNum)) { currentIterationNumber => for {
    _ <- generateUsers(email)
    _ = print(s">>>>>>>>>>>>>>>>>>>>>>>>>>>STEP 1")
    _ <- genDemoMatches(email, currentIterationNumber)
    _ = print(s">>>>>>>>>>>>>>>>>>>>>>>>>>>STEP 2")
    _ <- genIntMatches(email, currentIterationNumber)
    _ = print(s">>>>>>>>>>>>>>>>>>>>>>>>>>>STEP 3")
    _ <- genSchoolWorkMatches(email, currentIterationNumber)
    _ = Thread.sleep(15000)
    _ = print(s">>>>>>>>>>>>>>>>>>>>>>>>>>>STEP 4")
    numberOfMatches <- incrementNumberOfMatches(email)
    _ = println(s"current number of matches: $numberOfMatches")
    _ = println(s"current Iteration: $currentIterationNumber")
  } yield ()
  }

更新

如果您每次都敦促检查某些条件,一种方法是使用来自scalaz library 的很酷的 monadic 东西。它有 scala.Future 的 monad 定义,因此我们可以在需要时将单词 monadic 替换为 asynchronous

例如StreamT.unfoldM 可以创建条件单子(异步)循环,即使我们不需要结果集合的元素,我们仍然可以将其仅用于迭代。

让我们定义你的

def generateAll(email: String, iterationNumber: Int): Future[Unit] = for {
  _ <- generateUsers(email)
  _ <- genDemoMatches(email, iterationNumber)
  _ <- genIntMatches(email, iterationNumber)
  _ <- genSchoolWorkMatches(email, iterationNumber)
} yield ()

然后迭代步骤

def generateStep(email: String, limit: Int)(iterationNumber: Int): Future[Option[(Unit, Int)]] =
  if (MatchData.numberOfCurrentMatches(email) >= limit) Future(None)
  else for {
    _ <- generateAll(email, iterationNumber)
    _ <- incrementNumberOfMatches(email)
    next = iterationNumber + 1
  } yield Some((), next)

现在我们得到的函数简化为

import scalaz._
import scalaz.std.scalaFuture._

def matchGeneration(email: String, itNum: Int): Future[Unit] =
  StreamT.unfoldM(0)(generateStep(email, 150) _).toStream.map(_.force: Unit)

看起来同步方法MatchData.numberOfCurrentMatches 正在对incrementNumberOfMatches 中的异步修改做出反应。请注意,通常它可能会导致灾难性的结果,您可能需要将该状态移动到 actor 或类似的东西中

【讨论】:

  • 这是非常有用的信息,谢谢,但我的问题是我不知道每个用户需要多少次迭代才能为自己生成 150 个假匹配。如果用户 A 有 30 个要匹配的字段,则需要 5 次迭代,而如果用户 B 有 50 个要匹配的字段,则只需 3 次迭代。
  • 这看起来不可思议,但我需要一点时间来吸收......哈哈。我将不得不仔细研究一下!
  • 是的,我认为我的逻辑中的主要缺陷是变量计数器的使用以及不同线程之间变量的期货和不一致值的使用
猜你喜欢
  • 2017-08-19
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-05-17
  • 2019-03-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多