【问题标题】:Akka Streams: Copy element to another stream via alsoToAkka Streams:通过alsoTo 将元素复制到另一个流
【发布时间】:2020-08-15 06:51:35
【问题描述】:

我想使用alsoTo 将元素从一个Source 复制到另一个,但它不像我预期的那样工作。从 Java InputStream 创建 Akka Source 并进行一些转换并使用 alsoTo 创建 s1 的副本的代码示例:

import java.io.ByteArrayInputStream
import java.nio.charset.StandardCharsets

import akka.actor.ActorSystem
import akka.stream.IOResult
import akka.stream.scaladsl.{Sink, Source, StreamConverters}

import scala.concurrent.duration._
import scala.concurrent.{Await, ExecutionContext, Future}

object Main {
  def main(args: Array[String]): Unit = {
    implicit val system: ActorSystem = ActorSystem("AkkaStreams_alsoTo")
    implicit val executionContext: ExecutionContext = scala.concurrent.ExecutionContext.Implicits.global

    val byteStream = new ByteArrayInputStream("a,b,c\nd,e,f\ng,h,j".getBytes(StandardCharsets.UTF_8))
    try {
      val s1: Source[List[List[String]], Future[IOResult]] = StreamConverters.fromInputStream(() => byteStream)
        .map { bs =>
          val rows = bs.utf8String.split("\n").toList
          val valuesPerRow = rows.map(row => row.split(",").toList)
          valuesPerRow
        }
      // A copy of s1?
      val s2: Source[List[List[String]], Future[IOResult]] = s1.alsoTo(Sink.collection)
      println("s1.runForeach: ")
      Await.result(s1.runForeach(println), 20.seconds)

      println("s2.runForeach: ")
      Await.result(s2.runForeach(println), 20.seconds)

      println("Done")
      system.terminate()
    }
    finally {
      byteStream.close()
    }
  }
}

它产生以下输出:

s1.runForeach: 
List(List(a, b, c), List(d, e, f), List(g, h, j))
s2.runForeach: 
Done

如您所见,s2.runForeach 不打印任何元素。这种行为的原因是什么——是因为它读取Java InputStream时的副作用?

我正在使用 Akka Streams v2.6.8。

【问题讨论】:

    标签: scala akka akka-stream


    【解决方案1】:

    我想使用alsoTo 将元素从一个Source 复制到另一个,但它没有按预期工作。

    alsoTo 不会将元素从Source 复制到另一个Source;它有效地从Source/Flow 复制元素并将它们发送到另一个Sink(alsoTo 方法的参数)。因此,您的期望是错误的。

    如您所见,s2.runForeach 不打印任何元素。这种行为的原因是什么——是因为它读取Java InputStream时的副作用?

    因为byteStream 是val,所以s1.runForeach(println) 和s2.runForeach(println) “共享”此实例,即使它们是两个不同的 Akka Stream 蓝图。因此,当s1.runForeach(println)被调用时,byteStream被消耗掉,而当s2.runForeach(println)之后被执行时,InputStream中就没有任何东西可供s2.runForeach(println)打印了。

    将byteStream 更改为def,并打印以下内容:

    s1.runForeach: 
    List(List(a, b, c), List(d, e, f), List(g, h, j))
    s2.runForeach: 
    List(List(a, b, c), List(d, e, f), List(g, h, j))
    Done
    

    这就解释了为什么s2.runForeach(println) 在这种特殊情况下不打印任何内容,但实际上并没有真正显示alsoTo。您的设置存在缺陷,因为 s2.runForeach(println) 仅打印来自 Source 的元素,而忽略了 alsoTo(Sink.collection) 的具体化值。

    查看alsoTo 行为的一种简单方法如下:

    val byteStream = new ByteArrayInputStream("a,b,c\nd,e,f\ng,h,j".getBytes(StandardCharsets.UTF_8))
    
    val s1: Source[List[List[String]], Future[IOResult]] =
      StreamConverters.fromInputStream(() => byteStream)
        .map { bs =>
          val rows = bs.utf8String.split("\n").toList
          val valuesPerRow = rows.map(row => row.split(",").toList)
          valuesPerRow
        }
    
    val stream = s1.alsoTo(Sink.foreach(println)).runWith(Sink.foreach(println))
                                              // ^ same thing as .runForeach(println)
    Await.ready(stream, 5.seconds)
    println("Done")
    system.terminate()
    

    运行上述打印以下内容:

    List(List(a, b, c), List(d, e, f), List(g, h, j))
    List(List(a, b, c), List(d, e, f), List(g, h, j))
    Done
    

    one Source 中的元素被发送到both Sinks。

    如果你想使用Sink.collection...

    val (result1, result2) =
      s1.alsoToMat(Sink.collection)(Keep.right).toMat(Sink.collection)(Keep.both).run()
     // ^ note the use of alsoToMat in order to retain the materialized value
    
    val res1 = Await.result(result1, 5.seconds)
    val res2 = Await.result(result2, 5.seconds)
    println(s"res1: $res1")
    println(s"res2: $res2")
    println("Done")
    system.terminate()
    

    ...打印...

    res1: List(List(List(a, b, c), List(d, e, f), List(g, h, j)))
    res2: List(List(List(a, b, c), List(d, e, f), List(g, h, j)))
    Done
    

    同样,一个Source 中的元素被发送到两个Sinks。

    【讨论】:

      猜你喜欢
      • 2016-07-06
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2010-11-12
      • 1970-01-01
      • 1970-01-01
      • 2017-01-08
      • 1970-01-01
      相关资源
      最近更新 更多