Scala 并发编程进阶:Future、ExecutionContext 与结构化并发

从 Future 的组合子语义与 ExecutionContext 线程池模型出发,剖析急切求值、无法取消、异常吞噬与阻塞传播四大缺陷,给出 Promise 手动完成与并发原语的实战写法,并对比结构化并发、虚拟线程与 Java 25 的 StructuredTaskScope,最后落到 Cats Effect 与 ZIO 的取舍,以及生产环境的超时、重试、熔断与监控方案。

Scala 的并发故事有两层:一层是标准库的 Future 与 ExecutionContext,它足够简单,也足够多坑;另一层是 Cats Effect 与 ZIO 代表的惰性可取消效果系统,它把并发从「提交任务」重构成「描述配方」。本文先把 Future 这条主线的语义、线程池模型与四大缺陷讲透,再沿着 Promise、并发原语、并行集合走到 Scala 3 结构化并发与虚拟线程,最后给出与 CE 和 ZIO 的取舍标准,以及一套生产可用的超时、重试、熔断、监控与测试方案。

前置:/scala-functional-effects/(效果系统入门)、/scala-cats-effect-deep-dive/(Cats Effect 并发模型)。

目录

1. Future 基础与组合子

scala.concurrent.Future 自 Scala 2.10 起进入标准库,表示一个未来某个时刻会完成的计算,终态只有成功与失败两种。理解 Future 的第一要义是:它在构造的那一刻就已经开始执行,而不是在你调用 map 的时候。所以 Future { ... } 返回的是句柄而不是配方,后续组合子只是在句柄上挂回调。

要点清单
- Future.apply 立即把任务提交给 ExecutionContext,返回句柄而非配方
- map / flatMap / filter / recover / recoverWith / transform 是核心组合子
- sequence 把 List[Future[A]] 翻转为 Future[List[A]],遇首个失败即整体失败
- traverse 等价于 map 加 sequence,但只遍历一次,更省一次分配
- zip 并行等待两个独立结果,比 flatMap 嵌套快一倍左右
- andThen 只做副作用,不改变结果,其回调抛出的异常会被忽略并上报
import scala.concurrent.Future
import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.*

case class User(id: Long, name: String)

def fetchUser(id: Long): Future[User] =
  Future { Thread.sleep(50); User(id, s"user-$id") }

def fetchOrders(id: Long): Future[List[Long]] =
  Future { Thread.sleep(30); List(id * 10, id * 10 + 1) }

// zip 并行等待,总耗时约等于较慢的那个
val profile: Future[(User, List[Long])] = fetchUser(1).zip(fetchOrders(1))

// traverse 一次遍历,比 list.map(fetchUser).sequence 少一次分配
val names: Future[List[String]] =
  Future.traverse(List(1L, 2L, 3L))(fetchUser(_).map(_.name))

// recover 把失败分支收敛成同类型的成功值
val safe: Future[String] = fetchUser(42).map(_.name).recover {
  case e: RuntimeException => s"fallback: ${e.getMessage}"
}

工程要点:sequence 与 traverse 在任一元素失败时立刻返回失败,但其余 Future 仍在后台跑完且无法取消;若希望「部分成功也拿到结果」,改用 Future.foldLeft 累积,或先把每个元素 recover 成 Either,把失败降级为数据。

2. ExecutionContext 与线程池模型

ExecutionContext 本质是 Executor 的包装,决定回调在哪个线程上跑。默认的 global 是并行度等于 CPU 核数的 ForkJoinPool,通过 -Dscala.concurrent.context.maxThreads、minThreads、numThreads 三个系统属性调整。它的问题不在性能,而在混用:CPU 密集任务与阻塞 IO 任务共用一个池,阻塞任务占住工作线程,会引发线程饥饿级联。

要点清单
- global 是 ForkJoinPool,默认并行度等于 CPU 核数
- blocking { } 只是向 FJP 申请临时扩容,不改变任务本身的阻塞性质
- 多 EC 隔离故障域:CPU 密集与 IO 密集必须分开,避免互相拖垮
- MDC 不会自动跨线程传播,需要包装 Executor 手动复制上下文
- 不要在 EC 线程上用 Await.result 等另一个 Future,会死锁
import java.util.concurrent.{Executors, ThreadFactory, Executor}
import java.util.concurrent.atomic.AtomicInteger
import scala.concurrent.ExecutionContext

