【问题标题】:Improving performance of fs2 stream involving file transformation提高涉及文件转换的 fs2 流的性能
【发布时间】:2020-10-29 00:40:06
【问题描述】:

我有这样的东西(这是https://github.com/typelevel/fs2 的一个例子,我的补充,我用 cmets 标记):

import cats.effect.{Blocker, ExitCode, IO, IOApp, Resource}
import fs2.{io, text, Stream}
import java.nio.file.Paths

object Converter extends IOApp {

  val converter: Stream[IO, Unit] = Stream.resource(Blocker[IO]).flatMap  { blocker =>
    def fahrenheitToCelsius(f: Double): Double =
      (f - 32.0) * (5.0/9.0)

    io.file.readAll[IO](Paths.get("testdata/fahrenheit.txt"), blocker, 4096)
      .balanceAvailable // my addition
      .map ( worker => // my addition
        worker // my addition
          .through(text.utf8Decode)
          .through(text.lines)
          .filter(s => !s.trim.isEmpty && !s.startsWith("//"))
          .map(line => fahrenheitToCelsius(line.toDouble).toString)
          .intersperse("\n")
          .through(text.utf8Encode)
          .through(io.file.writeAll(Paths.get("testdata/celsius.txt"), blocker))
      ) // my addition
      .take(4).parJoinUnbounded // my addition
  }

  def run(args: List[String]): IO[ExitCode] =
    converter.compile.drain.as(ExitCode.Success)
}

如果fahrenheit.txt 和例如一样大。 300mb 原始代码的执行需要几分钟。看来我的代码并没有更快。我怎样才能提高它的性能?运行时有大量的未使用 CPU 电源,磁盘是SSD,所以我不知道为什么它这么慢。我不确定我是否正确使用了balance

【问题讨论】:

  • 如果你有很多线程写入同一个文件,你不会赢太多,特别是因为这个工作不是真正的 cpu 绑定。
  • 我已经在我的硬件上检查了dd if=fahrenheit.txt of=fahrenheit2.txt 需要 5 秒(所以它只是一个没有任何转换的线程,纯 IO 操作)。虽然我同意提高性能并不像增加线程那么容易,但我认为可以提供缓冲区或其他可以减少所需时间的机制的组合。
  • 也许你可以增加读取缓冲区大小。反正 fs2 可能没有 dd 快。

标签: scala fs2 cats-effect


【解决方案1】:

罪魁祸首是text.utf8Encode,它不必要地每行发出一个块。当有很多短线时,这是一个巨大的开销,就像在示例中一样(每行一个温度值,108199750 行)。最近解决了(拉取请求:https://github.com/typelevel/fs2/pull/2096)。下面我提供了一个基于此 PR 的内联解决方案(只要有人使用没有此修复的版本就很有用):

import cats.effect.{Blocker, ExitCode, IO, IOApp, Resource}
import fs2.{io, text, Stream, Pipe, Chunk}
import java.nio.file.Paths
import java.nio.charset.Charset

object Converter extends IOApp {

  val converter: Stream[IO, Unit] = Stream.resource(Blocker[IO]).flatMap  { blocker =>
    def fahrenheitToCelsius(f: Double): Double =
      (f - 32.0) * (5.0/9.0)

    def betterUtf8Encode[F[_]]: Pipe[F, String, Byte] =
      _.mapChunks(c => c.flatMap(s => Chunk.bytes(s.getBytes(Charset.forName("UTF-8")))))

    io.file.readAll[IO](Paths.get("testdata/fahrenheit.txt"), blocker, 4096)
      .through(text.utf8Decode)
      .through(text.lines)
      .filter(s => !s.trim.isEmpty && !s.startsWith("//"))
      .map(line => fahrenheitToCelsius(line.toDouble).toString)
      .intersperse("\n")
      // .through(text.utf8Encode) // didn't finish, could be an hour
      .through(betterUtf8Encode) // 2 minutes
      .through(io.file.writeAll(Paths.get("testdata/celsius.txt"), blocker))
  }

  def run(args: List[String]): IO[ExitCode] =
    converter.compile.drain.as(ExitCode.Success)
}

这是 2 分钟和可能一个小时或更长时间之间的差异,在这种情况下...

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-02-04
    • 2013-01-22
    • 2022-08-04
    • 2021-06-01
    • 1970-01-01
    • 1970-01-01
    • 2021-04-30
    相关资源
    最近更新 更多