【问题标题】:fs2.Stream observeAsync does not execute a given sink asynchronouslyfs2.Stream observeAsync 不会异步执行给定的接收器
【发布时间】:2018-11-14 09:52:56
【问题描述】:

我正在试验fs2.Stream 并发功能,但对它的工作原理有一些误解。我想通过一些接收器并行发送流内容。这是我尝试过的:

object TestParallelStream extends App {
  val secondsOnStart = TimeUnit.MILLISECONDS.toSeconds(System.currentTimeMillis())
  val stream = fs2.Stream.emits(List(1, 2, 3, 4, 5, 6, 7, 8, 9)).covary[IO]
  val sink: fs2.Sink[IO, Int] = _.evalMap(i => IO {
    println(s"[${TimeUnit.MILLISECONDS.toSeconds(System.currentTimeMillis()) - secondsOnStart} second]: $i")
    Thread.sleep(5000)
  })
  val executor = Executors.newFixedThreadPool(4)
  implicit val cs: ContextShift[IO] = IO.contextShift(ExecutionContext.fromExecutor(executor))


  stream.observeAsync(3)(sink).compile.drain.unsafeRunSync() //1
  executor.shutdown()
}

//1 打印以下内容:

[1 second]: 1
[6 second]: 2
[11 second]: 3
[16 second]: 4
[21 second]: 5
[26 second]: 6
[31 second]: 7
[36 second]: 8
[41 second]: 9

从输出可以看出,每个元素都是通过sink顺序发送的。

但是如果我修改接收器如下:

// 5 limit and parEvalMap
val sink: fs2.Sink[IO, Int] = _.parEvalMap(5)(i => IO { 
  println(s"[${TimeUnit.MILLISECONDS.toSeconds(System.currentTimeMillis()) - secondsOnStart} second]: $i")
  Thread.sleep(5000)
})

输出是:

[1 second]: 3
[1 second]: 2
[1 second]: 4
[1 second]: 1
[6 second]: 5
[6 second]: 6
[6 second]: 7
[6 second]: 8
[11 second]: 9

现在我们有 4 个元素一次通过接收器并行发送(尽管将 3 设置为 observerAsync 的限制)。

即使我只用observe 替换observerAsync,我也得到了相同的并行化效果。

您能否解释一下水槽的实际工作原理?

【问题讨论】:

    标签: scala functional-programming scala-cats fs2


    【解决方案1】:

    observe 用于通过 多个 接收器传递流元素。它不会改变接收器本身的并发行为。

    你会这样使用它:

    stream.observeAsync(n)(sink1).to(sink2)
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2020-08-09
      • 1970-01-01
      • 2021-07-25
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-10-19
      相关资源
      最近更新 更多