一个进程里真正稀缺的不是 CPU,而是那些必须被归还的外部资源:数据库连接、文件句柄、网络套接字、线程池、临时目录、后台 fiber。它们的共同点是「获取与释放必须成对,且释放必须发生在所有退出路径上」——成功、异常、被取消、超时,一条都不能漏。手工写 try/finally 在单层逻辑里没问题,一旦逻辑变成多层嵌套、跨 fiber 并发、被上层取消,遗漏几乎是必然的。
Cats Effect 的 Resource[F, A] 把「获取」与「释放」建模成一个可组合的值:释放逻辑跟着值的类型走,组合时自动嵌套,取消时由运行时保证执行。本文从语义讲起,给出释放顺序的准确规则,讨论它与手工清理的取舍边界,最后落到连接池、后台任务与优雅停机的完整实现。
1. 从 bracket 到 Resource
最底层的原语是 bracket:给定「获取」「使用」「释放」三段,无论使用段是正常结束、抛错还是被取消,释放段都会执行一次。Resource 是它的可组合形态——把三段拆成「一段描述」和一个 use 消费点,从而能在不执行的情况下自由拼接。
import cats.effect.*
import scala.concurrent.duration.*
// bracket:一次性、不可组合
val readOnce: IO[String] =
IO.blocking(new java.io.FileInputStream("/etc/hosts"))
.bracket { in => IO.blocking(scala.io.Source.fromInputStream(in).mkString) } { in =>
IO.blocking(in.close())
}
// Resource:描述与使用分离,可拼接
val readRes: Resource[IO, String] =
Resource.fromAutoCloseable(IO.blocking(new java.io.FileInputStream("/etc/hosts")))
.evalMap(in => IO.blocking(scala.io.Source.fromInputStream(in).mkString))
val program: IO[String] = readRes.use(s => IO.pure(s.take(100)))
bracket 的语义要点:
□ acquire:获取资源,本身可被取消(取消点)
□ use:使用资源,可正常返回、抛错、被取消
□ release:清理,无论何种退出路径都执行一次
□ release 必须「不可取消」:运行时用 uncancelable 包住
□ release 抛错:不会掩盖 use 的原始错误(除非 use 正常结束)
□ release 执行时机:紧随 use 结束,而非程序结束
工程要点:bracket 与 Resource 的差别是**「何时决定如何使用」**。资源获取需要组合、需要被多个使用者共享、需要在 for 推导里逐步构建时,用 Resource;只有一个直白的「打开—用完—关闭」时,bracket 就够。别为了「看起来函数式」把单次资源也包成 Resource。
2. Resource.make 与释放顺序
Resource.make(acquire)(release) 是最常用的构造器。它的释放顺序遵循后进先出(LIFO):Resource 组合得越晚,释放得越早——和栈的行为一致,这也是嵌套资源最自然的清理次序。
import cats.effect.*
final case class Conn(id: Int)
final case class Tx(id: Int)
def mkConn(id: Int): Resource[IO, Conn] =
Resource.make(IO.println(s"acquire conn $id").as(Conn(id)))(c =>
IO.println(s"release conn ${c.id}"))
def mkTx(c: Conn): Resource[IO, Tx] =
Resource.make(IO.println(s"begin tx on ${c.id}").as(Tx(c.id)))(t =>
IO.println(s"commit/rollback tx ${t.id}"))
// 组合:先拿连接,再开事务;释放时先提交事务,再还连接
val layered: Resource[IO, (Conn, Tx)] =
for
c <- mkConn(1)
t <- mkTx(c)
yield (c, t)
layered.use(_ => IO.println("do work"))
// 输出:acquire conn 1 / begin tx on 1 / do work / commit tx 1 / release conn 1
释放顺序规则:
□ for 推导里「后 acquire 的」先 release(LIFO)
□ Resource.flatMap 的嵌套:内层 Resource 的 release 先于外层
□ Resource.both(a, b):b 先 release,再 a
□ parZip / parTupled:两侧并发 acquire,release 顺序不保证
□ 多个独立 Resource 用 for 串起来 = 显式指定了释放次序
工程要点:释放顺序在依赖场景里是正确性的一部分——事务必须在连接归还之前结束,子进程必须在父进程之前杀掉,临时文件必须在目录删除之前清理。用 for 推导把依赖关系写清楚,别把两个独立的 Resource 用 both 平铺,那样释放顺序就丢了语义。
3. 组合器:eval、both、parZip 与递归
Resource 的组合器覆盖了绝大多数组装需求:eval 提升一个 F[A] 为「无释放动作的资源」,both 组合两个资源,parZip 并发获取,evalMap 在资源上追加效果,>> 丢弃左值。
import cats.effect.*
import cats.syntax.all.*
import scala.concurrent.duration.*
val pool: Resource[IO, String] =
Resource.make(IO.println("open pool").as("pool"))(_ => IO.println("close pool"))
val warmup: Resource[IO, Unit] =
Resource.eval(IO.println("warm up connections")) // 只有效果,无释放
val service: Resource[IO, String] =
for
p <- pool
_ <- warmup // 无释放资源,可插在任意位置
_ <- Resource.eval(IO.println(s"service ready on $p"))
yield p
// 并发获取两个独立资源(各自独立、无依赖时才用)
val two: Resource[IO, (String, String)] =
(pool, pool).parTupled
// 递归资源:如「按目录逐层打开」,深度由数据决定
def walk(dir: java.io.File): Resource[IO, List[java.io.File]] =
Resource.eval(IO.blocking(dir.listFiles().toList)).flatMap { children =>
children.foldLeft(Resource.pure[IO, List[java.io.File]](Nil)) { (acc, f) =>
if f.isDirectory then (acc, walk(f)).mapN(_ ++ _)
else acc.map(_ :+ f)
}
}
工程要点:parZip 只在两个资源互相独立时使用。若第二个资源的获取依赖第一个,必须用 flatMap 串起来,否则第二个可能在第一个完成前就开始获取,或者释放顺序被并发打乱。
4. Scope 与资源作用域
Resource 的生命周期由 use 划定,但有些场景需要把作用域显式地「开一个口子」:手动控制分配与释放、把资源绑定到某个 fiber 的存活期、或者跨多个 use 复用同一个 Resource 值。
import cats.effect.*
import cats.effect.kernel.Resource.ExitCase
import scala.concurrent.duration.*
object ScopeDemo extends IOApp.Simple:
val res: Resource[IO, String] =
Resource.makeCase(IO.println("open").as("handle")) {
case (h, ExitCase.Succeeded) => IO.println(s"$h closed (success)")
case (h, ExitCase.Errored(e)) => IO.println(s"$h closed (error: ${e.getMessage})")
case (h, ExitCase.Canceled) => IO.println(s"$h closed (canceled)")
}
def run: IO[Unit] =
res.allocatedCase.flatMap { case (h, release) =>
IO.println(s"using $h") *> release(ExitCase.Succeeded)
}
// useForever:资源与整个应用同生共死(服务器、后台 fiber)
val server: Resource[IO, Unit] =
Resource.make(IO.println("start server"))(_ => IO.println("stop server"))
.useForever
生命周期相关 API:
□ use(f):作用域限定在 f 内,f 结束即释放
□ useForever:资源活到程序结束(服务器、常驻消费者)
□ allocated:手动拿到 (资源, 释放函数),自己负责调用释放
□ allocatedCase:释放时能拿到 ExitCase(成功/失败/取消)
□ makeCase:release 可区分退出原因,做针对性清理
□ scope / Scope[F]:把资源绑定到某段 fiber 的生命周期
工程要点:优先用 use,它把「忘记释放」这件事从可能变成不可能。allocated 只在框架层使用(比如把资源交给一个不接收 Resource 的老式 API),用完后必须在所有退出路径上调用释放函数,此时你已经退回到手工管理的风险区。
5. 与 try/finally 的取舍
try/finally 不是「错的」,它在同步、单线程、无取消语义的代码里完全够用。问题出在三个维度:取消、异步释放、组合。
| 维度 | try/finally | Resource |
|---|---|---|
| 异常路径 | 覆盖 | 覆盖 |
| 取消路径 | 不覆盖(fiber 被取消时不一定执行) | 覆盖(运行时保证) |
| 异步释放 | 需要阻塞等待 | 原生支持 F[Unit] 释放 |
| 多层嵌套 | 手工维护顺序 | 组合自动嵌套 |
| 资源复用 | 需自行抽象 | 值可复用、可参数化 |
| 释放抛错 | 会掩盖原始异常 | 保留原始异常 |
| 纯函数性 | 副作用混在语法里 | 描述与执行分离 |
// 反例:取消发生时 finally 可能不执行,连接泄漏
def bad(ec: ExecutionContext): Unit =
val conn = pool.getConnection()
try process(conn)
finally conn.close() // 线程被中断/取消时不可靠
// 正例:释放走 blocking 池,取消安全,顺序由 Resource 保证
val good: Resource[IO, java.sql.Connection] =
Resource.make(IO.blocking(pool.getConnection))(c => IO.blocking(c.close()))
工程要点:判断标准是**「释放动作是否需要异步、是否可能被取消」**。纯同步、局部、生命周期极短的资源(一个临时 StringBuilder 之类)用 try/finally 无妨;凡涉及 IO、网络、跨 fiber、跨函数边界的资源,一律走 Resource。
6. 连接池的完整落地
连接池是 Resource 最典型的生产用例:池本身是一个长生命周期资源,借出的连接是一个短生命周期资源,两者层级不同,释放顺序天然由嵌套决定。
import cats.effect.*
import com.zaxxer.hikari.{HikariConfig, HikariDataSource}
import java.sql.Connection
import scala.concurrent.duration.*
final case class DbPool(ds: HikariDataSource):
def connection: Resource[IO, Connection] =
Resource.make(IO.blocking(ds.getConnection))(c => IO.blocking(c.close()))
object DbPool:
def create(url: String, maxPool: Int): Resource[IO, DbPool] =
Resource.make {
IO.blocking {
val cfg = new HikariConfig()
cfg.setJdbcUrl(url)
cfg.setMaximumPoolSize(maxPool)
cfg.setConnectionTimeout(3000) // 借连接超时,避免无限等待
cfg.setLeakDetectionThreshold(20000) // 泄漏检测:20 秒未归还即告警
new DbPool(new HikariDataSource(cfg))
}
}(p => IO.blocking(p.ds.close()))
// 启动时探活:连接池能建立不等于数据库可用
def verified(p: DbPool): Resource[IO, DbPool] =
Resource.eval(
p.connection.use(c => IO.blocking(c.isValid(2)))
.flatMap(ok => IO.raiseUnless(ok)(new RuntimeException("db not reachable")))
).as(p)
object App extends IOApp.Simple:
val pool: Resource[IO, DbPool] =
DbPool.create("jdbc:postgresql://localhost:5432/app", 16).flatMap(DbPool.verified)
def run: IO[Unit] =
pool.use { p =>
p.connection.use { conn =>
IO.blocking {
val st = conn.prepareStatement("select count(*) from orders")
val rs = st.executeQuery()
rs.next(); rs.getInt(1)
}.flatMap(n => IO.println(s"orders = $n"))
}
}
工程要点:maximumPoolSize 必须与下游数据库的 max_connections 以及应用实例数联合计算——实例数 × 池上限 不能超过数据库连接上限。连接借出必须走 Resource,绝不要手工 getConnection 后忘记 close;leakDetectionThreshold 是发现这类遗漏最便宜的开关。
7. 常见陷阱与规避
陷阱一:Resource 值被复用但被当作「每次调用都新建」
Resource 是描述,不是实例;use 几次就 acquire/release 几次(正确)
但如果把它提升成 val 在对象里共享,语义仍是「每次 use 都新建」
陷阱二:release 里抛错吞掉了 use 的错误
release 抛错且 use 也失败时,use 的错误优先保留,release 的错误被抑制
对策:release 里用 handleErrorWith 记录日志,别让它抛
陷阱三:acquire 被取消留下半初始化状态
acquire 里「注册回调 + 返回句柄」不是原子操作
对策:用 Resource.make 的 acquire 内部保证原子,或 uncancelable 包裹
陷阱四:把 Resource 当缓存
Resource 不缓存已分配的值;需要复用请用 Resource.eval + Ref,或 allocated
陷阱五:释放顺序与依赖相反
内层资源被外层依赖时,用 flatMap 而非 both/parZip
陷阱六:release 走 compute 池
JDBC close / 文件 flush 是阻塞调用,必须 IO.blocking
// 陷阱三的规避:acquire 内部保证「分配 + 注册」原子
def atomicRes: Resource[IO, String] =
Resource.make {
IO.uncancelable { poll =>
for
h <- IO.println("allocate").as("handle")
_ <- poll(IO.sleep(50.millis)) // 这里才允许取消
yield h
}
}(_ => IO.println("release"))
工程要点:Resource.make 的 acquire 默认已经是不可取消的(CE3 内部实现),你只有在 acquire 里显式插入了「可取消的等待」(比如等一个信号量)时才需要自己用 uncancelable + poll 划出可取消点。
8. 后台任务与优雅停机
常驻的消费循环、定时任务、连接保活 fiber 都是「与程序同生命周期」的资源。它们应当用 Resource 建模:启动即 acquire,停机即 release,且 release 要能等待在途工作完成。
import cats.effect.*
import scala.concurrent.duration.*
def backgroundWorker(name: String): Resource[IO, Unit] =
for
stop <- Resource.eval(Deferred[IO, Unit])
fiber <- Resource.make(
(IO.println(s"$name started") *>
stop.get.flatMap(_ => IO.println(s"$name draining"))).start
)(f => stop.complete(()).void *> f.join.timeoutTo(30.seconds, IO.unit))
yield ()
object Graceful extends IOApp.Simple:
val app: Resource[IO, Unit] =
for
_ <- backgroundWorker("consumer")
_ <- Resource.eval(IO.println("serving traffic"))
yield ()
def run: IO[Unit] =
// IOApp 收到 SIGTERM 会取消 run 的 fiber,
// 取消沿 Resource 的 use 边界传播 → 触发 release → 等待在途任务
app.useForever
工程要点:优雅停机的关键是**「停机信号要能传到 release」**。IOApp 在收到 SIGTERM 时会取消主 fiber,取消会触发 Resource 的释放;因此 release 里要显式 join 在途 fiber 并加超时上限,避免一个卡死的任务把整个进程的停机拖过 Kubernetes 的 terminationGracePeriodSeconds。
9. 测试与可观测
Resource 的可测试性来自它的纯描述特性:可以只 acquire 不 use,也可以断言 acquire/release 的次数与顺序。
import cats.effect.*
import cats.effect.std.Queue
// 用 Queue 记录生命周期事件,断言顺序
def traced(log: Queue[IO, String]): Resource[IO, String] =
Resource.make(log.offer("acquire").as("r"))(_ => log.offer("release"))
val check: IO[Unit] =
for
log <- Queue.unbounded[IO, String]
_ <- traced(log).use(_ => log.offer("use"))
ev <- log.take.replicateA(3).map(_.toList)
_ <- IO(assert(ev == List("acquire", "use", "release"), s"顺序错误: $ev"))
yield ()
// 测试「失败路径也会释放」
val failCase: IO[Unit] =
for
log <- Queue.unbounded[IO, String]
_ <- traced(log).use(_ => IO.raiseError(new RuntimeException("boom"))).attempt
ev <- log.take.replicateA(2).map(_.toList)
_ <- IO(assert(ev == List("acquire", "release")))
yield ()
可观测性抓手:
- 池指标:active/idle/total/pending 连接数(Hikari 暴露为 JMX)
- 借出耗时:getConnection 的 p99,突增说明池过小或下游慢
- 泄漏:leakDetectionThreshold 触发的告警次数
- 生命周期日志:acquire/release 打点,配对统计(差值为泄漏量)
- 停机:release 耗时分布,超过 grace period 的实例数
工程要点:把「acquire 与 release 的配对」当成一个可观测指标来做:按资源类型统计 acquire 次数与 release 次数,两者长期不等就是泄漏。这比事后翻日志找未关闭的连接高效得多。
10. 速查表与一句话记忆
| 问题 | 一句话答案 |
|---|---|
| 用 bracket 还是 Resource | 单次直白用 bracket,需组合/复用/分层用 Resource |
| 释放顺序 | LIFO:后 acquire 的先 release |
| 依赖资源怎么组合 | 用 flatMap(for 推导),不要用 both/parZip |
| 资源活到程序结束 | useForever |
| 手动控制分配释放 | allocated / allocatedCase(框架层专用) |
| 取消时会释放吗 | 会,运行时用 uncancelable 包住 release |
| release 能抛错吗 | 尽量不抛,抛错会被抑制但会掩盖问题 |
| 连接池怎么建模 | 池是长生命周期 Resource,连接是短生命周期 Resource |
| 优雅停机 | 取消触发 use 边界释放,release 里 join 在途 fiber 并设超时 |
一句话记忆:Resource = 获取与释放绑定成可组合的值 + 嵌套自动 LIFO 释放 + 取消路径也保证执行 + 释放不可取消且不掩盖原始异常——池与连接分层嵌套、后台 fiber 用 useForever 活到停机、停机靠取消传播到 release;凡涉及 IO、跨 fiber、跨边界的资源,永远别退回手工 try/finally。
延伸阅读
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。