【问题标题】:Is it possible to extract the substream key in akkastreams?是否可以在akkastreams中提取子流密钥?
【发布时间】:2020-11-16 19:31:45
【问题描述】:

我似乎找不到任何关于此的文档,但我知道 AkkaStreams 在内存中调用 groupBy 时将用于将流分组为子流的键存储。是否可以从子流中提取这些密钥?假设我从我的主流中创建了一堆子流,将它们通过一个折叠来计算每个子流中的对象,然后将计数存储在一个类中。我可以让子流的密钥也传递给该类吗?或者有没有更好的方法来做到这一点?我需要计算每个子流的每个元素,但我还需要存储计数所属的组。

【问题讨论】:

    标签: scala akka akka-stream


    【解决方案1】:

    stream-cookbook 中有一个很好的例子:

    val counts: Source[(String, Int), NotUsed] = words
    // split the words into separate streams first
      .groupBy(MaximumDistinctWords, identity)
      //transform each element to pair with number of words in it
      .map(_ -> 1)
      // add counting logic to the streams
      .reduce((l, r) => (l._1, l._2 + r._2))
      // get a stream of word counts
      .mergeSubstreams
    

    然后是:

    val words = Source(List("Hello", "world", "let's", "say", "again", "Hello", "world"))
    counts.runWith(Sink.foreach(println))
    

    将打印:

    (world,2)
    (Hello,2)
    (let's,1)
    (again,1)
    (say,1)
    

    我想到的另一个例子是用余数计算数字。所以以下,例如:

    Source(0 to 101)
      .groupBy(10, x => x % 4)
      .map(e => e % 4 -> 1)
      .reduce((l, r) => (l._1, l._2 + r._2))
      .mergeSubstreams.to(Sink.foreach(println)).run()
    

    将打印:

    (0,26)
    (1,26)
    (2,25)
    (3,25)
    

    【讨论】:

    • 这个策略可以满足我的需要!但是,我需要两个异步执行此操作的管道。当我实现这一点时,出于某种原因,在将管道合并在一起之后,当应该有数千个项目时,我只会得到 3 个项目。如果我用折叠来做,所有条目都会被输出,但其中一半会错过键值。知道为什么会这样吗?
    • @WolfDeWulf,没有看到你的代码真的很难知道。但是,每个 SO 问题应该只有一个问题。如果我的回答回答了您的问题,您可以接受它,并针对您的非工作代码创建另一个问题。
    猜你喜欢
    • 2019-07-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-08-29
    • 2011-10-26
    • 1970-01-01
    • 2011-01-24
    相关资源
    最近更新 更多