val counter = new AtomicInteger(0)
val factory: ThreadFactory = r => {
  val t = new Thread(r, s"app-ec-${counter.incrementAndGet()}")
  t.setDaemon(true)
  t
}

// 固定线程池:适合阻塞 IO,线程数按下游连接池与超时估算
val blockingEc: ExecutionContext =
  ExecutionContext.fromExecutor(Executors.newFixedThreadPool(32, factory))

// 虚拟线程:JDK 21 起,每任务一线程,阻塞成本极低
val virtualEc: ExecutionContext =
  ExecutionContext.fromExecutor(Executors.newVirtualThreadPerTaskExecutor())

// MDC 不会跨线程传播,包装 Executor 手动复制上下文
def mdcAware(delegate: Executor): Executor = (command: Runnable) => {
  val ctx = Option(org.slf4j.MDC.getCopyOfContextMap)
  delegate.execute { () =>
    val prev = Option(org.slf4j.MDC.getCopyOfContextMap)
    ctx.foreach(org.slf4j.MDC.setContextMap)
    try command.run() finally prev.fold(org.slf4j.MDC.clear())(org.slf4j.MDC.setContextMap)
  }
}

工程要点:blocking { } 只是提示 FJP 临时扩容,扩容上限受 maxThreads 约束,极端负载下仍会排队;真正稳妥的做法是给阻塞调用一个专属的固定线程池,并在其上设置明确的超时。

3. Future 的四大缺陷

Future 的设计目标是「轻量异步句柄」,不是「并发编程模型」。把它当模型用,会在四个地方反复摔跤。认清这四点,是判断什么时候该换 Cats Effect 或 ZIO 的分水岭。

四大缺陷
- 急切求值:一构造就提交,无法复用、无法重试、无法只描述不执行
- 无法取消:没有取消协议,只能靠 AtomicBoolean 之类的协作式标志位自救
- 异常吞噬:失败只在 onComplete 与 recover 处可见,忘挂回调就静默丢失
- 阻塞传播:一个阻塞任务占住工作线程,容易引发线程饥饿与级联超时
- 衍生问题:全局隐式 EC 让依赖关系不可见,测试难以注入确定性调度
import java.util.concurrent.atomic.AtomicBoolean

// 缺陷一:急切求值。下面两行都会立即执行,没有延迟语义
val a = Future { expensiveCall() }
val b = Future { expensiveCall() }

// 缺陷二:协作式取消。Future 本身不支持取消,只能自查标志位
def longTask(flag: AtomicBoolean): Future[Unit] = Future {
  while (!flag.get()) { step() }
}

// 缺陷三:异常吞噬。这个 Future 的失败没有任何人观察
def fireAndForget(): Unit = Future { risky() }

// 缺陷四:阻塞传播。在 EC 线程上等待,可能直接耗尽线程池
def deadlockProne(): Unit = Await.result(Future { blockingCall() }, 5.seconds)

工程要点:给每个「发射后不管」的 Future 挂一个 onComplete 做日志与指标,否则失败就是黑洞;重试语义不要靠递归重建 Future 实现,那会重复提交,正确做法是把副作用包成 () => Future[A] 交给重试器调用。

4. Promise 与手动完成

Promise[A] 是 Future 的写端,与 Future[A] 构成一对:Future 是只读句柄,Promise 是一次性可写句柄。trySuccess、tryFailure、tryComplete 保证只有第一次调用生效,这天然适配「多路竞争、先到先得」的场景,比如超时、竞速、把回调式 API 桥接成 Future。

要点清单
- Promise 与 Future 一一对应,promise.future 拿到只读端
- tryComplete 系列返回 Boolean,表示本次是否真的写入成功
- 典型用途:超时、first-completed 竞速、回调式 API 桥接
- Promise 线程安全,可在任意线程完成,无需额外加锁
- 用 Promise 实现超时时务必取消定时器,避免线程与内存泄漏
import scala.concurrent.{Future, Promise}
import scala.concurrent.duration.*
import java.util.concurrent.{Executors, TimeUnit, TimeoutException}

private val scheduler = Executors.newSingleThreadScheduledExecutor()

// 给任意 Future 加超时:谁先完成谁生效
def withTimeout[A](f: Future[A], d: FiniteDuration)(using ec: ExecutionContext): Future[A] = {
  val p = Promise[A]()
  val task = scheduler.schedule(
    () => p.tryFailure(new TimeoutException(s"timed out after $d")),
    d.toMillis, TimeUnit.MILLISECONDS)
  f.onComplete { r => p.tryComplete(r); task.cancel(false) }
  p.future
}

