Cats Effect:在parUnorderedTraverse之后,共享的 `mutable.Map` 仍然为空

编程语言 2026-07-09

我正在尝试从许多并行 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导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。

相关文章