纯函数式效果系统实战:Cats Effect IO 与 ZIO

深入 Scala 效果系统:Cats Effect IO 与 ZIO 的 IO Monad、异步与并发原语、资源安全(bracket/Finalizer)、错误处理、依赖注入(ZLayer/Runtime),用两套主流实现落地生产级纯函数式程序。

引言

「纯函数式程序」听上去很学术,落到生产只有一个问题:副作用(网络、数据库、时间)怎么办? 效果系统(Effect System)的答案是把副作用包装成值——IO[User] 不是「拿到一个 User」,而是「一份能拿到 User 的程序描述」。你在纯代码里组合这份描述,最后才在唯一的边界执行它。Cats Effect 与 ZIO 是 JVM 上两套主流实现,本文讲透核心模型并给出可直接运行的写法。

前置:/scala-functional-programming/(Monad/类型类)。配合 Web 见 /scala-web-http-apps/(http4s 即基于 Cats Effect)。


目录


1. 为什么需要效果系统

普通代码的问题:副作用与纯逻辑纠缠,无法测试、无法组合、无法确定失败行为。

// 问题代码:读文件 + 解析 + 写日志,全在一块不可测
def loadConfig(path: String): Config = {
  val raw = scala.io.Source.fromFile(path).mkString   // 副作用
  val cfg = parse(raw)                                 // 纯
  println(s"loaded: $cfg")                            // 副作用
  cfg
}

效果系统改造:只有 run 时刻才发生副作用,中间全是可测的纯值:

import cats.effect.IO

def loadConfig(path: String): IO[Config] =
  IO.blocking(Source.fromFile(path).mkString)   // 描述「读文件」
    .map(parse)                                  // 纯变换
    .flatMap(cfg => IO(println(s"loaded: $cfg")).as(cfg)) // 组合描述

// 测试时:替换成内存字符串即可
def loadConfigTest: IO[Config] = IO.pure("key=value").map(parse)

核心心智:IO[A] 是「程序清单」,不是「值」。写完不执行,执行有且仅有 unsafeRunSync 一处。


2. IO Monad:副作用包装成值

创建 IO:

import cats.effect.IO

val now: IO[Long] = IO(System.currentTimeMillis())   // 延迟执行
val pure: IO[Int] = IO.pure(42)                       // 纯值立即封装
val deferred: IO[Int] = IO.defer(IO(1 + 1))          // 惰性
val blocking: IO[Unit] = IO.blocking(Thread.sleep(1000)) // 阻塞区(用专用线程池)

// 立即执行(仅测试/入口使用)
val result: Int = pure.unsafeRunSync()

分类:

构造器执行语义
IO.pure值已就绪,立即
IO.apply / IO(...)惰性,运行时空闲线程执行
IO.blocking阻塞 IO,走阻塞线程池(防饿死)
IO.defer延迟创建,用于自递归
IO.canceled取消型效果

3. 组合子:flatMap、for 与 mapN

顺序组合(for 推导,背后是 flatMap):

def fetchUser(id: Long): IO[Option[User]] = IO(...)
def fetchProfile(id: Long): IO[Profile] = IO(...)

val program: IO[UserView] = for {
  userOpt <- fetchUser(1L)
  user    <- userOpt match {
               case Some(u) => IO.pure(u)
               case None    => IO.raiseError(UserNotFound(1L))
             }
  profile <- fetchProfile(user.id)
} yield UserView(user, profile)

并行组合(mapN / parMapN):

import cats.syntax.parallel._

// 两个不相关的请求并行发,再合并
val parallel: IO[UserView] =
  (fetchUser(1L), fetchProfile(1L)).parMapN { (u, p) => UserView(u, p) }
组合子语义
map纯变换
flatMap / for顺序依赖
parMapN并行无依赖
race竞速取先完成
*> / <*丢弃前/后结果

4. 异步与并发原语

并发协调原语(cats.effect.concurrent → CE3 的 cats.effect):

import cats.effect.{IO, Ref, Deferred, Semaphore}

// Ref:并发安全的可变状态(替代 var)
val counter: IO[Ref[IO, Int]] = Ref.of[IO, Int](0)
val inc: IO[Unit] = counter.flatMap(_.update(_ + 1))

// Deferred:一次性栅栏(等另一个 fiber 完成)
val gate: IO[Deferred[IO, Int]] = Deferred[IO, Int]
val waitThenUse = gate.flatMap(_.get)                    // 阻塞直到被 fulfill

// Semaphore:限流(最多 N 并发)
val permit = Semaphore[IO](5).flatMap(_.acquire)         // 获取许可

Fiber 并发:

// 启动两个 fiber 并行处理,join 等待全部完成
val batch: IO[(Int, Int)] =
  for {
    fa <- IO(heavyTaskA).start       // 启动 fiber A
    fb <- IO(heavyTaskB).start       // 启动 fiber B
    a  <- fa.join                    // 等待 A
    b  <- fb.join
  } yield (a, b)

Fiber 是「轻量线程」(万级无压力),由效果系统调度,可取消、可超时、可组合——比裸 Future 更可控。


5. 资源安全:bracket 与 Finalizer

数据库连接、文件句柄必须保证释放,即使中途异常——bracket 三参数搞定:

import cats.effect.IO

def withConnection[A](conn: Connection)(use: Connection => IO[A]): IO[A] =
  IO.blocking(conn.open)     // acquire
    .bracket { c =>
      use(c)                 // use
    } { c =>
      IO.blocking(c.close)   // release(无论成败)
    }

CE3 更现代写法 Resource:

import cats.effect.Resource

def connResource: Resource[IO, Connection] =
  Resource.make(IO.blocking(open()))(c => IO.blocking(c.close))