工程要点:withTimeout 只让调用方提前拿到失败,被超时的那个 Future 仍在后台继续跑,它持有的连接、线程、内存都不会自动释放;真正要止损,必须在下游客户端上设置连接超时与读超时。

5. 并发原语与内存模型

当你需要的不是异步编排而是共享可变状态,Future 帮不上忙,得回到 JVM 并发原语。JMM 的核心是 happens-before 关系:volatile 写与读、synchronized 的解锁与加锁、AtomicXxx 的 CAS 操作,都会建立跨线程的可见性保证。缺了它,代码可能永远读到缓存里的旧值。

要点清单
- volatile 保证可见性与有序性,不保证复合操作的原子性
- AtomicInteger 与 AtomicReference 基于 CAS,适合低争用计数
- LongAdder 在高争用计数下显著优于 AtomicLong,代价是读取弱一致
- ConcurrentHashMap 的 computeIfAbsent 是原子的,函数体内不要再操作同一张表
- synchronized 在 Scala 中写作 lock.synchronized { },是可重入互斥锁
- 无锁不等于无竞争,CAS 自旋失败会浪费 CPU,高争用时锁反而更快
import java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.atomic.{AtomicLong, LongAdder}

class Metrics {
  private val requests = new LongAdder          // 高争用计数首选
  private val latency = new AtomicLong(0L)      // 需要精确读改写时用 Atomic
  private val gauges = new ConcurrentHashMap[String, Long]()

  def record(ms: Long): Unit = { requests.increment(); latency.addAndGet(ms) }

  def gauge(name: String, v: Long): Long =
    gauges.merge(name, v, (a: Long, b: Long) => a + b)
}

@volatile private var running = true
def stop(): Unit = running = false              // volatile 写,建立 happens-before
def loop(): Unit = while (running) { work() }   // volatile 读,可见性有保证

private val lock = new Object
def transfer(from: Account, to: Account, amount: Long): Unit = lock.synchronized {
  from.withdraw(amount)
  to.deposit(amount)
}

工程要点:LongAdder.sum() 在并发写入时返回近似值,只适合做监控指标;任何需要「精确读到最新值再据此决策」的地方必须用 AtomicLong 或锁,否则会出现检查与使用之间的竞态。

6. 并行集合与 Par 演进

.par 曾是 Scala 2.9 到 2.12 的招牌特性,2.13 起被拆分为独立的 scala-parallel-collections 模块,Scala 3 沿用这一安排。它基于 ForkJoin 与 Spliterator 拆分,适合「纯函数、无共享状态、计算密集」的批量处理;反过来说,只要任务里有 IO 或阻塞,它就会污染全局的公共 ForkJoinPool。

要点清单
- Scala 2.13 起必须显式依赖 scala-parallel-collections,导入 CollectionConverters
- 默认使用公共 ForkJoinPool,无法隔离,一个阻塞任务拖垮全局
- 拆分粒度由 Spliterator 决定,Vector 拆分均匀,List 拆分代价高
- 归约操作必须满足结合律,否则并行结果与串行结果不一致
- 顺序敏感的 foldLeft 在 Par 上语义会变,需用 fold 显式给结合律
- 现代替代:Future.traverse 配专用 EC,或流式库的并行算子
// build.sbt: "org.scala-lang.modules" %% "scala-parallel-collections" % "1.0.4"
import scala.collection.parallel.CollectionConverters.*

val xs = (1 to 1000000).toVector
val sum: Long = xs.par.map(_ * 2L).sum   // 结合律成立,并行安全

// 现代写法:显式控制并发上限,失败可观察,可加超时
def batchCompute(xs: Vector[Long])(using ec: ExecutionContext): Future[Long] =
  Future.traverse(xs.grouped(10000).toList)(c => Future(c.map(_ * 2L).sum)).map(_.sum)

工程要点:并行集合的收益只在单次任务耗时远大于拆分开销时才显现,几十万元素的简单映射往往被拆分与合并开销吃掉;线上代码优先用 Future.traverse 加显式并发上限,把并行度、超时与失败语义握在自己手里。

7. 结构化并发与虚拟线程

