You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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方法,它支持读取当前状态并返回新状态的原子操作

如果一定要用自定义类封装状态

如果坚持用自定义类,需要满足两个条件:

  1. 只创建一个类实例,所有Fiber共享它
  2. 对可变状态的操作必须是线程安全的

示例代码:

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.02 04:57:26