Akka Streams 与响应式流:图 DSL、背压与流式实战

深入 Akka Streams:Reactive Streams 规范与背压协议(request 信号/Publisher/Subscriber)、Source/Flow/Sink 算子、图 DSL(Fan-out/Fan-in/连接器)、背压策略(buffer/throttle/溢出策略)、监督与重启(recover/restart)、Kafka/数据库/网络集成、ETL 与事件处理实战、materialization 与异步边界、以及与函数式流(fs2)的对比选型。

Akka Streams 是 Reactive Streams 规范在 JVM 上的黄金实现,它的精髓不只是「Source→Flow→Sink 连起来」,而是图 DSL 与背压协议——元素按需流动、慢者不被快者淹没、图可以分叉汇合还能优雅失败。本文在前置篇(基础算子与反压)之上深入实现层:Reactive Streams 的 request/onNext 协议、Fan-out/Fan-in 连接器、buffer/throttle/溢出策略、监督重启语义、materialization 的「运行句柄」、异步边界,以及 Kafka/数据库/网络集成的完整实战。最后与 fs2 对比,帮你在「图式流」与「函数式流」之间选型。

前置:/scala-stream-processing/(流处理与反压基础)、/scala-akka-cluster/(Akka 生态与分布式)、/actor-model-detailed-explanation/(Actor 模型)、/scala-concurrency-atomics/(并发底层)。

目录

1. 响应式流规范与背压语义

Reactive Streams 是一个「四个接口 + 一条协议」的标准,Akka Streams 是其完整实现:

四个接口:
□ Publisher[T]   :可被订阅的数据源
□ Subscription  :连接两者的「管道」(背压控制端)
背压协议(本质是「信用系统」):
□ 下游调用 request(n):声明「我能再收 n 个」
□ 上游按信用放行 onNext,最多放 n 个
□ 无 request 则上游不推送 → 内存不会被冲爆
关键规则:
□ request 是可累计的信用额度(replenish)
□ 非法信号:错误协议实现直接失败(fail fast)
一次背压走位(文字时序):
Subscriber: request(2)
Publisher : onNext(a)  → 剩余额度 1
Publisher : onNext(b)  → 剩余额度 0(停住,等再 request)
Subscriber: request(5)
Publisher : onNext(c) ...

工程要点:背压的本质是**「下游用 request 发放信用,上游按信用放行」**——不是「缓冲区多大」的问题,而是「生产者无信用就停产」的协议保证。理解这一点,就不会用「无限缓冲」糊弄背压:正确做法是让慢路径通过 request 语义自然拖慢快路径,内存占用有上限。

2. Source、Flow 与 Sink:核心算子

Source/Flow/Sink 是「图 DSL」的三个基础形状,它们的区别在输入输出口数量:

形状(Shape):
□ Source[+Out, +Mat] :一个输出口(上游)
□ Flow[-In, +Out, +Mat] :一个输入 + 一个输出(变换段)
□ Sink[-In, +Mat]:一个输入口(下游终点)
常用构造:
□ Flow.map/filter/scan/grouped/sliding/prepend
□ Sink.foreach/seq/ignore/fold/headOption
Mat(materialized value):运行后产生的「句柄」
□ Source.queue(...).preMaterialize → (queue, source)
□ Sink.foreach 的 mat 是 Future[Done]
□ 用 toMat/keep 选择保留哪一侧句柄
import akka.actor.ActorSystem
import akka.stream.scaladsl.*
import akka.Done
import scala.concurrent.Future
given system: ActorSystem = ActorSystem("demo")
// 一段管道的三种写法(等价)
val p1: Future[Done] =
  Source(1 to 5).map(_ * 2).filter(_ % 3 != 0).runForeach(println)
val p2: Future[Done] =
  Source(1 to 5).via(Flow[Int].map(_ * 2).filter(_ % 3 != 0)).runWith(Sink.ignore)
// 保留 mat:拿到队列句柄,向流中动态推送
val (queue, source) = Source.queue[Int](8, OverflowStrategy.backpressure).preMaterialize()
queue.offer(42) // 返回 Future[QueueOfferResult]

工程要点:三个形状的核心是**「口的数量决定角色」**——Source 只有输出、Sink 只有输入、Flow 是中间变换;而 Mat 是「运行后你能拿到的句柄」(队列、Future、取消信号),是 Akka Streams「可编程入口」的关键。先记形状,再记算子,最后记 Mat——三件事构成了使用 Akka Streams 的全部骨架。

3. 图 DSL:Fan-out、Fan-in 与连接

单线管道不够,Akka Streams 用 RunnableGraph 图 DSL 表达分叉、汇合与环路:

