大多数流式框架把「背压」做成一个需要配置的旋钮:缓冲区多大、水位多高、满了之后丢还是阻塞。fs2 走的是另一条路——它把流建模成拉取(pull)驱动的惰性值:下游不主动拉取,上游就不会继续产生元素。背压因此不是外挂的机制,而是模型本身的性质,不需要任何缓冲参数就成立。
这种设计带来的直接收益是:同一个 Pipe 可以串在任意位置、parEvalMap 的并发度可以精确控制、资源释放可以跟着流的作用域走、错误可以在流的任意层级被拦截和重试。代价是它要求你按「拉取」而不是「推送」的直觉去思考——本文的多数篇幅都在讲这件事在实践中的具体表现。
1. Stream 与 Pipe 的语义
Stream[F, O] 是「在效果 F 中产生零到多个 O 的描述」。它本身不执行任何东西,只有 compile 之后才变成 F[...]。Pipe[F, I, O] 是一个类型别名,等价于 Stream[F, I] => Stream[F, O]——它把「一段可复用的处理逻辑」提升成一等值。
import cats.effect.*
import fs2.*
// Stream 是描述:下面两行都不产生任何副作用
val s1: Stream[IO, Int] = Stream.emits(List(1, 2, 3))
val s2: Stream[Pure, Int] = Stream(1, 2, 3) // 纯流,无需 F
// Pipe 是可复用的处理段
val double: Pipe[IO, Int, Int] = _.map(_ * 2)
val onlyEven: Pipe[IO, Int, Int] = _.filter(_ % 2 == 0)
// 组合:through 把 Pipe 串进流
val pipeline: Stream[IO, Int] = s1.through(double).through(onlyEven)
val io: IO[List[Int]] = pipeline.compile.toList // 此刻才执行
// Pipe 可以直接拼成一个复合 Pipe
val combined: Pipe[IO, Int, Int] = double.andThen(onlyEven)
类型速查:
□ Stream[F, O] 效果 F 中产生 O 的惰性流
□ Stream[Pure, O] 不涉及效果(可用 .covary[IO] 提升)
□ Pipe[F, I, O] = Stream[F, I] => Stream[F, O]
□ Sink[F, I] = Pipe[F, I, Unit](等价于管道末端)
□ Chunk[O] 批量元素,流的内部传输单元(减少逐元素开销)
□ Pull[F, O, R] 底层拉取原语,Stream 由它构建
工程要点:把「一段逻辑」写成 Pipe 而不是「一个接收 Stream 返回 Stream 的方法」,好处是它能被 .andThen 拼接、能被独立测试(用 Stream.emits(...).through(p).compile.toList 断言),也能在多个管道里复用。
2. 构造与常用组合子
fs2 的构造器覆盖了「已知集合」「重复求值」「时间驱动」「状态展开」四类来源。选择哪个构造器,决定了流的终止性与资源占用。
import cats.effect.*
import fs2.*
import scala.concurrent.duration.*
val fromList: Stream[Pure, Int] = Stream.emits(List(1, 2, 3))
val infinite: Stream[Pure, Int] = Stream.iterate(0)(_ + 1) // 无限
val fromEff: Stream[IO, String] = Stream.eval(IO.println("once").as("x"))
val repeat: Stream[IO, Long] = Stream.repeatEval(IO.realTime) // 无限重复
// 时间驱动:按固定间隔产生一个值(心跳、轮询、指标上报)
val tick: Stream[IO, Unit] = Stream.fixedRate[IO](1.second).void
// 状态展开:把「下一步」编码成 Option
val countdown: Stream[Pure, Int] =
Stream.unfold(5)(n => if n <= 0 then None else Some((n, n - 1)))
// 从资源构造:流结束时自动释放
val lines: Stream[IO, String] =
Stream.resource(Resource.fromAutoCloseable(
IO.blocking(scala.io.Source.fromFile("/etc/hosts"))))
.flatMap(src => Stream.fromIterator[IO](src.getLines(), 64))
常用组合子按用途分组:
转换:map / filter / collect / flatMap / evalMap / mapChunks
批量:chunkN(n) / chunkAll / grouped(n) / groupWithin(n, d)
限速:metered(rate) / awakeEvery / fixedRate / fixedDelay
去重:distinct / distinctBy / changes(相邻去重)
累积:scan / mapAccumulate / fold / foldMonoid
终结:compile.toList / .drain / .count / .last
工程要点:Stream.repeatEval 与 Stream.fixedRate 是无限流,单独 compile.toList 会挂死。它们必须与 take、interruptAfter、through(sink) 或 merge 一起使用。写测试时给无限流加 .take(n) 是最省事的护栏。
3. 背压的真实表现
背压的验证方法是:让上游产生速度远快于下游处理速度,观察内存是否增长。在 fs2 里,这个实验的结论是内存平稳——因为上游的 Pull 只有在下游请求时才会被求值。
import cats.effect.*
import fs2.*
import scala.concurrent.duration.*
// 上游每秒产 1000 个,下游每个耗时 10ms(每秒只能处理 100 个)
val fastProducer: Stream[IO, Int] =
Stream.range(0, 1_000_000).covary[IO].metered(1.millis)
val slowConsumer: Pipe[IO, Int, Int] =
_.evalMap(i => IO.sleep(10.millis).as(i))
val balanced: Stream[IO, Int] = fastProducer.through(slowConsumer)
// 结果:整体耗时由下游决定,内存不增长,上游被自然拖慢
什么时候需要显式缓冲:
□ 上游有突发性(burst),下游处理速率稳定
→ buffer(n) 吸收抖动,n 取「突发峰值 × 持续时间 / 处理速率」
□ 上游是网络/磁盘,单次拉取成本高
→ prefetch 提前拉取,隐藏 IO 延迟
□ 下游是批处理,希望减少批次次数
→ chunkN(n) 或 grouped(n)
□ 并行处理时希望打散边界
→ parEvalMap(n) 内部自带小缓冲
不要做的事:
✗ 用 buffer 解决「下游太慢」——那只是把问题推迟成内存问题
✗ 用无界队列做「解耦」——那是把背压彻底关掉
工程要点:buffer(n) 与 prefetch 是性能优化,不是背压的修复手段。判断标准很简单:加上 buffer 之后内存是否随运行时长增长。若增长,说明下游处理能力确实不足,应当增加并发度或提高单次处理效率,而不是继续加大缓冲。
4. 并发原语
fs2 的并发原语都是「结构化」的:并发的子流会随主流结束而结束,不会留下孤儿 fiber。选哪个原语取决于「顺序是否需要保持」与「并发单元是什么」。
import cats.effect.*
import fs2.*
import scala.concurrent.duration.*
val src: Stream[IO, Int] = Stream.range(0, 100).covary[IO]
// parEvalMap:限并发、保序
val ordered: Stream[IO, Int] = src.parEvalMap(8)(slow)
// parEvalMapUnordered:限并发、不保序(吞吐更高)
val unordered: Stream[IO, Int] = src.parEvalMapUnordered(8)(slow)
// parJoin:把一个「流的流」并发展开,总并发受 maxConcurrent 限制
val sharded: Stream[IO, Int] =
Stream.emits(List(Stream.range(0, 50), Stream.range(50, 100)).map(_.covary[IO]))
.parJoin(2)
// merge:两个流并发合流,谁先产生谁先出
val merged: Stream[IO, Int] = src.merge(Stream.range(200, 250).covary[IO])
// concurrently:让「伴随流」与主流并行运行,主流结束则伴随流被取消
val withMetrics: Stream[IO, Int] =
src.concurrently(Stream.fixedRate[IO](1.second).evalMap(_ => report()))
def slow(i: Int): IO[Int] = IO.sleep(10.millis).as(i)
def report(): IO[Unit] = IO.println("tick")
并发原语对照:
□ parEvalMap(n) 逐元素并行,保序,最常用
□ parEvalMapUnordered(n) 逐元素并行,不保序,吞吐更高
□ parJoin(n) 并发展开子流,适合分片/多源
□ merge / mergeHaltBoth 合流,顺序不确定
□ concurrently 伴随流(指标、心跳、日志刷盘)
□ broadcast(n) 一份输入广播到 n 个消费者
□ hold / Topic 与外部(如 HTTP 接口)共享最新值
工程要点:parEvalMap 的并发度必须由下游资源上限反推——如果下游是数据库,n 不能超过连接池的 maximumPoolSize;如果下游是外部 API,n 不能超过对方的配额。把 n 设成「越大越快」是生产事故的常见来源。
5. 错误处理与重试
fs2 的错误处理原则是:在能处理错误的层级处理它。attempt 把错误变成值以便继续流,handleErrorWith 换一条流,retry 系列在元素级重试,而不可恢复的错误应当让整个流失败而不是吞掉。
import cats.effect.*
import fs2.*
import scala.concurrent.duration.*
final case class Record(id: String, payload: String)
def parse(r: String): IO[Record] =
IO(Record(r.takeWhile(_ != ':'), r)).adaptError(e =>
new IllegalArgumentException(s"bad record: $r"))
// 元素级:解析失败不中断流,转成 Either 继续
val tolerant: Pipe[IO, String, Either[String, Record]] =
_.evalMap(r => parse(r).attempt.map(_.left.map(_.getMessage)))
// 元素级重试:只对可恢复错误重试,带退避
val retrying: Pipe[IO, Record, Record] =
_.evalMap(r => persist(r).retryWithDelay(3, 200.millis))
// 流级:整条流失败时切换降级源(打印后返回空流,不中断应用)
val fallback: Stream[IO, Record] =
Stream.emits(List("a:1", "b:2")).covary[IO].through(parsePipe)
.handleErrorWith(e => Stream.exec(IO.println(s"degraded: $e")))
// 只对特定错误兜底
val selective: Stream[IO, Record] =
parsePipe(Stream.emits(List("x:1")).covary[IO])
.handleErrorWith {
case _: IllegalArgumentException => Stream.empty // 脏数据,跳过
case e => Stream.raiseError[IO](e) // 其余继续抛
}
def persist(r: Record): IO[Unit] = IO.unit
def parsePipe: Pipe[IO, String, Record] = _.evalMap(parse)
错误处理决策表:
错误性质 处理方式 示例
可恢复(网络/超时) 元素级重试 + 退避 JDBC 连接超时、HTTP 502
不可恢复(格式错误) 转 Either / 死信,继续流 反序列化失败、校验不通过
致命(配置/权限) 让流失败,快速退出 认证失败、表不存在
上游中断 不重试,让取消传播 取消异常不可重试
工程要点:retryWithDelay 只对幂等操作安全。若元素是「扣款」「发券」这类非幂等动作,重试前必须先做幂等键去重,否则一次网络抖动就是一次重复扣款。判据是「同一个元素执行两次的结果是否与执行一次相同」。
6. 资源与流的生命周期
Stream 与 Resource 的融合点是 Stream.resource 与 Stream.bracket:资源在流开始消费时获取,在流结束(正常/失败/取消)时释放。这使得「每处理一批就打开关闭一个连接」这类需求变得自然。
import cats.effect.*
import fs2.*
// 把 Resource 织入流:资源随流的作用域存活
val toFile: Pipe[IO, String, Unit] =
_.flatMap { line =>
Stream.resource(
Resource.make(
IO.blocking(new java.io.PrintWriter("out.txt"))
)(w => IO.blocking(w.close()))
).flatMap(w => Stream.eval(IO.blocking(w.println(line))))
}
// 用 chunks.evalMap 表达「每个 chunk 共享一次资源」
val perChunk: Pipe[IO, Int, Int] =
_.chunks.evalMap(chunk => IO.blocking(chunk.toList.sum))
// 长生命周期资源:整个流共用一份
def withPool[A](use: String => Stream[IO, A]): Stream[IO, A] =
Stream.resource(poolResource).flatMap(use)
def poolResource: Resource[IO, String] =
Resource.make(IO.println("open pool").as("pool"))(_ => IO.println("close pool"))
资源与流的组合方式:
□ Stream.resource(r):资源作用于「它之后」的流片段
□ Stream.bracket(acquire)(use)(release):更底层的三段式
□ flatMap 中的 Resource:每个元素各自获取释放(细粒度,开销大)
□ 外层 Resource.use { ... 整个流 ... }:整个流共用一份(粗粒度,推荐)
□ Stream.exec(f):只在流末尾执行一次的效果(如 flush)
工程要点:细粒度资源(每个元素一个连接)在元素速率高时会成为瓶颈——获取连接的开销可能比处理本身还大。正确的分层是:池用外层 Resource 管,连接按 chunk 借还,用 chunks.evalMap 让一批元素共享一次连接。
7. 与文件、消息队列的集成
fs2-io 提供了文件与网络的流式读写,fs2-kafka 提供了 Kafka 的消费与生产。两者的共同点是都遵守「资源随流生命周期」的约定。
import cats.effect.*
import fs2.*
import fs2.io.file.{Files, Path}
import fs2.io.file.Path as FsPath
// 文件:流式读取大文件,内存恒定
val fileStream: Stream[IO, String] =
Files[IO].readAll(FsPath("/var/log/app.log"))
.through(fs2.text.utf8.decode)
.through(fs2.text.lines)
// 文件:流式写入(追加),自动关闭
val sink: Pipe[IO, String, Unit] =
_.intersperse("\n").through(fs2.text.utf8.encode)
.through(Files[IO].writeAll(FsPath("out.txt")))
// 目录遍历:递归列出所有文件,逐个处理
val allLogs: Stream[IO, FsPath] =
Files[IO].walk(FsPath("/var/log"))
.filter(_.extName == ".log")
import fs2.kafka.*
import scala.concurrent.duration.*
// Kafka:消费 → 转换 → 生产,位点随流提交
def kafkaPipe(consumer: ConsumerSettings[IO, String, String],
producer: ProducerSettings[IO, String, String]): Stream[IO, Unit] =
KafkaConsumer.stream(consumer)
.subscribeTo("orders")
.records
.map(_.value.toUpperCase)
.through(KafkaProducer.pipe(producer))
工程要点:文件流必须走 Files[IO] 而不是 java.io 的阻塞读——fs2.io.file 把读操作放在 blocking 池并做分块,直接 new FileInputStream 会阻塞 compute 池。Kafka 侧的位点提交要与处理完成绑定,具体策略见下方延伸阅读里的专门文章。
8. 有状态聚合与窗口
fs2 的有状态组合子都是局部的:状态属于某一段流,随流结束而消失。要做跨分区、可恢复的状态,必须把状态外置或依赖有状态流引擎。
import cats.effect.*
import fs2.*
import scala.concurrent.duration.*
final case class Trade(symbol: String, qty: BigDecimal)
// scan:保留累积值并逐个输出(适合做移动指标)
val runningTotal: Stream[Pure, BigDecimal] =
Stream.emits(List(Trade("A", 10), Trade("A", 5), Trade("B", 2)))
.scan(BigDecimal(0))((acc, t) => acc + t.qty)
// mapAccumulate:显式状态 + 输出
val withPrev: Stream[Pure, (BigDecimal, Trade)] =
Stream.emits(List(Trade("A", 10), Trade("A", 5)))
.mapAccumulate(BigDecimal(0)) { (acc, t) =>
val next = acc + t.qty
(next, (next, t))
}
// groupWithin:按「数量或时间」双触发窗口
val windowed: Stream[IO, Map[String, BigDecimal]] =
Stream.emits(List(Trade("A", 10), Trade("B", 3))).covary[IO]
.repeat.take(1000)
.groupWithin(100, 2.seconds) // 100 个或 2 秒,先到者触发
.map(chunk => chunk.groupBy(_.symbol).view.mapValues(_.map(_.qty).sum).toMap)
// changes:相邻去重(去抖动)
val debounced: Stream[Pure, Int] =
Stream.emits(List(1, 1, 2, 2, 2, 3)).changes
状态组合子的容量与代价:
□ scan / mapAccumulate 状态大小由用户控制,可能无界增长
□ groupWithin 窗口内元素全在内存,size 决定内存上限
□ fold / foldMonoid 仅在流结束时输出,无中间结果
□ changes 只保留上一个元素,恒定内存
□ distinct 内部维护已见集合,无界增长(大流慎用)
工程要点:groupWithin 的窗口大小是内存上限的直接来源——groupWithin(100_000, 1.second) 意味着最坏情况下十万个元素同时在内存里。按元素平均大小反推一个安全的 size,比按「想攒多大批」来设更靠谱。distinct 在长流上会持续吃内存,能用 changes(相邻去重)就不要用 distinct。
9. 测试与性能
fs2 的测试优势在于流是纯值:可以只测 Pipe,也可以把 compile 的耗时压到微秒级。
import cats.effect.*
import fs2.*
import munit.CatsEffectSuite
class PipeSuite extends CatsEffectSuite:
test("onlyEven 过滤奇数") {
val p: Pipe[IO, Int, Int] = _.filter(_ % 2 == 0)
Stream.emits(List(1, 2, 3, 4)).covary[IO].through(p).compile.toList
.map(result => assertEquals(result, List(2, 4)))
}
test("失败路径也释放资源") {
val released = scala.collection.mutable.ListBuffer.empty[String]
val res = Resource.make(IO(released += "acquire").as("r"))(_ => IO(released += "release"))
Stream.resource(res).flatMap(_ => Stream.raiseError[IO](new RuntimeException("boom")))
.compile.drain.attempt.map { _ =>
assertEquals(released.toList, List("acquire", "release"))
}
}
性能抓手:
- chunk 大小:默认 64KB 左右,chunkN 让下游按批处理减少逐元素开销
- 并发度:parEvalMap 的 n 由下游资源上限反推
- 缓冲:buffer/prefetch 只用于吸收抖动,不用于「解决慢下游」
- 阻塞调用:一律 IO.blocking,否则拖死 compute 池
- 度量:统计每秒元素数、单元素耗时 p99、chunk 平均大小
- 反模式:在 evalMap 里做逐元素的同步 IO(应当 chunkN 后批量)
工程要点:性能问题十有八九出在**「在元素级做了本应批量做的事」**:逐元素一次数据库 round-trip、逐元素一次 HTTP 请求、逐元素一次文件 flush。先用 chunkN 聚合,再用 evalMap 批量处理,通常就能把吞吐提高一个数量级。
10. 速查表与一句话记忆
| 问题 | 一句话答案 |
|---|---|
| 背压怎么实现 | 拉取模型天然背压,不需要缓冲参数 |
| 保序并发 | parEvalMap(n) |
| 吞吐优先并发 | parEvalMapUnordered(n) |
| 并发度怎么定 | 由下游资源上限反推,不是越大越好 |
| 解析失败怎么办 | 转 Either 或死信,别中断整条流 |
| 重试的前提 | 操作必须幂等,配指数退避 |
| 资源怎么管 | Stream.resource 或外层 Resource 包住整个流 |
| 批量处理 | chunks.evalMap,一批共享一次连接 |
| 窗口 | groupWithin(n, d),n 决定内存上限 |
| 阻塞调用 | 必须 IO.blocking |
一句话记忆:fs2 = 惰性拉取流(背压免费)+ Pipe 可复用处理段 + parEvalMap 控并发(保序)/Unordered(追吞吐)+ 元素级错误转值或死信 + 资源随流作用域释放 + chunks 批量降低 round-trip——核心心法是把「推送」的直觉换成「拉取」,把「缓冲」当成优化而非解药。
延伸阅读
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。