背压
安全地处理高速生产者。
背压 是 CoddyKit 上的免费 Scala for Backend Engineering & Functional Programming 课时。 这是第 3 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Scala for Backend Engineering & Functional Programming 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Scala for Backend Engineering & Functional Programming 课程共包含 4 节课。
什么是背压
背压是一种流控制机制,可以防止快速生产者压垮缓慢的消费者。消费者不会让缓冲区无限增长或直接丢弃数据,而是会发出信号,说明自己能够处理多少数据。
Akka Streams 实现了 Reactive Streams 标准,其中需求向上游传递,元素向下游传递。
需求驱动的流
只有在下一个阶段发出需求信号后,每个阶段才会发出元素。Sink 请求 N 个元素;该需求会一直向上游传播,直到 Source 正好产生所请求的数量。
这种基于拉取的协议意味着生产者不会推送超出消费者处理能力的数据。
背压为何对流水线很重要
如果没有背压,一个快速的 Kafka 消费者向缓慢的数据库供给数据时,可能会积累数百万条正在处理的记录,耗尽内存并导致进程崩溃。
背压会自然地将上游速度限制在最慢阶段的速度,从而在负载下保持稳定的内存使用量。
// Fast source, slow sink: backpressure slows the source
val g =
Source(1 to 1000000)
.map(_ * 2)
.to(slowDatabaseSink)内部缓冲
在异步边界之间,Akka Streams 会维护一个小型内部缓冲区(默认包含 16 个元素)。它可以吸收短暂的突发流量,因此各阶段不必在每个元素上严格同步。
缓冲区填满后,背压就会生效,上游会停止产生数据,直到缓冲区腾出空间。
import akka.stream.Attributes
val buffered =
Flow[Int]
.map(identity)
.addAttributes(Attributes.inputBuffer(initial = 32, max = 32))使用溢出策略的显式缓冲
buffer 运算符会插入一个指定大小的显式缓冲区,并使用 OverflowStrategy 决定缓冲区已满时的处理方式。
这样,您可以在内存占用和解耦生产者与消费者速度的能力之间进行权衡。
import akka.stream.OverflowStrategy
val withBuffer =
Source(1 to 1000)
.buffer(size = 100, OverflowStrategy.backpressure)溢出策略
可用策略包括 backpressure(减慢上游)、dropHead/dropTail(丢弃最旧或最新的元素)、dropBuffer、dropNew,以及 fail(以错误终止)。
丢弃策略适合传感器读数等实时数据,因为过时的值可以安全地丢弃。
import akka.stream.OverflowStrategy
val latestWins =
liveTicks.buffer(1, OverflowStrategy.dropHead)
val strict =
liveTicks.buffer(50, OverflowStrategy.fail)使用 conflate 进行汇总
当消费者速度较慢时,conflate 会使用组合函数将待处理元素合并为一个元素,而不是将它们全部缓冲起来。
例如,可以将许多数值更新合并为它们的总和,使消费者始终看到所错过数据的聚合结果。
val summarized =
fastMetrics
.conflate((acc, next) => acc + next)
// Slow downstream receives summed batches使用 expand 填满足需求
expand 与 conflate 相反:当下游的需求速度快于上游的产生速度时,它会根据最近看到的值合成额外元素。
这适合以稳定速率持续发出最近一次读数。
val repeated =
sensor.expand(last => Iterator.continually(last))
// Downstream always gets the latest sensor value异步边界
默认情况下,融合的阶段会在同一个 actor 上运行,彼此之间没有缓冲。插入 async 后,某个阶段会运行在独立的 actor 上,从而增加缓冲并支持流水线式并行。
异步边界就是背压缓冲区实际存在的位置。
val pipelined =
Source(1 to 1000)
.map(slowStep).async
.map(anotherSlowStep).async
.to(Sink.ignore)使用 throttle 进行显式速率控制
throttle 会施加一个明确的最大速率,并向上游产生背压以遵守该速率。即使消费者能够处理得更快,它也能保护受到速率限制的外部服务。
突发参数允许短时间内超过稳定速率。
import scala.concurrent.duration._
val limited =
requests
.throttle(
elements = 100, per = 1.second, maximumBurst = 20,
akka.stream.ThrottleMode.Shaping)观察背压
您可以通过观察上游是否变慢,或测量缓冲区占用情况来检测背压。log 运算符和 Akka 的流属性有助于追踪流水线发生停滞的位置。
持续处于满载状态的缓冲区表明,最慢阶段正在限制吞吐量。
val traced =
Source(1 to 100)
.log("after-source")
.map(_ * 2)
.log("after-map")
.to(Sink.ignore)快速检查
思考 Akka Streams 如何防止快速生产者淹没缓慢的消费者。
回顾
背压是 Akka Streams 由需求驱动的骨干机制:消费者向上游发出需求信号,因此生产者无法压垮消费者,同时将内存使用量控制在有限范围内。
您了解了内部缓冲、带溢出策略的显式 buffer 运算符、使用 conflate 进行汇总、使用 expand 填满足需求、异步边界,以及通过 throttle 进行有意的速率控制。下一步,您将运行一条完整的流水线。
常见问题解答
「背压」课时是免费的吗?
是的 — 「背压」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Scala for Backend Engineering & Functional Programming 课程的其余内容,请升级到 CoddyKit PRO。 Scala for Backend Engineering & Functional Programming 课程共包含 4 节课。
「背压」这节课中我会学到什么?
安全地处理高速生产者。 你通过在浏览器中直接运行的动手代码来练习 Scala for Backend Engineering & Functional Programming,全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 Scala for Backend Engineering & Functional Programming 需要有经验吗?
无需任何先前经验。CoddyKit 上的 Scala for Backend Engineering & Functional Programming 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 3 节课,共 4 节。
「背压」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 Scala for Backend Engineering & Functional Programming 课中编写并运行代码吗?
能。每节 Scala for Backend Engineering & Functional Programming 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。