【问题标题】:How do I create a Flow with a different input and output types for use inside of a graph?如何创建具有不同输入和输出类型的流以在图形内部使用?
【发布时间】:2023-03-31 16:09:02
【问题描述】:

我正在通过在内部构建图表来制作自定义水槽。这是我的代码的广泛简化,以证明我的问题:

def mySink: Sink[Int, Unit] = Sink() { implicit builder =>

    val entrance = builder.add(Flow[Int].buffer(500, OverflowStrategy.backpressure))
    val toString = builder.add(Flow[Int, String, Unit].map(_.toString))
    val printSink = builder.add(Sink.foreach(elem => println(elem)))

    builder.addEdge(entrance.out, toString.in)
    builder.addEdge(toString.out, printSink.in)

    entrance.in
}

我遇到的问题是,虽然创建具有相同输入/输出类型且只有一个类型参数且没有值参数的 Flow 是有效的,例如:Flow[Int](在整个文档中)它是只提供两个类型参数和零值参数是无效的。

根据reference documentation for the Flow object apply 我正在寻找的方法定义为

def apply[I, O]()(block: (Builder[Unit]) ⇒ (Inlet[I], Outlet[O])): Flow[I, O, Unit]

通过将 FlowGraph.Builder 传递给给定的创建函数来创建流。

create 函数应返回一对 Inlet 和 Outlet,它们对应于创建的 Flows 输入和输出端口。

当我试图创建一个我认为非常简单的流程时,我似乎需要处理另一个级别的图形构建器。有没有一种更简单、更简洁的方法来创建一个 Flow 来改变它的输入和输出的类型,而不需要弄乱它的内部端口?如果这是解决此问题的正确方法,那么解决方案应该是什么样的?

奖励:为什么创建一个不改变其输入类型与输出类型的 Flow 很容易?

【问题讨论】:

  • “只提供两个类型参数和零值参数是无效的” 这种流的语义是什么?您是否正在考虑更短的Flow[Int].map(_.toString)
  • 我的理解是Flow[Int].map(_.toString)是无效的,因为Flow[Int]表示从IntInt的流。但是,您的地图函数 (_.toString) 的类型是 Int => String。这是不正确的吗? (您的评论和用新鲜的眼光看文档让我怀疑我不正确)

标签: akka-stream


【解决方案1】:

如果您想同时指定流的输入和输出类型,您确实需要使用您在文档中找到的 apply 方法。但是,使用它的方式与您已经使用的方式几乎完全相同。

Flow[String, Message]() { implicit b =>
  import FlowGraph.Implicits._

  val reverseString = b.add(Flow[String].map[String] { msg => msg.reverse })
  val mapStringToMsg = b.add(Flow[String].map[Message]( x => TextMessage.Strict(x)))

  // connect the graph
  reverseString ~> mapStringToMsg

  // expose ports
  (reverseString.inlet, mapStringToMsg.outlet)
}

您不仅返回入口,还返回一个包含入口和出口的元组。现在我们可以将这个流程与特定的 Source 或 Sink 一起使用(例如在另一个构建器中,或直接与 runWith 一起使用)。

【讨论】:

  • 整个答案都很好,但我专门寻找的部分是Flow[String].map[Message]( x => TextMessage.Strict(x))。谢谢!
猜你喜欢
  • 2022-10-02
  • 1970-01-01
  • 2017-07-15
  • 2019-02-07
  • 2020-11-12
  • 2012-12-04
  • 1970-01-01
  • 2019-04-26
相关资源
最近更新 更多