【问题标题】:Put elements in stream and return an object将元素放入流中并返回一个对象
【发布时间】:2020-02-06 22:01:37
【问题描述】:

在 akka 中,我想将元素放入流中并返回一个对象。我知道这些元素可能是运行图表的来源。但是我怎样才能在运行时放置元素并返回一个对象呢?


import akka.actor.ActorSystem
import akka.stream.QueueOfferResult.{Dropped, Enqueued, Failure, QueueClosed}
import akka.stream.{ActorMaterializer, OverflowStrategy}
import akka.stream.scaladsl.{Keep, Sink, Source}

import scala.Array.range
import scala.util.Success

object StreamElement {

  implicit val system = ActorSystem("StreamElement")
  implicit val materializer = ActorMaterializer()
  implicit val executionContext = system.dispatcher

  def main(args: Array[String]): Unit = {
    val (queue, value) = Source
    .queue[Int](10, OverflowStrategy.backpressure)
    .map(x => {
      x * x
    })
    .toMat(Sink.asPublisher(false))(Keep.both)
    .run()

    range(0, 10)
      .map(x => {
        queue.offer(x).onComplete {
          case Success(Enqueued) => {
          }

          case Success(Dropped) => {}
          case _ => {
             println("others")
          }
        }
      })
    }
}

我怎样才能得到返回的值?

【问题讨论】:

  • 我很好奇你为什么要使用 Scala 集合而不是 Source 来提供 Stream 元素。特别是,由于您已经编写了一个 Stream,其中包含要在 Source Queue 和发布者 Sink 中捕获的物化值,我认为一个很好的用例是将 Scala 集合包装在 Source 中并创建一个订阅者 Source 来收集想要的值。
  • 谢谢。你的意思是“.runWith(Sink.actorRefWithAck)”?
  • 其实我对 Akka-stream 很陌生。对于我的情况,请您也举个例子吗?谢谢。
  • 很难将示例代码显示为 cmets。请在下面查看我的答案。

标签: scala akka akka-stream


【解决方案1】:

实际上,您希望返回每个元素的 int 值。 因此,您可以创建流,然后每次都连接到源和接收器。


package tech.parasol.scala.akka

import akka.actor.ActorSystem
import akka.stream.QueueOfferResult.{Dropped, Enqueued, Failure, QueueClosed}
import akka.stream.{ActorMaterializer, OverflowStrategy}
import akka.stream.scaladsl.{Flow, Keep, Sink, Source}

import scala.Array.range
import scala.util.Success

object StreamElement {

  implicit val system = ActorSystem("StreamElement")
  implicit val materializer = ActorMaterializer()
  implicit val executionContext = system.dispatcher

  val flow = Flow[Int]
    .buffer(16, OverflowStrategy.backpressure)
    .map(x => x * x)

  def main(args: Array[String]): Unit = {
    range(0, 10)
      .map(x => {
        Source.single(x).via(flow).runWith(Sink.head)
      }.map( v => println("v ===> " + v)
      ))
  }

}


【讨论】:

  • 创建源和接收器是否存在性能问题?
  • 我之前测试过,性能还可以。你也可以测试一下。
【解决方案2】:

我不清楚为什么在您的示例代码中没有将 Scala 集合作为 Source 馈送到 Stream。假设您已经组合了一个 Stream,其中包含要在 Source Queue 和发布者 Sink 中捕获的物化值,您可以使用 Source.fromPublisher 创建订阅者 Source 来收集想要的值,如下所示:

import akka.actor.ActorSystem
import akka.stream.scaladsl._
import akka.stream._

implicit val system = ActorSystem("system")
implicit val materializer = ActorMaterializer()  // Not needed for Akka 2.6+

val (queue, pub) = Source
  .queue[Int](10, OverflowStrategy.backpressure)
  .map(x =>  x * x)
  .toMat(Sink.asPublisher(false))(Keep.both)
  .run()

val fromQueue = Source(0 until 10).runForeach(queue.offer(_))

val source = Source.fromPublisher(pub)

source.runForeach(x => print(x + " "))
// Output:
// 0 1 4 9 16 25 36 49 64 81 

【讨论】:

  • 如果是请求->响应方式,例如http请求.via(flow)则为http响应。以你的例子,怎么能做到这一点?
  • @YouXiang-Wang,建议的方法只是利用 OP 已经拥有的 SourceQueue/publisher(尽管我可能误解了 OP 想要什么)。我不会在您的用例中采用这种方法。
  • 基本上,对于http请求->响应场景,我们必须使用akka ask模式。那么我们是否有另一种没有询问模式的方法来处理这种 http 请求 -> 响应场景?我也在寻找这种方法。 :-)
  • 我不清楚您的具体业务逻辑要求是什么。也许考虑发布一个问题以及示例请求/响应/流程要求。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多