连接器(Junctions):
Fan-out(一分多):
□ Balance[T]      :元素轮流/按空闲分发给下游(负载均衡)
□ Partition[T]    :按谓词把元素分到指定分支
Fan-in(多合一):
□ Merge[T]        :合并多个输入流(任意顺序)
□ Zip[L, R]       :一对一配对(最慢者决定节奏)
□ Concat          :按顺序衔接多个流(一个结束才开始下一个)
组合语法:
□ 用 GraphDSL.Implicits 的 ~> 连接
□ 入口/出口用 SourceShape/SinkShape 暴露
选型直觉:
□ 广播复制 → Broadcast;分发干活 → Balance
□ 合并异步事件 → Merge;对齐多源 → Zip
□ 分组路由 → Partition
import akka.stream.scaladsl.*
import akka.stream.*
// 图:偶数走 A 分支,奇数走 B 分支,再汇合
val g: RunnableGraph[Future[Done]] = RunnableGraph.fromGraph(GraphDSL.create() {
  implicit b =>
    import GraphDSL.Implicits.*
    val src   = b.add(Source(1 to 10))
    val part  = b.add(Partition[Int](2, _ % 2 == 0 ? 0 | 1))
    val merge = b.add(Merge[Int](2))
    val sink  = b.add(Sink.foreach(println))
    src  ~> part
    part.out(0) ~> merge.in(0)      // 偶数分支
    part.out(1) ~> merge.in(1)      // 奇数分支
    merge ~> sink
    ClosedShape
})
g.run()

工程要点:图 DSL 的准则是**「分叉用 Broadcast/Balance/Partition,汇合用 Merge/Zip/Concat」**——分叉时想清楚「复制 vs 分配 vs 路由」;汇合时想清楚「合并 vs 对齐」。任何分叉汇合都要回到背压语义:Broadcast 会等最慢分支,Balance 按空闲分发,Zip 由最慢输入决定节奏。环路必须加缓冲,否则信用协议会死锁。

4. 背压策略:buffer、throttle 与溢出策略

背压的「默认行为」是传导(慢者拖慢上游),但工程上常要局部缓冲或主动限速:

buffer + OverflowStrategy(溢出策略):
□ backpressure :满了挂起上游(最安全,全传导)
□ dropHead      :满则丢最旧(滑动窗口,保新鲜)
□ dropTail      :满则丢最新(保历史)
□ fail          :满了让流失败(背压即崩溃信号)
throttle(主动限速):
□ throttle(elements, per, mode)
  - Shaping  :均摊放行(平滑输出速率)
  - Enforcing:超出就 fail 或丢弃(严格执行速率)
组合策略:
□ 慢下游但不想丢 → 无界 buffer(内存换吞吐,谨慎)
□ 可丢的实时数据 → dropHead(永远消费最新)
□ 外部 API 限速  → throttle + 重试
import akka.stream.OverflowStrategy
import scala.concurrent.duration.*
// 实时监控:保最新,满则丢最旧(可丢数据)
val fast = Source.fromIterator(() => Iterator.from(1))
  .buffer(100, OverflowStrategy.dropHead)
// 主动限速:每秒最多 5 个,平滑输出
val limited = Source(1 to 100)
  .throttle(5, 1.second, ThrottleMode.Shaping)

工程要点:背压策略的准则是**「关键数据用 backpressure、实时数据用 dropHead、外部限速用 throttle」**——先问「能不能丢」再选策略:不能丢就让它挂起传导,能丢就选保新鲜的丢弃策略。throttle(Shaping) 是「平滑速率」,Enforcing 是「严格执行」,前者用于保护下游、后者用于契约限速。

5. 流式错误处理:recover、restart 与监督

流的失败不是「整个应用崩掉」,而是在图的特定位置被处理或重启:

错误处理层次:
□ recover/recoverWith:把失败替换成兜底元素/流
□ mapError:把异常翻译成领域错误(向下游继续推)
□ 监督(Supervision):算子级失败策略
  - ActorAttributes.supervisionStrategy(decider)
  - decider: (throwable) => Resume / Restart / Stop
    Resume  :跳过坏元素继续(元素级)
    Restart :重启该算子,丢状态但继续流
    Stop    :终止整条流(默认)
□ restartSource/restartFlow:整个子图重启(连接中断自愈)
import akka.stream.Supervision
import akka.stream.scaladsl.*
// 元素级跳过坏数据:坏元素丢弃,流继续
val resilient = Source(1 to 100)
  .map { x =>
    if (x % 10 == 0) throw new RuntimeException("bad")
    x * 2
  }
  .withAttributes(ActorAttributes.supervisionStrategy { _ => Supervision.Resume })
// 连接中断自动重连:restartSource 按退避重启
val reconnecting = RestartSource.onFailuresWithBackoff(
  minBackoff = 1.second, maxBackoff = 30.seconds, randomFactor = 0.2
) { () => Source.queue[Int](8, OverflowStrategy.backpressure) }

