【问题标题】:Conditionally skip flow using akka streams使用 akka 流有条件地跳过流
【发布时间】:2016-02-22 09:12:17
【问题描述】:

我正在使用 akka 流,并且我需要有条件地跳过图表的一部分,因为流无法处理某些值。具体来说,我有一个接受字符串并发出 http 请求的流程,但是当字符串为空时服务器无法处理这种情况。但我只需要返回一个空字符串。有没有一种方法可以做到这一点,而不必通过 http 请求就知道它会失败?我基本上有这个:

val source = Source("1", "2", "", "3", "4")
val httpRequest: Flow[String, HttpRequest, _]
val httpResponse: Flow[HttpResponse, String, _]
val flow = source.via(httpRequest).via(httpResponse)

我唯一能想到的就是在我的 httpResponse 流中捕获 400 错误并返回一个默认值。但我希望能够避免因事先知道会失败的请求而访问服务器的开销。

【问题讨论】:

  • 您的示例无法编译。 httpRequest 的输出是 HttpRequest 类型,httpResponse 的输入是 HttpResponse 类型,因此它们不能与“via”链接在一起。

标签: akka akka-stream akka-http


【解决方案1】:

你可以使用flatMapConcat:

(警告:从未编译,但你会明白它的要点)

val source = Source("1", "2", "", "3", "4")
val httpRequest: Flow[String, HttpRequest, _]
val httpResponse: Flow[HttpResponse, String, _]
val makeHttpCall: Flow[HttpRequest, HttpResponse, _]
val someHttpTransformation = httpRequest via makeHttpCall via httpResponse
val emptyStringSource = Source.single("")
val cleanerSource = source.flatMapConcat({
  case "" => emptyStringSource
  case other => Source.single(other) via someHttpTransformation
})

【讨论】:

  • Viktor,可以为每个“过滤器”创建Source.single,以任何方式对流产生负面影响吗?
  • 请原谅我的题外话问题,但map({ case _ => ??? }) 更喜欢map{ case _ => ??? } 吗?除非必要(在较长的参数列表中),否则我倾向于不将大括号嵌套在括号内。我的前同事曾经这样做({ case _ => ???}),我曾经删除多余的括号...
  • 当我使用明确的点时,我倾向于使用括号。
【解决方案2】:

Viktor Klang 的解决方案简洁明了,elegant。我只是想演示一个使用 Graphs 的替代方案。

您可以将字符串源拆分为两个流,并为一个流过滤有效字符串,另一个流过滤无效字符串。然后合并结果(“cross the streams”)。

基于documentation

val g = RunnableGraph.fromGraph(FlowGraph.create() { implicit builder: FlowGraph.Builder[Unit] =>
  import FlowGraph.Implicits._

  val source = Source(List("1", "2", "", "3", "4"))
  val sink : Sink[String,_] = ???

  val bcast = builder.add(Broadcast[String](2))
  val merge = builder.add(Merge[String](2))

  val validReq =   Flow[String].filter(_.size > 0)
  val invalidReq = Flow[String].filter(_.size == 0)

  val httpRequest: Flow[String, HttpRequest, _] = ???
  val makeHttpCall: Flow[HttpRequest, HttpResponse, _] = ???
  val httpResponse: Flow[HttpResponse, String, _] = ???
  val someHttpTransformation = httpRequest via makeHttpCall via httpResponse

  source ~> bcast ~> validReq ~> someHttpTransformation ~> merge ~> sink
            bcast ~>      invalidReq                    ~> merge
  ClosedShape
})

注意:此解决方案会拆分流,因此接收器可能会以与基于输入的预期不同的顺序处理字符串值结果。

【讨论】:

  • 是的!当您可以重新排序操作时,这是另一种选择(因为 http-flow 和 non-http 并行操作)
  • 这是完美的。我没有考虑并行过滤器。关于订购的好点,但是akka-http的请求流无论如何都是无序的(至少https池是)
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2017-10-23
  • 1970-01-01
  • 2021-06-03
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多