Cats Effect:parUnorderedTraverse后共享mutable.Map仍为空的问题解析
并行Fiber中累积可变Map为空的原因及解决方法
我尝试从多个并行IO Fiber中将条目累积到单个共享容器中,该容器封装了一个mutable.Map。通过parUnorderedTraverse执行100次插入操作后,预期map包含100个条目,但实际得到的是空map。
版本信息: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
已尝试操作:将traverse包装在.start中并通过.join等待,结果相同。
我了解到可以使用Ref,但在使用它之前,我想理解为什么基于class的方法会失败。在Cats Effect中,跨并行fiber共享单个可变容器的正确方式是什么?
核心原因:实例未真正共享
你定义的rr0io = IO(new RR(mutable.Map.empty))是一个IO动作,每次执行这个动作(比如processItem里的rr0io.map,或者最后rr1 <- rr0io),都会新建一个独立的RR对象:
- 并行执行的100个
processItem,每个都在操作自己的RR实例,根本没有共享同一个容器 - 最后打印时调用
rr0io,又新建了一个空的RR,所以输出必然是空Map
另外,就算解决了实例共享问题,直接在多个Fiber中操作mutable.Map还会有线程安全问题:mutable.Map不是线程安全的,并行修改会导致数据竞争,可能出现丢失更新、异常甚至数据损坏。
Cats Effect中跨Fiber共享可变状态的正确方式
在Cats Effect中,必须使用线程安全的、支持原子操作的状态容器,最常用的就是Ref。
修正后的代码示例
import cats.effect.{ExitCode, IO, IOApp, Ref} import cats.syntax.parallel._ import scala.collection.mutable object Rep0Fixed extends IOApp { val app: IO[ExitCode] = for { // 用Ref封装可变Map,保证线程安全的原子操作 ref <- Ref.of(mutable.Map.empty[String, Int]) _ <- (1 to 100).toList.parUnorderedTraverse { id => ref.update(_ += (s"l0-$id" -> id)) } finalMap <- ref.get _ <- IO.println(s"res: $finalMap, size=${finalMap.size}") } yield ExitCode.Success override def run(args: List[String]): IO[ExitCode] = app }
关键说明
Ref是Cats Effect提供的线程安全引用,所有操作(update、get等)都是原子性的,彻底避免数据竞争- 只创建一次
Ref实例,所有并行Fiber共享同一个引用,确保状态是累积的 - 如果需要更复杂的状态操作,可以用
modify方法,它支持读取当前状态并返回新状态的原子操作
如果一定要用自定义类封装状态
如果坚持用自定义类,需要满足两个条件:
- 只创建一个类实例,所有Fiber共享它
- 对可变状态的操作必须是线程安全的
示例代码:
import cats.effect.{ExitCode, IO, IOApp, Ref} import cats.syntax.parallel._ import scala.collection.mutable class RR(ref: Ref[mutable.Map[String, Int]]) { def addItem(name: String, newInt: Int): IO[Unit] = ref.update(_ += (name -> newInt)) } object Rep0WithClass extends IOApp { val app: IO[ExitCode] = for { mapRef <- Ref.of(mutable.Map.empty[String, Int]) rr <- IO(new RR(mapRef)) // 只创建一个RR实例 _ <- (1 to 100).toList.parUnorderedTraverse { id => rr.addItem(s"l0-$id", id) } finalMap <- mapRef.get _ <- IO.println(s"res: $finalMap, size=${finalMap.size}") } yield ExitCode.Success override def run(args: List[String]): IO[ExitCode] = app }
内容的提问来源于stack exchange,提问作者Dmitry Petrushchenko
相关产品推荐
相关产品推荐