工程要点:流式错误的准则是**「元素级用 Resume、子图级用 Restart、全局用 Stop + 报警」**——数据管道里「一个坏元素不该杀掉整条流」,用监督策略跳过;外部连接(Kafka/WebSocket)用 RestartSource.onFailuresWithBackoff 自愈。recover 处理「预期内失败」,监督处理「算子级意外」,两者配合让流「活着且正确」。

6. 与外部系统集成:Kafka、数据库与网络

Akka Streams 的价值一半在外部系统桥接——且桥接也要讲背压:

Kafka(Alpakka Kafka):
□ Source.committerSink 消费:自动提交 offset(背压敏感)
□ 用 committable 消息保证「处理完才提交」
□ 暂停/恢复:consumer 默认按下游需求拉取(天然背压)
数据库:
□ Doobie/Slick 配合:把查询变成 Source/Flow
网络:
□ Source.tcp / Flow.tcp:原始 TCP 流式收发
□ WebSocket:Source/Flow/Sink 三合一(协议本身双向)
□ HTTP(Akka HTTP):响应体即 Source[ByteString]
桥接原则:
□ 用 mapAsync 控制对下游的并发冲击(限流)
□ 外部「拉取型」天然配合背压;「推送型」要自己缓冲
import akka.stream.scaladsl.*
import akka.stream.alpakka.kafka.scaladsl.*
import akka.stream.alpakka.kafka.{CommitterSettings, ConsumerMessage, Subscriptions}
// 消费 Kafka:处理完才提交 offset(不丢消息)
def kafkaPipeline: Source[ConsumerMessage.CommittableMessage[String, String], _] =
  CommittableSource(consumerSettings, Subscriptions.topics("orders"))
    .map { msg =>
      process(msg.record.value())
      msg
    }
    .mapAsync(4)(msg => msg.committableOffset.commitScaladsl())

工程要点:外部集成的准则是**「拉取型天然背压、推送型自己缓冲、写外部要 mapAsync 限流」**——Kafka 消费用 committable + 提交机制保证「处理完才确认」;DB 写用 mapAsync(n) 限制并发批次;网络流把字节当元素流。任何桥接都回到同一句话:别让外部系统的节奏打穿你的流。

7. 流式应用实战:ETL 与事件处理

把以上拼成一个真实 ETL/事件处理管线:

事件处理管线(订单事件 → 聚合 → 落库):
① Source:Kafka 订单事件(committable)
② 清洗:filter 非法事件、map 领域对象、回填维表(mapAsync)
③ 聚合:grouped / groupBy + 窗口聚合(按订单号、按分钟)
④ 转换:金额换算、状态机流转(纯函数段)
⑤ Sink:批量写 DB + 失败分流到死信队列
工程要点:
□ 每个段可独立测试(Source 用内存、Sink 用收集)
□ 错误隔离:清洗失败 → 死信队列;聚合失败 → 重试
批/流一体视角:
□ ETL 在流里的形态:Source(文件/DB) → map → Sink
□ 与 Spark 差异:Akka 是「逐条流式」,Spark 是「微批」
import akka.stream.scaladsl.*
// 订单事件聚合:按订单号分组 → 求和 → 写库
val orderTotals = Source
  .fromIterator(() => rawEvents.iterator)
  .filter(_.isValid)
  .groupBy(1024, e => e.orderId)
  .fold((e: Event) => e.copy(amount = e.amount + e.amount))
  .mergeSubstreams
  .mapAsync(8)(saveToDB)
  .to(Sink.ignore)

工程要点:流式应用实战的准则是**「清洗→变换→聚合→落库 每段可测、失败分流、监控在位」**——把管线切成纯函数段与 IO 段,纯段可离线测,IO 段配重试与死信。groupBy + fold 是流式聚合的经典组合,mergeSubstreams 把分组流收回主线。监控背压水位与错误率,让「慢在哪」随时可见。

8. 性能与资源:materialization 与异步边界

运行流的性能要害是 materialization 与异步边界(async):

materialization:
□ run() 的返回值就是 mat(Future/队列/取消信号)
□ killSwitch:手动终止流的句柄(SharedKillSwitch)
异步边界(async):
□ 默认一个 Actor 跑整条流(顺序执行)
□ .async:把一段放到独立 Actor,引入并行与缓冲
□ 分叉汇合(Broadcast/Merge)天然跨 Actor
□ 目的:CPU 并行 + 隔离慢段
性能要点:
□ 控制 Actor 数量:每个 async 一个 Actor,太多则调度开销
□ 批量化:grouped(n) 减少每元素开销
□ 缓冲适度:太大延迟高、太小抖动大
import akka.stream.scaladsl.*
// async 边界:让重计算段与 IO 段并行
val pipelined = Source(1 to 1000)
  .map(heavyCompute).async     // 段 A:独立 Actor 跑 CPU
  .mapAsync(4)(networkCall)    // 段 B:IO 并行
  .runForeach(println)
