【问题标题】:Transforming the inner elements of a stream of collections转换集合流的内部元素
【发布时间】:2020-11-17 21:22:58
【问题描述】:

我最近在空闲时间学习使用 Akka Streams(包括 Scala 和 java),并且想知道如何实现以下场景。

我有一个连续的非常大的集合流进入我的管道,我想让管道转换每个集合中的元素。

将 Collection 转换为其元素流很容易,但我还需要将 1 个 Collection 中的所有转换后的元素重新收集到 1 个新 Collection 中(仅包含以前也在原始集合中的转换后的对象) .因此,我必须知道 1 Collection 的特定元素流何时被处理,因为这样我就可以发出转换后的集合,以便在通用管道中进一步处理。

【问题讨论】:

  • 你可以使用折叠。您可以包含一些代码来处理它
  • 我知道 fold 能够从元素流创建一个新集合,但是由于我有一个集合流,我不断地转换为每个集合的元素流不确定何时停止“折叠”以防止下一个集合中的元素也被折叠到当前集合中。我会尝试在我的帖子中添加一些虚拟代码。
  • 我不确定我是否在关注。如果transform 将单个CustomObject 转换为单个TransformedCustomObjecttransformationPipeline 应该如何从单个CustomObject 多个TransformedCustomObject s 创建?要么是 1:1 映射,要么是一对多。

标签: java scala akka-stream


【解决方案1】:

根据评论者的建议,您可以在 transformationPipeline 中使用 fold 来组装 List 类型的元素。要在运行 Stream 时维护 List 边界,请使用 flatMapConcat 而不是 mapConcat,如下面的简单示例所示:

def transform(s: String): Int = s.length

val transformationPipeline: Flow[String, List[Int], NotUsed] = Flow[String].
  fold(List.empty[Int])((ls, s) => transform(s) :: ls).
  map(_.reverse)

val flow: Flow[List[String], List[Int], NotUsed] = Flow[List[String]].
  flatMapConcat(Source(_).via(transformationPipeline))

Source(List("a", "bb") :: List("cc", "ddd", "e") :: Nil).
  via(flow).
  runForeach(println)
// List(1, 2)
// List(2, 3, 1)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-11-28
    • 2020-06-30
    • 2019-04-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多