【问题标题】:How to assemble an Akka Streams sink from multiple file writes?如何从多个文件写入中组装 Akka Streams 接收器?
【发布时间】:2016-09-05 01:29:55
【问题描述】:

我正在尝试将基于 akka 流的流集成到我的 Play 2.5 应用程序中。这个想法是您可以流式传输照片,然后将其作为原始文件、缩略图版本和水印版本写入磁盘。

我设法使用类似这样的图表来完成这项工作:

val byteAccumulator = Flow[ByteString].fold(new ByteStringBuilder())((builder, b) => {builder ++= b.toArray})
                                    .map(_.result().toArray)

def toByteArray = Flow[ByteString].map(b => b.toArray)

val graph = Flow.fromGraph(GraphDSL.create() {implicit builder =>
  import GraphDSL.Implicits._
  val streamFan = builder.add(Broadcast[ByteString](3))
  val byteArrayFan = builder.add(Broadcast[Array[Byte]](2))
  val output = builder.add(Flow[ByteString].map(x => Success(Done)))

  val rawFileSink = FileIO.toFile(file)
  val thumbnailFileSink = FileIO.toFile(getFile(path, Thumbnail))
  val watermarkedFileSink = FileIO.toFile(getFile(path, Watermarked))

  streamFan.out(0) ~> rawFileSink
  streamFan.out(1) ~> byteAccumulator ~> byteArrayFan.in
  streamFan.out(2) ~> output.in

  byteArrayFan.out(0) ~> slowThumbnailProcessing ~> thumbnailFileSink
  byteArrayFan.out(1) ~> slowWatermarkProcessing ~> watermarkedFileSink

  FlowShape(streamFan.in, output.out)
})

graph

}

然后我使用这样的累加器将它连接到我的播放控制器:

val sink = Sink.head[Try[Done]]

val photoStorageParser = BodyParser { req =>
     Accumulator(sink).through(graph).map(Right.apply)
}

问题是我的两个已处理文件接收器没有完成,并且两个已处理文件的大小都为零,但不是原始文件。我的理论是累加器只等待我的扇出的一个输出,所以当输入流完成并且我的 byteAccumulator 吐出完整的文件时,到处理完成时,播放已经从输出中获得了物化值.

所以,我的问题是:
就我的方法而言,我是否走在正确的轨道上? 运行这样的图表的预期行为是什么? 如何将所有水槽组合在一起形成一个最终水槽?

【问题讨论】:

  • 我也认为是处理后没有合并流的原因。你试过Sink.combinedoc.akka.io/docs/akka/2.4.4/scala/stream/…)吗?
  • 是的,我尝试了 Sink.combine,但它统一了多个接收器以像扇出一样发送。我想我正在寻找一个粉丝,但似乎你不能只使用接收器的来源!
  • 这似乎是一个类似的例子:doc.akka.io/docs/akka/2.4.4/scala/stream/…。也许您必须返回 SinkShape 而不是 FlowShape 来声明您的流已完成?
  • 我有一个用例,我下载.gz 文件然后提取它们。当然,每次下载都应该转到一个单独的文件,以下载的文件名命名。有什么建议可以写到Sink,其名称取决于输入,然后提取文件?像往常一样,Akka Streams 文档没有提供 combine 的示例。

标签: scala akka akka-stream playframework-2.5


【解决方案1】:

好的,经过一点帮助(安德烈亚斯在正确的轨道上),我已经找到了这个解决方案:

val rawFileSink = FileIO.toFile(file)
val thumbnailFileSink = FileIO.toFile(getFile(path, Thumbnail))
val watermarkedFileSink = FileIO.toFile(getFile(path, Watermarked))

val graph = Sink.fromGraph(GraphDSL.create(rawFileSink, thumbnailFileSink, watermarkedFileSink)((_, _, _)) {
  implicit builder => (rawSink, thumbSink, waterSink) => {
    val streamFan = builder.add(Broadcast[ByteString](2))
    val byteArrayFan = builder.add(Broadcast[Array[Byte]](2))

    streamFan.out(0) ~> rawSink
    streamFan.out(1) ~> byteAccumulator ~> byteArrayFan.in

    byteArrayFan.out(0) ~> processorFlow(Thumbnail) ~> thumbSink
    byteArrayFan.out(1) ~> processorFlow(Watermarked) ~> waterSink

    SinkShape(streamFan.in)
  }
})

graph.mapMaterializedValue[Future[Try[Done]]](fs => Future.sequence(Seq(fs._1, fs._2, fs._3)).map(f => Success(Done)))

之后,从 Play 中调用它非常容易:

val photoStorageParser = BodyParser { req =>
  Accumulator(theSink).map(Right.apply)
}

def createImage(path: String) = Action(photoStorageParser) { req =>
  Created
}

【讨论】:

  • 谢谢,我刚刚有一个类似的任务,不知道如何等待所有物化的期货。您的解决方案帮助很大,而且很有效!
  • 嗨!组合的可变数量的水槽怎么样?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-08-14
  • 1970-01-01
相关资源
最近更新 更多