【问题标题】:Print Source[ByteString, NotUsed] values to console将 Source[ByteString, NotUsed] 值打印到控制台
【发布时间】:2021-08-10 06:38:38
【问题描述】:

如何在控制台中打印源的值。

val someSource = Source.single(ByteString("SomeValue"))

我想从此源打印字符串“SomeValue”。我试过了:

someSource.to(Sink.foreach(println)) //This one prints RunnableGraph object

someSource.map(each => {
    val pqr = each.decodeString(ByteString.UTF_8)
    print(pqr)
}) // THis one prints res3: soneSource.Repr[Unit]  = Source(SourceShape(Map.out(169373838)))

如何打印最初用于创建单个对象源的原始字符串。

【问题讨论】:

    标签: scala akka akka-stream


    【解决方案1】:

    根据问题中的内容,我认为您可能正在使用 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)
    }
    

    【讨论】:

    • 非常感谢@sarveshseri 的解释。我是函数式编程的新手,但我清楚地理解了这里的问题。
    猜你喜欢
    • 1970-01-01
    • 2011-01-30
    • 2016-07-20
    • 2013-03-29
    • 1970-01-01
    • 2011-12-03
    • 1970-01-01
    相关资源
    最近更新 更多