// killSwitch:可随时干净终止
import akka.stream.KillSwitches
val (kill, done) = infiniteSource
  .viaMat(KillSwitches.single)(Keep.right)
  .toMat(Sink.ignore)(Keep.both)
  .run()
kill.shutdown()

工程要点:性能调优的准则是**「按段并行(async)、按批降开销、用 killSwitch 掌控生命周期」**——默认单 Actor 顺序流延迟最低;需要吞吐就把 CPU 段与 IO 段用 .async 分隔并行。Actor 数量与缓冲大小都是权衡:太多 Actor 调度开销、太大缓冲推高延迟。生产上先量测瓶颈段再加 async,别盲目并行。

9. 与函数式流对比:选择与配合

Akka Streams(图式)与 fs2(函数式)是 Scala 流的两大流派,选型看需求:

Akka Streams:
□ Reactive Streams 标准实现(可与任何 RS 库互操作)
□ 图 DSL:Fan-out/Fan-in/环路表达力强
□ Actor 后台:与 Akka 生态(集群/持久化)天然一体
□ 更适合:复杂拓扑、Akka 生态、Java/Scala 团队
fs2(函数式流):
□ 纯函数式,嵌合 Cats Effect(IO)
□ 流就是「可组合的惰性效果」,无 Actor
□ 更适合:与 CE/ZIO 效果栈统一、函数式纯度优先
对比维度:
□ 互操作:Akka 是 RS 标准,fs2 也有 RS 桥(但非核心)
□ 生命周期:Akka 图 + Actor;fs2 纯 F 组合
□ 背压:两者都是推送/信用,fs2 在 Effect 层
□ 迁移成本:Akka 存量用 Akka,新纯函数栈用 fs2
实际项目:
□ 大型事件平台(Akka 集群)→ Akka Streams
□ 效果栈(CE/ZIO)内做流 → fs2
选型速判:
□ 要图拓扑(分叉汇合/环路)+ Akka 生态 → Akka Streams
□ 要纯函数 + 与 Cats Effect 统一 → fs2

工程要点:选型的准则是**「看拓扑与生态」**——复杂分叉汇合与 Akka 生态选 Akka Streams,纯函数栈(CE/ZIO)内选 fs2。两者背压语义等价、可互操作,不必纠结「谁更正统」。实际项目常以一方为主,另一方只做边界桥接(比如 fs2 处理纯计算、Akka 接 Kafka 与集群)。

10. 速查表与一句话记忆

问题一句话答案
背压是什么下游 request 信用、上游按信用放行
Source/Flow/Sink输出口 / 变换段 / 输入口
Mat 是什么运行后拿到的句柄(队列/Future/killSwitch)
分叉怎么选Broadcast 复制、Balance 分配、Partition 路由
汇合怎么选Merge 合并、Zip 对齐、Concat 衔接
溢出策略关键 backpressure、实时 dropHead
限速throttle(Shaping 平滑 / Enforcing 严格)
错误怎么办元素 Resume、子图 Restart、全局 Stop+报警
外部桥接拉取天然背压、推送自缓冲、写用 mapAsync 限流
性能要害按段 async 并行、批量化、killSwitch 管控

一句话记忆:Akka Streams = 响应式流协议(request 信用制背压)+ 图 DSL(Fan-out/Fan-in 分叉汇合)+ 溢出策略(backpressure/dropHead/throttle)+ 监督重启(Resume/Restart/Stop + 退避自愈)+ 外部桥接(Kafka committable/mapAsync 限流)+ materialization(句柄 + async 并行 + killSwitch)——核心心法:让信用制背压管住内存,让监督策略管住失败,让图 DSL 管住拓扑。

延伸阅读

  • /scala-stream-processing/ — 流处理基础与反压机制
  • /scala-akka-cluster/ — Akka 分布式与集群
  • /actor-model-detailed-explanation/ — Actor 模型深入
  • /scala-bigdata-spark/ — 批处理与微批对比
  • /scala-observability-logging-tracing/ — 流式监控与追踪
  • /scala-functional-effects/ — 效果系统与 fs2 的基座
  • Kafka 专题 — 消息队列与流平台
  • 分布式系统专题 — 事件驱动架构

继续阅读

探索更多技术文章

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

全部文章 返回首页

「scala」更多文章

  1. Scala Native 与 GraalVM:AOT 编译、互操作与部署
  2. 函数式架构:六边形设计、纯核心与副作用外壳
  3. Tagless Final 与代数式设计:类型类、DSL 与多解释器