【发布时间】:2017-11-09 19:36:34
【问题描述】:
有两个表TableA 和TableB。
我需要将一些记录从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