Scala 流处理实战:Akka Streams、fs2 与反压机制

Scala 流处理三件套:Akka Streams(Source/Flow/Sink 图 DSL)、fs2(纯函数式流与 Effect 融合)、反压机制(推拉模型、缓冲与异步边界),覆盖错误恢复、并发流、Kafka 桥接、性能调优与流式测试。

引言

「流处理」不是大数据专属——日志采集、消息队列消费、文件管道、实时聚合,处处都在处理无限数据序列。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 与背压

流处理把「数据」抽象为无限序列,把「处理」抽象为流水线上的 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 Streamsfs2
构造流Source(seq)Stream(seq)
变换.map/.filter.map/.filter
带效果mapAsyncevalMap
并行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 集群与流处理基础设施

继续阅读

探索更多技术文章

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

全部文章 返回首页

「scala」更多文章

  1. 纯函数式效果系统实战:Cats Effect IO 与 ZIO
  2. Scala.js 与 Scala Native:跨平台编译、互操作与工程实践
  3. Scala 领域建模实战:ADT、类型驱动设计与模块化架构