【发布时间】:2016-05-06 09:22:24
【问题描述】:
软件版本:
- Akka 2.4.4
- 光滑 3.1.0
我想在 Slick 事务中处理来自 Akka 流的元素。 下面是一些简化的代码来说明一种可能的方法:
def insert(d: AnimalFields): DBIO[Long] =
animals returning animals.map(_.id) += d
val source: Source[AnimalFields, _]
val sourceAsTraversable = ???
db.run((for {
ids <- DBIO.sequence(sourceAsTraversable.map(insert))
} yield { ids }).transactionally)
到目前为止,我能想到的一个解决方案是阻止每个未来遍历元素:
class TraversableQueue[T](sinkQueue: SinkQueue[T]) extends Traversable[T] {
@tailrec private def next[U](f: T => U): Unit = {
val nextElem = Await.result(sinkQueue.pull(), Duration.Inf)
if (nextElem.isDefined) {
f(nextElem.get)
next(f)
}
}
def foreach[U](f: T => U): Unit = next(f)
}
val sinkQueue = source.runWith(Sink.queue())
val queue = new TraversableQueue(sinkQueue)
现在我可以将可遍历队列传递给DBIO.sequence()。但是,这违背了流式处理的目的。
我发现的另一种方法是:
def toDbioAction[T](queue: SinkQueue[DBIOAction[S, NoStream, Effect.All]]):
DBIOAction[Queue[T], NoStream, Effect.All] =
DBIO.from(queue.pull() map { tOption =>
tOption match {
case Some(action) =>
action.flatMap(t => toDbioAction(queue).map(_ :+ t))
case None => DBIO.successful(Queue())
}
}).flatMap(r => r)
使用这种方法,可以不阻塞地生成一系列DBIOActions:
toDbioAction(source.runWith(Sink.queue()))
有没有更好/更惯用的方法来达到预期的效果?
【问题讨论】:
-
我认为这与您将获得的解决方案一样好。您需要在最后阻塞才能在单个事务中运行;此外,在您致电
db.run()之前,Slick 实际上不会向数据库发送任何内容。
标签: akka slick akka-stream slick-3.0