Cats Effect:在parUnorderedTraverse之后,共享的 `mutable.Map` 仍然为空
我正在尝试从许多并行 IO 纤程中,将项累积到一个单一的共享容器中。我的容器包装了一个 mutable.Map。通过 parUnorderedTraverse 执行100次插入后,我希望映射包含100条记录;但我得到的是一个空映射。
版本:Scala 2.13.14,cats-effect 3.5.4。
最小示例:
import cats.effect.{ExitCode, IO, IOApp}
import cats.syntax.parallel._
import scala.collection.mutable
class RR(var m: mutable.Map[String, Int]) {
def addItem(name: String, newInt: Int): IO[Unit] =
IO { m = m ++ mutable.Map(name -> newInt) }
}
object Rep0 extends IOApp {
val rr0io = IO(new RR(mutable.Map.empty))
def processItem(id: Int) = rr0io.map(_.addItem(s"l0-$id", id))
val app: IO[ExitCode] = for {
_ <- (1 to 100).toList.parUnorderedTraverse(processItem)
rr1 <- rr0io
_ <- IO.println(s"res: ${rr1.m}, size=${rr1.m.size}")
} yield ExitCode.Success
override def run(args: List[String]): IO[ExitCode] = app
}
预期: size=100。 实际: res: HashMap(), size=0。
我尝试过的做法: 将遍历包裹在 .start 中,并通过 .join 进行等待 — 结果相同。
我看到有人建议使用 Ref,但我想在使用它之前理解为什么基于 class 的方法会失败。跨并行纤程在Cats Effect中共享一个可变容器的正确方式是什么?
解决方案
这个问题并非并发问题,而是 IO(new RR(mutable.Map.empty)) 每次运行时都会构造一个全新的 RR。每次 processItem 调用都会执行 rr0io,获取一个全新的 RR,对其进行修改,然后丢弃。最终的 rr1 <- rr0io 还会再创建一个空的。
有两种修复方法。
1.共享一个实例(最小改动):
import cats.effect.{ExitCode, IO, IOApp}
import cats.syntax.parallel._
import scala.collection.mutable
class RR(var m: mutable.Map[String, Int]) {
def addItem(name: String, newInt: Int): IO[Unit] =
IO { m = m ++ mutable.Map(name -> newInt) }
}
object Rep0 extends IOApp {
def run(args: List[String]): IO[ExitCode] = for {
rr <- IO(new RR(mutable.Map.empty)) // built ONCE, shared
_ <- (1 to 100).toList.parUnorderedTraverse { id =>
rr.addItem(s"l0-$id", id)
}
_ <- IO.println(s"res: ${rr.m},\n size=${rr.m.size}")
} yield ExitCode.Success
}
你可以在 这里的Scastie上试试这个示例。
注意:在并行访问下,这里仍然存在对 var m 的数据竞争;你偶尔会看到少于100条记录。用于快速演示可以,但不适合生产环境。
2. Ref(地道且无竞争条件):
import cats.effect.{ExitCode, IO, IOApp, Ref}
import cats.syntax.parallel._
object Rep1 extends IOApp {
def run(args: List[String]): IO[ExitCode] = for {
ref <- Ref.of[IO, Map[String, Int]](Map.empty)
_ <- (1 to 100).toList.parUnorderedTraverse { id =>
ref.update(_ + (s"l0-$id" -> id))
}
m <- ref.get
_ <- IO.println(s"res: $m,\n size=${m.size}")
} yield ExitCode.Success
}
你也可以在 这里的Scastie上试试这个。它总是打印 size=100。
你的纤程尝试也因为同样的根本原因而失败:如果每个纤程都在修改不同的 RR,调度顺序就不再重要。
站内所有文章版权归属LeftHeroAI导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。