【问题标题】:How to copy some records from TableA to TableB with slick-streaming and akka-streaming如何使用 slick-streaming 和 akka-streaming 将一些记录从 TableA 复制到 TableB
【发布时间】:2017-11-09 19:36:34
【问题描述】:

有两个表TableATableB

我需要将一些记录从TableA 复制到TableB。我使用slick-3.0并使用以下方式:

import akka.stream._
import akka.stream.scaladsl._
...

//{{ READ DATA FROM TABLE A
val q = TableA.filter(somePredicate).result
val source = Source.fromPublisher {
      db.stream(q.result).mapResult { r => 
        val record: RecordA = someTransformation(r) 
        record
      }
    }.grouped(50) // grouping because I want to write records in batch mode
//}}

//{{ WRITE DATA TO TABLE B
val f:Future[Done] = source.runWith(Sink.foreach { batch: Seq[RecordA] =>
      //TODO how to write batch to TableB asynchronously?
      val insertAction = TableB ++= batch  // insert batch to table
      val fInsert: Future[_] = db.run(insertAction)
      Await.result(fInsert, ...)           // #1 this works only with blocking
})
//}}

但我遇到了一个问题 - 如何将批处理异步写入TableB(请参阅 TODO)。现在上面的代码只适用于阻塞到内部未来(见#1评论)。是否有异步执行该任务的正确方法?

【问题讨论】:

  • 如果你不阻止内心的未来会发生什么?
  • @thwiegan,如果我不阻止内在的未来并返回它,那么它就不会完成
  • 这似乎是您的用例:stackoverflow.com/questions/36400152/… 我看不出与您的示例有什么不同
  • @thwiegan,感谢您的回答。但不幸的是,这个例子也不起作用。它只有在内心的未来被阻塞时才有效。否则数据不会保存到TableB

标签: scala slick akka-stream reactive-streams


【解决方案1】:

使用mapAsync 它期望返回一个future,并在下一阶段公开“解包”结果。

source.mapAsync(4){batch: Seq[RecordA] =>
      val insertAction = TableB ++= batch  // insert batch to table
      db.run(insertAction)
}).to(Sink.ignore).run

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-06-08
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多