val program: IO[Unit] = connResource.use { c =>
  runQuery(c)
}

释放语义表:

场景释放行为
正常完成release 执行
中途异常release 执行
协程被取消release 执行(可通过 onCancel 追加清理)
嵌套使用后开先关(栈式)

6. 错误处理:Either 与 Retry

显式错误路径:返回 IO[Either[DomainError, A]] 或 IO[A] + raiseError:

sealed trait DomainError
case object NotFound extends DomainError
case object RateLimit extends DomainError

val safe: IO[Either[DomainError, User]] =
  fetch(1L).attempt.map(_.left.map(_ => RateLimit))

Retry 与超时(CE 内置):

import cats.effect.{IO, Temporal}
import scala.concurrent.duration._

// 超时
val withTimeout: IO[User] = fetch(1L).timeout(5.seconds)

// 重试(模拟抖动退避)
def retryWithBackoff[A](io: IO[A], max: Int): IO[A] =
  io.retrySchedule(
    Schedule.exponential(100.millis).intersect(Schedule.recurs(max))
  )

生产经验:可重试的错误(超时/限流/5xx)与不可重试的错误(4xx 参数错/业务错)分开建模——别一股脑重试。


7. ZIO 的 R:ZLayer 依赖注入

ZIO 的签名是 ZIO[R, E, A]——多一个环境类型 R,把依赖编码进类型系统:

import zio._

trait UserRepo {
  def find(id: Long): Task[Option[User]]
}

// 依赖 UserRepo,产出 User
val program: ZIO[UserRepo, Throwable, Option[User]] =
  ZIO.serviceWithZIO[UserRepo](_.find(1L))

ZLayer 组合依赖:

object UserRepoLive {
  val layer: ZLayer[DB, Nothing, UserRepo] =
    ZLayer.fromFunction { (db: DB) =>
      new UserRepo {
        def find(id: Long): Task[Option[User]] =
          db.query("select * from users where id = ?", id)
      }
    }
}

// 组装:App 需要 DB,DB 从配置来
val fullLayer: ZLayer[Config, Nothing, UserRepo] =
  Config.live >>> DB.live >>> UserRepoLive.layer

依赖注入对比:

方式载体编译期检查
ZLayer类型安全、可组合✅ 缺依赖编译不过
手动传参函数参数手动
全局单例JVM 静态❌ 隐式共享状态

8. 调度与 Runtime

效果需要「执行环境」——CE 的 Runtime 与 ZIO 的 Runtime 负责调度 fiber 到线程池:

import cats.effect.unsafe.IORuntime

// CE3 默认全局 runtime(含计算池 + 阻塞池)
val runtime = IORuntime.global
val value: Int = program.unsafeRunSync()(runtime)

// 自定义配置:调线程池大小、队列
import cats.effect.unsafe.IORuntimeConfig
val cfg = IORuntimeConfig(
  cancelCheckThreshold = 1024,
  cpuStarvationCheckInitialDelay = 5.seconds
)

ZIO 推荐入口(不要 unsafeRun,用 main):

object App extends ZIOAppDefault {
  def run = for {
    _ <- Console.printLine("hello")
    user <- UserService.get(1L)
  } yield user
}

线程模型要点:效果系统的线程池是全局共享的少量工作线程,跑的是 fiber(轻量),不是每任务一线程——所以 IO.blocking 必须单独走阻塞池,否则计算池会被睡死。


9. 测试与效果系统

效果系统测试的杀手锏:不 mock,直接替换实现。

// 生产:查真库
val userRepo: UserRepo = new UserRepo { def find(id) = realDB.find(id) }

// 测试:换内存实现
val fakeRepo: UserRepo = new UserRepo {
  def find(id: Long) = ZIO.succeed(Some(User(id, "Test")))
}

// ZIO Test 支持虚拟时钟
test("超时测试用虚拟时间") {
  testClock
  assertZIO(program.timeout(1.minute))(isSome(...))
}

CE 测试:IO 直接在测试里运行(配合 /scala-testing-practice/ 的 MUnit CatsEffectSuite),无任何 mock 框架。


10. 效果系统选型速查表

需求Cats Effect IOZIO
核心模型IO[A]ZIO[R, E, A]
依赖注入手动/Cats 生态ZLayer 内置
错误Either / raiseErrorE 通道类型化
并发fiber + Ref/Deferred/Semaphorefiber + Ref/Queue/Promise
调度IORuntime(全局/自定义)ZIOAppDefault
测试MUnit CatsEffectSuiteZIO Test + TestClock
生态http4s、Doobie、FS2http4s、ZIO HTTP、Quill

一句话记忆:Cats Effect 更接近「纯函数式内核」,ZIO 多了内置 DI 与错误通道;有依赖图复杂度选 ZIO,想最小组件选 CE——两者都能把副作用锁进类型。


延伸阅读

  • /scala-functional-programming/ — Monad 与类型类,理解 IO 的数学基础
  • /scala-web-http-apps/ — http4s 全栈基于 Cats Effect
  • /scala-database-access/ — Doobie(CE)与 Quill(ZIO)数据库访问
  • /scala-testing-practice/ — 效果代码的测试写法
  • /actor-model-detailed-explanation/ — Fiber 与 Actor 两种并发模型对比

继续阅读

探索更多技术文章

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

全部文章 返回首页

「scala」更多文章

  1. Scala.js 与 Scala Native:跨平台编译、互操作与工程实践
  2. Scala 领域建模实战:ADT、类型驱动设计与模块化架构
  3. Scala 类型级编程实战:Shapeless、Match Types 与编译期计算