引言
「纯函数式程序」听上去很学术,落到生产只有一个问题:副作用(网络、数据库、时间)怎么办? 效果系统(Effect System)的答案是把副作用包装成值——IO[User] 不是「拿到一个 User」,而是「一份能拿到 User 的程序描述」。你在纯代码里组合这份描述,最后才在唯一的边界执行它。Cats Effect 与 ZIO 是 JVM 上两套主流实现,本文讲透核心模型并给出可直接运行的写法。
前置:/scala-functional-programming/(Monad/类型类)。配合 Web 见 /scala-web-http-apps/(http4s 即基于 Cats Effect)。
目录
- 1. 为什么需要效果系统
- 2. IO Monad:副作用包装成值
- 3. 组合子:flatMap、for 与 mapN
- 4. 异步与并发原语
- 5. 资源安全:bracket 与 Finalizer
- 6. 错误处理:Either 与 Retry
- 7. ZIO 的 R:ZLayer 依赖注入
- 8. 调度与 Runtime
- 9. 测试与效果系统
- 10. 效果系统选型速查表
- 延伸阅读
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 IO | ZIO |
|---|---|---|
| 核心模型 | IO[A] | ZIO[R, E, A] |
| 依赖注入 | 手动/Cats 生态 | ZLayer 内置 |
| 错误 | Either / raiseError | E 通道类型化 |
| 并发 | fiber + Ref/Deferred/Semaphore | fiber + Ref/Queue/Promise |
| 调度 | IORuntime(全局/自定义) | ZIOAppDefault |
| 测试 | MUnit CatsEffectSuite | ZIO Test + TestClock |
| 生态 | http4s、Doobie、FS2 | http4s、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 两种并发模型对比
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。