fs2 流处理实战:Stream、Pipe 与背压模型

流处理的难点从来不在「怎么读数据」,而在背压、并发、错误恢复与资源释放这几件相互牵制的事。fs2 把流建模成惰性的拉取式值:下游不拉取,上游就停,背压因此不是需要调参的配置项,而是模型本身的性质。本文讲清 Stream 与 Pipe 的语义、构造与组合子、并发原语的取舍、错误处理与重试策略、与文件与消息队列的集成,以及有状态聚合与窗口的实现方式。

大多数流式框架把「背压」做成一个需要配置的旋钮:缓冲区多大、水位多高、满了之后丢还是阻塞。fs2 走的是另一条路——它把流建模成拉取(pull)驱动的惰性值:下游不主动拉取,上游就不会继续产生元素。背压因此不是外挂的机制,而是模型本身的性质,不需要任何缓冲参数就成立。

这种设计带来的直接收益是:同一个 Pipe 可以串在任意位置、parEvalMap 的并发度可以精确控制、资源释放可以跟着流的作用域走、错误可以在流的任意层级被拦截和重试。代价是它要求你按「拉取」而不是「推送」的直觉去思考——本文的多数篇幅都在讲这件事在实践中的具体表现。

前置:流处理与响应式流基础 、Cats Effect 运行时与 fibers 。

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——核心心法是把「推送」的直觉换成「拉取」,把「缓冲」当成优化而非解药。

延伸阅读

继续阅读

探索更多技术文章

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

全部文章 返回首页

「scala」更多文章

  1. Scala CLI 与工具链现代化:指令声明、打包发布与 CI 集成
  2. 持久化与事件溯源:journal、快照与 CQRS 读模型
  3. JSON 编解码与性能:circe 派生、jsoniter-scala 代码生成与选型