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. 响应式流规范与背压语义
- 2. Source、Flow 与 Sink:核心算子
- 3. 图 DSL:Fan-out、Fan-in 与连接
- 4. 背压策略:buffer、throttle 与溢出策略
- 5. 流式错误处理:recover、restart 与监督
- 6. 与外部系统集成:Kafka、数据库与网络
- 7. 流式应用实战:ETL 与事件处理
- 8. 性能与资源:materialization 与异步边界
- 9. 与函数式流对比:选择与配合
- 10. 速查表与一句话记忆
- 延伸阅读
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 专题 — 消息队列与流平台
- 分布式系统专题 — 事件驱动架构
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。