引言
「流处理」不是大数据专属——日志采集、消息队列消费、文件管道、实时聚合,处处都在处理无限数据序列。Scala 生态有两套主流流库:Akka Streams(Reactive Streams 标准实现,图式 DSL)与 fs2(纯函数式流,天然嵌合 Cats Effect)。两者共享同一个灵魂:背压(Backpressure)——慢的消费者不能被快的生产者淹没。本文先讲清流的心智模型与反压机制,再分别深入两套流库的核心用法、错误处理、并发流、外部系统桥接与测试。
前置:/scala-functional-effects/(IO/Fiber 与并发)、/scala-bigdata-spark/(批处理与 DAG 对比)、/scala-collections/(集合与惰性求值)。分布式相关见 /scala-akka-cluster/。
目录
- 1. 流处理的心智模型:数据流、Stage 与背压
- 2. Akka Streams 核心:Source、Flow 与 Sink
- 3. 反压机制:推拉模型、缓冲与异步边界
- 4. fs2 流处理:纯函数式流与 Effect 融合
- 5. 错误处理与恢复:监督、重试与限流
- 6. 并发与并行流:mapAsync 与平衡
- 7. 外部系统桥接:Kafka、文件与数据库
- 8. 性能调优:缓冲、批量与背压水位
- 9. 流式测试与调试:TestKit、虚拟时钟与日志
- 10. 速查表与一句话记忆
- 延伸阅读
1. 流处理的心智模型:数据流、Stage 与背压
流处理把「数据」抽象为无限序列,把「处理」抽象为流水线上的 Stage:
Source(源头) → Flow(变换) → Flow(过滤) → Sink(出口)
三个核心概念:
① Source:数据源头(集合、文件、消息队列、定时器、无穷流)
② Flow:变换(map/filter/flatMap……纯数据变换)
③ Sink:出口(收集、打印、写库、聚合)
为什么需要背压:
- 生产者速度 >> 消费者速度 → 内存爆掉(无界缓冲)
- 反压 = 下游慢时,上游减速,而不是无脑堆积
- 三种策略:缓冲(限界)、丢弃、回传暂停
流 vs 集合 vs Spark DAG:
| 维度 | 集合 | 流(Akka/fs2) | Spark |
|---|---|---|---|
| 数据规模 | 有界 | 有界或无限 | 大且有界 |
| 求值 | 立即 | 惰性+运行时 | 惰性+执行计划 |
| 背压 | 无 | 有(核心) | 有(shuffle 缓冲) |
| 延迟 | 批 | 近实时 | 批 |
流的三不变式:无副作用在变换里、错误沿管道传播、慢则减速。
记忆:流 = 无限数据 + 流水线 Stage + 背压;Source 进、Flow 变、Sink 出,慢消费者让上游减速是流的核心约束。
2. Akka Streams 核心:Source、Flow 与 Sink
最小管道:
import akka.actor.ActorSystem
import akka.stream.scaladsl._
implicit val system: ActorSystem = ActorSystem("demo")
Source(1 to 10) // ① 源头:有限集合
.map(_ * 2) // ② 变换
.filter(_ % 4 == 0)
.runForeach(println) // ③ 出口:run 才真正执行
三个关键点:
① 惰性:构建 Source/Flow 只是蓝图,run 才触发执行
② 可复用:同一个流定义可 run 多次,每次独立
③ 可组合:.via(flow) 串联、Graph DSL 并联
Graph DSL(分叉/合并):
import akka.stream.{FlowShape, GraphDSL}
val graph = GraphDSL.create() { implicit b =>
import GraphDSL.Implicits._
val bcast = b.add(Broadcast[Int](2)) // 一进两出
val merge = b.add(Merge[String](2)) // 两进一出
Source(1 to 5) ~> bcast
bcast.out(0) ~> Flow[Int].map(n => s"a$n") ~> merge.in(0)
bcast.out(1) ~> Flow[Int].map(n => s"b$n") ~> merge.in(1)
FlowShape(bcast.in, merge.out)
}
常用 Source:
Source(List/Seq)、Source.fromIterator、Source.tick(定时)、
Source.repeat、Source.unfold(状态流)、Source.queue(手动喂)
记忆:Akka Streams 三件——Source 源头、Flow 变换、Sink 出口;惰性构建 + run 执行,Graph DSL 处理分叉合并;流蓝图可复用。
3. 反压机制:推拉模型、缓冲与异步边界
Reactive Streams 的契约:
- 下游通过 request(n) 向上游声明需求(推拉结合)
- 上游最多发 n 个,发完等下一批需求
- 全程只有 4 个信号:onSubscribe / onNext / onError / onComplete
背压的三种表现形式:
① 默认:逐级回传 —— 下游慢 → 上游停 → 源头减速
② 显式缓冲:.buffer(n, OverflowStrategy) → 下游满就丢/报错
③ 异步边界:async → 各段独立线程,之间以有界缓冲衔接
.buffer 策略:
Source(1 to 100000)
.buffer(100, OverflowStrategy.dropHead) // 满则丢最旧
.map(expensiveTransform)
async 异步边界:让重变换不阻塞源头:
Source(1 to 1000)
.map(cheapStep)
.async
.map(expensiveStep)
.runForeach(println)
背压调试信号:流量控制节流、无界缓冲导致 OOM、async 边界填满都提示背压位置。
记忆:背压契约是 request(n) 声明需求;默认逐级回传、buffer 限界丢弃、async 划异步边界;慢在哪段,哪段就该 buffer 或分流。
4. fs2 流处理:纯函数式流与 Effect 融合
fs2 与 Akka Streams 的根本差异:fs2 的 Stream 是纯函数式数据结构——构建不产生副作用,效果用 Effect(IO)表达,天然可组合、可测试。
最小示例:
import fs2.{Stream}
import cats.effect.{IO, IOApp}
object Demo extends IOApp.Simple {
def run: IO[Unit] =
Stream(1, 2, 3)
.map(_ * 2)
.covary[IO] // 从纯流升级为带 Effect 的流
.evalMap(n => IO.println(s"got $n"))
.compile
.drain
}
fs2 核心操作:
① 构造:Stream(1,2,3)、Stream.emits(seq)、Stream.iterate/range
② 变换:map/filter/flatMap/zipWith/scan
③ 效果:evalMap(每个元素带 IO)、eval(无元素副作用)
④ 收尾:compile.toList / .toVector / .drain / .fold
惰性与推拉:fs2 的流在 compile 前不执行任何效果;消费时按 pull 拉取。
fs2 与 IO 深度整合:
val process: Stream[IO, Unit] =
Stream.range(0, 100)
.evalMap(i => IO.sleep(10.millis) >> IO.println(i))
.interruptAfter(1.second) // 限时
fiber 并发内置:parEvalMap 并行处理元素(见第 6 节)。
记忆:fs2 的流是纯数据结构,效果交给 IO;构造-变换-evalMap-compile 五段式;compile 前无副作用、interruptAfter 限时。
5. 错误处理与恢复:监督、重试与限流
流里的错误不是异常而是数据流的一部分——必须显式决定「遇到错怎么办」。
Akka Streams 错误处理:
// 遇错即停(默认)
Source(1 to 10)
.map(n => if (n == 5) throw new RuntimeException("boom") else n)
.runForeach(println)
// 恢复:遇错替换为备用值,继续流
Source(1 to 10)
.map(n => if (n == 5) throw new RuntimeException() else n)
.recover { case _ => -1 }
.runForeach(println)
// 跳过出错元素
.recoverWithRetries(Int.MaxValue, { case _ => Source.empty })
重试与退避:
def call(n: Int) = // 可能失败的 Effect
Source(1 to 10)
.mapAsync(4)(call.retry) // 结合 retry 退避
fs2 错误处理:
Stream(1 to 10)
.map(n => if (n == 5) throw new RuntimeException() else n)
.handleErrorWith { _ => Stream(-1) } // 遇错改道
.attempt // 错误打包成 Either
限流(throttle):防压垮下游/外部 API:
Source(1 to 100)
.throttle(10, 1.second, 1, ThrottleMode.shaping)
记忆:流错误是数据不是异常——recover 换值、recoverWithRetries 跳过、handleErrorWith 改道、attempt 转 Either;外部调用加 retry 退避与 throttle 限流。
6. 并发与并行流:mapAsync 与平衡
串行太慢 → 并行处理。流里并发的核心是「以受控并发度把元素送进异步任务」。
Akka Streams mapAsync:
Source(1 to 100)
.mapAsync(parallelism = 8)(n => futureCall(n)) // 保持输出顺序
.runForeach(println)
// mapAsyncUnordered:不保序,吞吐更高
Source(1 to 100)
.mapAsyncUnordered(8)(n => futureCall(n))
fs2 parEvalMap:
Stream(1 to 100)
.parEvalMap(8)(n => futureCall(n)) // 并发度 8,保持顺序
.compile.toList
并发度怎么选:
- IO 密集(外部调用):并发度 = 下游容量 / 单任务耗时(经验 8~32)
- CPU 密集(计算):并发度 ≈ 核数
- 过高并发 → 打满线程池/连接池、背压失效
Balance 分流:多消费者平分负载:
val balance = GraphDSL.create() { implicit b =>
import GraphDSL.Implicits._
val bal = b.add(Balance[Int](3))
Source(1 to 100) ~> bal
(0 until 3).foreach(i => bal.out(i) ~> Flow[Int].map(n => s"w$i:$n").to(Sink.foreach(println)))
ClosedShape
}
记忆:mapAsync/parEvalMap 带并行度参数并行处理;IO 密集 8~32、CPU 密集核数;要保序用 mapAsync、拼吞吐用 Unordered;Balance 分流多消费者。
7. 外部系统桥接:Kafka、文件与数据库
流的价值在与外部系统双向打通。
Akka Streams + Kafka(Alpakka Kafka):
import akka.kafka.scaladsl._
val source = Consumer.plainSource(settings, subscription) // 消费
.map(_.value)
val producer = Producer.plainSink(producerSettings) // 生产
文件流:
// 读大文件(按行,惰性)
FileIO.fromPath(Paths.get("data.log"))
.via(Framing.delimiter(ByteString("\n"), 8192))
.map(_.utf8String)
.runForeach(println)
// 写文件
Source.single(ByteString("hello"))
.runWith(FileIO.toPath(Paths.get("out.txt")))
数据库流:Doobie/Slick 都提供流式 ResultSet:
// fs2 + Doobie:流式读表,不一次性载入内存
import doobie._
import fs2.Stream
val rows: Stream[ConnectionIO, Row] = sql"SELECT * FROM logs".query[Row].stream
Http4s 流式响应:把流直接作为 HTTP 响应体。
记忆:桥接三件——Kafka 用 Alpakka(Consumer.plainSource/Producer.plainSink)、文件用 FileIO+Framing、数据库用 Doobie 的 query.stream 流式读;Http4s 支持流响应体。
8. 性能调优:缓冲、批量与背压水位
流性能的两大矛盾:吞吐 vs 内存、延迟 vs 吞吐。
调优清单:
□ 批量化:.grouped(n) / .chunkN(n) 减少逐元素开销
□ 缓冲区大小:buffer(n) 匹配下游消费能力,别无界
□ 异步边界:heavy 变换 .async 隔离,避免阻塞源头
□ 并行度:mapAsync 按 IO/CPU 选并发度(见第 6 节)
□ 避免逐元素 IO:evalMap 里做小批量再写
□ 复用 Materializer/线程池:别每 run 建一个
Akka Streams 批量示例:
Source(1 to 10000)
.grouped(100) // 攒 100 个一批
.mapAsync(2)(batch => writeBatch(batch)) // 批量写库
背压水位观察:monitor/调试日志看 buffer 占用,满则调并发或扩缓冲。
fs2 chunk 优化:
Stream(1 to 10000)
.chunkN(100) // 按 chunk 批量处理
.evalMapChunk(batch => IO(writeBatch(batch)))
记忆:调优四板斧——grouped/chunkN 批量、buffer 定界、async 隔离重变换、mapAsync 匹配并发;背压水位满则降并发或扩缓冲,杜绝无界堆积。
9. 流式测试与调试:TestKit、虚拟时钟与日志
流测试最怕「跑起来才知道」——要可控、可复现、快。
Akka Streams TestKit:
import akka.stream.testkit.scaladsl._
val probe = TestSink.probe[Int](system)
Source(1 to 10).filter(_ % 2 == 0).runWith(probe)
probe
.request(2)
.expectNext(2, 4)
.expectComplete()
fs2 测试:纯流可直接断言 compile 结果:
val result = Stream(1 to 10).filter(_ % 2 == 0).compile.toList
assert(result == List(2, 4, 6, 8, 10))
虚拟时钟:定时器/流用 TestKit/cats-effect 的测试调度器推进时间,不用真等:
// cats-effect 测试调度器
val io = Stream.tick[IO](1.second).take(3).compile.toList
// 在测试里 advanceTime 快速推进
调试日志:
Source(1 to 5)
.map(_ * 2).log("afterMap").addAttributes(Attributes.logLevels(onElement = Logging.InfoLevel))
常见坑:
□ 忘了 run 而流不执行
□ 在测试里用真定时器拖慢 → 虚拟时钟
□ 无限流 + take 忘了限 → 永不结束
□ Materializer 泄漏 → 显式关闭
记忆:流测试三件——TestKit 的 probe 逐元素断言、fs2 纯流直接 compile 断言、虚拟时钟加速定时;调试用 log 级 Attributes;无限流必配 take 限界。
10. 速查表与一句话记忆
| 场景 | Akka Streams | fs2 |
|---|---|---|
| 构造流 | Source(seq) | Stream(seq) |
| 变换 | .map/.filter | .map/.filter |
| 带效果 | mapAsync | evalMap |
| 并行 | mapAsync(8) | parEvalMap(8) |
| 批量 | .grouped(100) | .chunkN(100) |
| 错误处理 | .recover / .recoverWithRetries | .handleErrorWith / .attempt |
| 限流 | .throttle | .metered |
| 停止 | run / runForeach | .compile.drain / toList |
一句话记忆:流处理 = Source 进、Flow 变、Sink 出,背压让慢下游回传减速;Akka Streams 用图 DSL 和 Materializer、fs2 用纯流 + Effect;错误是数据、并发用 mapAsync/parEvalMap 控并发度、批量用 grouped/chunkN、桥接用 Alpakka/FileIO/Doobie;测试用 TestKit 或直接 compile 断言——把「处理无限数据」从烧内存变成可控的流水线。
延伸阅读
- /scala-functional-effects/ — IO/Fiber 是 fs2 与流并发的底层
- /scala-bigdata-spark/ — 批处理 DAG 与流式管道的对比
- /scala-akka-cluster/ — 分布式流与 Actor 模型的配合
- /scala-database-access/ — Doobie 流式查询与事务
- /scala-testing-practice/ — 流式代码的测试体系
- [[distributed-systems]] — 消息队列与流式架构
- [[infra]] — Kafka 集群与流处理基础设施
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。