【发布时间】: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