根据问题中的内容,我认为您可能正在使用 Scala 控制台或 Scala 工作表。
在 Scala 控制台或工作集中,它打印当前语句中创建的事物的表示。例如,
scala> val i = 5
val i: Int = 5
scala> val s = "ssfdf"
val s: String = ssfdf
但是,当你在这里使用println 之类的东西时会发生什么,
scala> val u = println("dfsd")
dfsd
val u: Unit = ()
它还执行 println,然后打印出该 println 创建的值 u 实际上是一个 Unit。
这可能是您的困惑所在,因为您在Sink.foreach 中的println 在这种情况下不起作用。
这是因为这种情况更像是下面的情况,您实际上是在定义一个函数。
scala> val f1 = (s: String) => println(s)
val f1: String => Unit = $Lambda$1062/0x0000000800689840@1796b2d4
您在这里没有使用println,您只是定义了一个将使用println 的函数(String => Unit 或Function1[String, Unit] 的实例)。
所以,控制台只是打印出此处创建的值 f1 是 String => Unit 类型。
你需要调用这个函数来实际执行println,
scala> f1.apply("dsfsd")
dsfsd
同样,someSource.to(Sink.foreach(println)) 将创建一个 RunnableGraph 类型的值,因此 scala 控制台将打印类似 val res0: RunnableGraph... 的内容。
您现在需要运行此图表才能实际执行它。
但与前面的函数示例相比,graph 的执行是在线程池上异步执行的,这意味着它可能无法在某些版本的 Scala 控制台或工作台中工作(取决于线程池生命周期的管理方式)。所以,如果你这样做,
scala> val someSource = Source.single(ByteString("SomeValue"))
val someSource: akka.stream.scaladsl.Source[akka.util.ByteString,akka.NotUsed] = Source(SourceShape(single.out(369296388)))
scala> val runnableGraph = someSource.to(Sink.foreach(println))
val runnableGraph: akka.stream.scaladsl.RunnableGraph[akka.NotUsed] = RunnableGraph
scala> runnableGraph.run()
如果它有效,那么您将看到以下内容,
scala> runnableGraph.run()
val res0: akka.NotUsed = NotUsed
ByteString(83, 111, 109, 101, 86, 97, 108, 117, 101)
但您可能会看到一些与控制台由于某种原因未能完成图表运行有关的错误。
您实际上需要具体化Sink,这将在运行图表时重新生成Future[Done]。然后你将不得不使用Await 等待Future[Done]。
您必须将所有这些放入一个普通的 Scala 文件并作为 Scala 应用程序执行。
import akka.{Done, actor}
import akka.actor.typed.ActorSystem
import akka.actor.typed.scaladsl.Behaviors
import akka.stream.scaladsl.{Keep, Sink, Source}
import akka.util.ByteString
import scala.concurrent.duration.Duration
import scala.concurrent.{Await, Future}
object TestAkkaStream extends App {
val actorSystem = ActorSystem(Behaviors.empty, "test-stream-system")
implicit val classicActorSystem = actorSystem.classicSystem
val someSource = Source.single(ByteString("SomeValue"))
val runnableGraph = someSource.toMat(Sink.foreach(println))(Keep.right)
val graphRunDoneFuture: Future[Done] = runnableGraph.run()
Await.result(graphRunDoneFuture, Duration.Inf)
}