结构化并发要解决的是「父任务不结束,子任务必须被回收」这一约束。Scala 3 用 boundary 与 break 提供轻量的作用域跳转,Java 21 引入预览版、Java 25 正式转正的 StructuredTaskScope 给出标准答案;再叠加虚拟线程,阻塞式写法重新变得可接受。

要点清单
- boundary 定义作用域边界,break 从任意深度跳出并返回一个值
- 相比非局部 return 抛异常的实现,boundary 由编译器支持,无异常开销
- Java 21 的 StructuredTaskScope 为预览 API,Java 25 正式转正
- ShutdownOnFailure 遇首个失败即取消兄弟任务,ShutdownOnSuccess 用于竞速
- 虚拟线程让一任务一线程可行,sleep 不再昂贵,但不解决 CPU 密集
- 虚拟线程遇 synchronized 会钉住载体线程,需换成 ReentrantLock
import scala.util.boundary, boundary.break

// 用 boundary 实现提前返回,编译期支持,无异常开销
def firstPositive(xs: List[Int]): Int = boundary {
  for x <- xs do if x > 0 then break(x)
  -1
}

// 结构化并发:任一子任务失败,其余自动取消,且必须全部收敛
import java.util.concurrent.StructuredTaskScope
final case class Profile(user: User, orders: List[Long])

def loadProfile(id: Long): Profile = {
  val scope = new StructuredTaskScope.ShutdownOnFailure()
  try {
    val user   = scope.fork(() => fetchUser(id))
    val orders = scope.fork(() => fetchOrders(id))
    scope.join()            // 等全部子任务结束或首个失败
    scope.throwIfFailed()   // 有失败则抛出,兄弟任务已被取消
    Profile(user.get(), orders.get())
  } finally scope.close()   // 作用域退出时确保所有子任务终止
}

工程要点:StructuredTaskScope 要求 JDK 21 以上且 21 到 24 需要 --enable-preview,生产上更稳妥的过渡路径是先用虚拟线程池加显式超时,等 Java 25 普及后再切到作用域 API;boundary 解决控制流,结构化并发解决生命周期,两者不在同一层面。

8. 与 Cats Effect 及 ZIO 的对比取舍

Future 与效果系统的根本差别不在 API,而在求值模型:Future 是急切的完成句柄,IO 与 ZIO 是惰性的、可取消的、可组合的配方。前者上手五分钟,后者需要理解 Fiber、作用域与类型化错误通道。选择标准很简单:只要业务需要取消、资源安全、重试与调度统一抽象,就应该付出学习成本。

| 维度 | Future | Cats Effect IO | ZIO |
| 求值时机 | 急切,构造即执行 | 惰性,unsafeRunSync 才跑 | 惰性,由 Runtime 驱动 |
| 取消支持 | 无,只能协作式标志位 | fiber.cancel 协作式取消 | Fiber.interrupt |
| 资源管理 | 手动 try finally | Resource 自动释放 | Scope 与 acquireRelease |
| 错误模型 | Throwable 单通道 | Throwable,可叠加类型化 | 类型化 E 通道,R 为依赖 |
| 并发原语 | 仅 Promise 与锁 | Ref Deferred Queue | Ref Promise Queue Semaphore |
| 调度与重试 | 手写递归与定时器 | Temporal 统一超时重试 | Schedule 组合子 |
| 学习成本 | 低 | 中高 | 中高 |

// Cats Effect:重试、超时、并发全部由 Temporal 统一表达
import cats.effect.IO
import scala.concurrent.duration.*

def fetch(id: Long): IO[String] = IO { blockingCall() }

val program: IO[String] = for {
  a <- fetch(1).timeout(2.seconds).handleError(_ => "fallback")
  b <- fetch(2).retry(3, 200.millis)
  c <- (fetch(3), fetch(4)).parTupled.map(_._1)
} yield s"$a-$b-$c"

工程要点:迁移不要一步到位,可以先把重试、超时、并发上限这三件事从散落的手写 Future 里收敛成统一工具层,再逐步把边界处替换成 IO;CE 与 ZIO 内部同样跑在 JVM 线程池上,它们买的是语义而不是免费的吞吐。

9. 生产实践:超时重试熔断与监控

线上并发代码的成败取决于四件事:每个外部调用都有超时,每次失败都有上限与退避,每个下游都有熔断保护,每个异步链路都有指标与追踪。这些都应该收敛到一个工具层,而不是散落在业务代码里。

