【问题标题】:How to stop processing further elements in akka streams?如何停止处理 akka 流中的其他元素?
【发布时间】:2020-11-19 01:23:07
【问题描述】:

我有一个整数列表 {2,4,6,8,9,10,12}

为了简化我的问题,我的目标是获取所有偶数,直到遇到奇数。所以我的结果应该是 -> {2,4,6,8, 9}

另外,我有一个演员说一个数字是偶数还是奇数(为简单起见)

我做了以下事情:-

CompletionStage<List<Integer>> result = Source.from(integerList)
      .ask(oddEvenActor, OddEvenResponse.class, Timeout.apply(1, TimeUnit.SECONDS))
      .map(oddEvenResult -> if(oddEvenResult.isOdd()){
                                //stop processing further elements
                            }
                            else {
                                return oddEvenResult.number();
                            })
     .runWith(Sink.seq(), materializer)

那么,一旦我遇到一个奇怪的元素,我怎样才能停止对其他元素的处理呢?

一旦流完成,CompletionStage“结果”应该包含 2,4,6,8,9。

我查看了 statefulMapConcat (https://doc.akka.io/docs/akka/current/stream/operators/Source-or-Flow/statefulMapConcat.html) 但是,这仍然会处理 9 之后的其他元素,因为仍然会“询问”演员

当然我可以做到以下几点:-

  1. 有一个 resultList 变量(全局),我执行一个 resultList.add(oddEvenResult.number()),然后在遇到奇数时抛出异常。我必须编写一个自定义异常类来搭载这个全局结果列表。

  2. 按照@Jeffrey Chung 的建议使用 takeWhile,但 OddEvenActor 仍然被“要求”处理元素 10 和 12。这是没有意义的。

有没有更简洁的方法来实现这一点?

【问题讨论】:

  • 你如何“停止处理”?

标签: java scala akka akka-stream akka-http


【解决方案1】:

使用takeWhile。在 Scala 中,这将类似于以下内容:

implicit val timeout: akka.util.Timeout = 3.seconds

val result: Future[Seq[Int]] =
  Source(List(2, 4, 6, 8, 9, 10, 12))
    .ask[OddEvenResponse](oddEvenActor)
    .takeWhile(resp => !resp.isOdd, true)
    .map(_.number)
    .runWith(Sink.seq)

注意在调用takeWhile 时使用了inclusive 布尔标志,如果要保留第一个奇数,这是必需的。

Java 等价物看起来很相似。

【讨论】:

  • 这给出了所需的输出,但仍然“要求”actor 处理其他元素 10 和 12。如何根本不处理其他元素?
【解决方案2】:

如果你想保留你的演员,你可以实现它,例如如下(在 Scala 中):

implicit val system = ActorSystem("StopOnOdd")
implicit val materializer = ActorMaterializer()

class StopOnOdd extends Actor with ActorLogging {
  override def receive: Receive = {
    case x: Int if x % 2 == 0 =>
      log.info(s"Just received an even int: $x")
      sender() ! x
    case x: Int if x % 2 == 1 =>
      log.info(s"Just received an odd number: $x Stop processing.")
      context.become(dontProcess)
    case _ =>
  }

  private def dontProcess: Receive = {
    case x =>
      log.info(s"Dropping $x because odd number was received.")
    }
  }


  def main(args: Array[String]): Unit = {
  val stopOnOdd = system.actorOf(Props[StopOnOdd], "simpleActor")

  val source = Source(List(2,4,6,8,9,10,12))
  implicit val timeout: Timeout = Timeout(2.seconds)
  val stopOnOddFlow = Flow[Int].ask[Int](parallelism = 1)(stopOnOdd)

  source.via(stopOnOddFlow).to(Sink.foreach[Int](number => println(s"Got number: $number"))).run()
}

输出是:

[INFO] [07/29/2020 16:37:14.618] [StopOnOdd-akka.actor.default-dispatcher-4] 
[akka://StopOnOdd/user/simpleActor] Just received an even int: 2
Got number: 2
[INFO] [07/29/2020 16:37:14.625] [StopOnOdd-akka.actor.default-dispatcher-2] 
[akka://StopOnOdd/user/simpleActor] Just received an even int: 4
Got number: 4
Got number: 6
[INFO] [07/29/2020 16:37:14.627] [StopOnOdd-akka.actor.default-dispatcher-4] 
[akka://StopOnOdd/user/simpleActor] Just received an even int: 6
Got number: 8
[INFO] [07/29/2020 16:37:14.627] [StopOnOdd-akka.actor.default-dispatcher-4] 
[akka://StopOnOdd/user/simpleActor] Just received an even int: 8
[INFO] [07/29/2020 16:37:14.628] [StopOnOdd-akka.actor.default-dispatcher-4] 
[akka://StopOnOdd/user/simpleActor] Just received an odd number: 9 Stop processing.

【讨论】:

  • 嗨,不,我不想惹 OddEvenActor。只是如果我得到 ODD 的结果,我只想停止处理流中的其他元素。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-12-05
  • 1970-01-01
  • 2022-01-15
  • 1970-01-01
相关资源
最近更新 更多