要点清单
- 超时必须同时设置连接超时与读超时,只用 Future 层超时无法真正止损
- 重试用指数退避加抖动,避免重试风暴;只重试幂等或可安全重放的操作
- 熔断器按下游维度隔离,半开状态放少量流量探测恢复
- 指标至少覆盖进行中任务数、队列长度、失败率与 P99 延迟
- 追踪上下文要显式跨线程传播,否则异步链路会断成两截
import scala.concurrent.{Future, Promise, ExecutionContext}
import scala.concurrent.duration.*
import java.util.concurrent.{Executors, TimeUnit}
import scala.util.Random

private val scheduler = Executors.newSingleThreadScheduledExecutor()

def delayed[A](d: FiniteDuration)(body: => Future[A])(using ec: ExecutionContext): Future[A] = {
  val p = Promise[A]()
  scheduler.schedule(() => p.completeWith(body), d.toMillis, TimeUnit.MILLISECONDS)
  p.future
}

// 指数退避加抖动,最多重试 max 次,仅对可重试异常生效
def retry[A](max: Int, base: FiniteDuration, isRetryable: Throwable => Boolean)
            (op: => Future[A])(using ec: ExecutionContext): Future[A] =
  op.recoverWith {
    case e if max > 0 && isRetryable(e) =>
      delayed(base + Random.nextLong(50).millis)(retry(max - 1, base * 2, isRetryable)(op))
  }

// 测试:确定性 EC 让回调同步执行,断言不再依赖 sleep
final class DeterministicExecutionContext extends ExecutionContext {
  def execute(r: Runnable): Unit = r.run()
  def reportFailure(t: Throwable): Unit = t.printStackTrace()
}

given ExecutionContext = new DeterministicExecutionContext

工程要点:retry 的参数必须写成 => Future[A] 而不是 Future[A],传值会把重试变成「重复等待同一个已失败的 Future」;同理,超时时间要按单次尝试计算,不要写成整个重试链路的总预算,否则重试次数越多总耗时越不可控。

10. 速查表与一句话记忆

| 问题 | 一句话答案 |
| Future 什么时候开始执行 | 构造那一刻,立即提交给 ExecutionContext |
| 为什么 Future 不能取消 | 它只是完成句柄,没有中断协议,只能靠协作式标志位 |
| global 线程池线程数是多少 | 默认等于 CPU 核数,可用 scala.concurrent.context.maxThreads 调整 |
| blocking 到底做了什么 | 提示 ForkJoinPool 临时扩容,不改变任务的阻塞性质 |
| 如何并行等待两个 Future | 用 zip,不要用 flatMap 嵌套,后者会串行化 |
| 什么时候该换 Cats Effect 或 ZIO | 需要真正的取消、资源安全与统一重试调度时 |
| 结构化并发的标准做法 | Java 25 的 StructuredTaskScope,或 CE 与 ZIO 的作用域 |
| 高争用计数用什么 | LongAdder,精确读改写用 AtomicLong |
| 虚拟线程的坑是什么 | synchronized 会钉住载体线程,CPU 密集任务无收益 |

一句话记忆:Future 是急切的完成句柄,IO 与 ZIO 是惰性的可取消配方;要取消与资源安全就换后者,要低成本并发就守住 ExecutionContext 的隔离与超时。

延伸阅读

  • /scala-functional-effects/ — 效果系统入门,理解 IO 为什么必须惰性
  • /scala-cats-effect-deep-dive/ — Fiber、作用域与取消语义的深入剖析
  • /scala-concurrency-atomics/ — 原子类、CAS 与无锁数据结构的底层细节
  • /scala-zio-program/ — ZIO 的类型化错误通道与依赖注入
  • /scala-testing-practice/ — 异步与并发代码的测试策略
  • /scala-observability-logging-tracing/ — 追踪上下文在异步链路上的传播
  • /scala-microservices-practice/ — 服务间调用的超时、重试与熔断落地
  • Java 企业级专题 — 虚拟线程与结构化并发在 Java 侧的工程实践

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「scala」更多文章

  1. Scala 云原生部署实战:容器化、健康检查与 Kubernetes 运维
  2. Scala Web 安全与鉴权实战:JWT、OAuth2 与安全加固
  3. Scala 缓存与 Redis 集成:Caffeine、Redis4cats 与缓存模式