如果你在 Swift 并发里用 AsyncStream 把回调事件包装成异步序列,大概率会用到 continuation.yield 来往流里塞数据。一个常见的疑问是:当消费者的 for await 循环处理得特别慢,生产者还在快速 yield,这时候生产者会被自动挂起等待吗?答案是不一定,而且大多数情况下并不会。yield 本身是同步方法,它只负责和内部缓冲区打交道,不会等待消费者把元素取走。要理解这一点,必须先看 yield 的返回值和 AsyncStream 的缓冲策略。

yield 的返回值与三种缓冲策略
创建 AsyncStream 时,makeStream 有个 bufferingPolicy 参数,它决定了内部缓冲区的大小和溢出行为。可选值有三个:unbounded、bufferingNewest 和 bufferingOldest。其中 unbounded 表示缓冲区没有上限,有多少元素就存多少;bufferingNewest 会保留最新的若干元素,缓冲区满时丢弃最旧的那个;bufferingOldest 则相反,缓冲区满时丢弃新来的元素。
与这些策略直接相关的是 continuation.yield 的返回值 YieldResult。它有三个枚举 case:enqueued 表示元素成功进入缓冲区;dropped 表示元素没有被保存下来;terminated 表示流已经结束,再 yield 也没有意义。不同缓冲策略下,同样的快速生产慢速消费会得到完全不同的 YieldResult。
下面这段代码使用 .bufferingOldest(2),缓冲区只能留两个元素,消费者每 100 毫秒才消费一个,而生产者每 10 毫秒 yield 一个。因为缓冲区很快被占满,后面新元素会直接被丢弃,yield 返回 dropped。
let (stream, continuation) = AsyncStream.makeStream(of: Int.self, bufferingPolicy: .bufferingOldest(2))
Task {
for i in 0..<10 {
let result = continuation.yield(i)
print("yield \(i): \(result)")
try? await Task.sleep(nanoseconds: 10_000_000)
}
continuation.finish()
}
Task {
for await value in stream {
try? await Task.sleep(nanoseconds: 100_000_000)
print("consume \(value)")
}
}
运行这个例子会发现,生产者打印的 result 里会频繁出现 dropped。这清楚说明 yield 没有因为消费者慢而阻塞生产者,它只是按照策略把元素入队或丢弃,然后立即返回。如果把缓冲策略改成 .bufferingNewest(2),生产者看到的几乎永远是 enqueued,因为旧元素被静默丢掉给新元素腾地方,生产者甚至不知道有数据丢失。这个差异在选择策略时需要特别留意。
为什么 yield 不能天然充当背压机制
从 API 签名就能看出端倪:yield 不是 async 方法,调用它不能用 await,也就无法把生产者任务挂起。它的职责范围非常小,就是向 AsyncStream 内部缓冲区投递一个元素,然后报告这次投递的结果。真正从缓冲区取数据的是消费者的 for await 循环,而生产者和消费者之间没有直接的执行流依赖。
如果使用 unbounded 策略,生产者不会遇到任何阻碍,元素会不断堆积在内部缓冲区里。消费者处理得越慢,缓冲区就膨胀得越大。在长时间运行的服务端程序或高频事件源中,这可能会造成内存持续上涨,最终触发内存压力甚至崩溃。所以当数据量不可控时,无限缓冲通常不是一个稳妥的选择。
有人会想,既然 .bufferingOldest 能让 yield 返回 dropped,那生产者不就能知道消费者慢了吗?这个判断只对了一半。知道 dropped 确实是一种信号,但它发生在元素已经被丢弃之后,严格来说属于事后通知,而不像背压那样在发送前就把速度降下来。生产者如果连续看到多次 dropped,可以选择主动休眠或降低生产频率,但这种做法是外部补偿,不是 AsyncStream 内建的背压能力。
与其他响应式框架对比会更清楚。例如 Combine 中有 Subscribers.Demand 来表示订阅者愿意接收多少元素,RxSwift 也有类似的背压策略。Swift 的 AsyncStream 设计目标主要是把既有同步回调接口桥接到 AsyncSequence,它本身没有为慢消费者提供完整的流量控制方案。
实现真正背压的几种可行方案
如果业务上不允许丢数据,又要求生产者在消费者跟不上时自动暂停,就需要在 AsyncStream 之外增加流量控制。第一种思路是用一个 actor 充当许可计数器。生产者每次 yield 前必须先获取许可,许可数量代表缓冲区还能容纳的未消费元素数;消费者每处理完一个元素,就释放一个许可。这样生产者会在许可不足时 await,从而真正停下来等待。
下面是一个简化版实现。生产者先通过 gate.acquire() 拿到许可,再执行 yield;消费者处理完后调用 gate.release() 归还许可。许可上限控制在一个比较小的值,例如 4,那么缓冲区里未被消费的元素最多不会超过 4 个。
actor BackpressureGate {
private let limit: Int
private var current = 0
init(limit: Int) {
self.limit = limit
}
func acquire() async {
while current >= limit {
await Task.yield()
}
current += 1
}
func release() {
current -= 1
}
}
let (stream, continuation) = AsyncStream.makeStream(of: Int.self, bufferingPolicy: .bufferingOldest(64))
let gate = BackpressureGate(limit: 4)
Task {
for i in 0..<20 {
await gate.acquire()
let result = continuation.yield(i)
print("yield \(i): \(result)")
}
continuation.finish()
}
Task {
for await value in stream {
try? await Task.sleep(nanoseconds: 50_000_000)
print("consume \(value)")
await gate.release()
}
}
不过这个示例里的 acquire 使用了忙等循环,会持续占用 actor 执行器,并不适合生产环境。更完善的实现可以借助 withCheckedContinuation 把等待者挂起,在 release 时恢复它们。原理是维护一个等待队列,每次释放许可时唤醒队首的任务。这样生产者不会空转,也不会浪费 CPU 资源。
第二种思路是直接使用 Apple 的 swift-async-algorithms 包提供的 AsyncChannel。与 AsyncStream 不同,AsyncChannel 的 send 方法是异步的:当缓冲区满时,发送者会被挂起,直到消费者取出元素腾出空间。这正好就是我们需要的内建背压能力。
import AsyncAlgorithms
let channel = AsyncChannel<Int>()
Task {
for i in 0..<10 {
await channel.send(i)
print("sent \(i)")
}
channel.finish()
}
Task {
for await value in channel {
try? await Task.sleep(nanoseconds: 100_000_000)
print("received \(value)")
}
}
使用 AsyncChannel 时,生产者的 await channel.send(i) 会在缓冲区达到上限时暂停,这样从根源上避免了无限堆积。它适合对数据完整性要求高、不能随意丢弃元素的场景。代价是需要引入额外的 Swift 包,并且在某些简单的回调桥接场景中会显得稍重。
第三种思路是折中方案:继续使用 AsyncStream 的有限缓冲,比如 .bufferingOldest(32),但生产者统计 yield 返回 dropped 的次数。如果连续多次丢弃,就执行 Task.sleep 主动降速;如果丢弃减少或消失,再逐步恢复生产速度。这种方式实现简单,不引入新依赖,但无法保证零丢失,适合允许丢失部分数据但希望尽量避免长时间过载的日志采集、指标上报等场景。
在实际工程中如何选择,取决于数据的重要程度和延迟容忍度。如果每个元素都必须被处理,那就不应该依赖 yield 的返回值来事后补救,而应该选择 AsyncChannel 或自建带许可队列的方案。如果可以接受溢出时丢弃一部分数据,那么 .bufferingNewest 适合只关心最新状态的场景,.bufferingOldest 适合更看重先到先处理的场景。无论哪种情况,都要避免把 unbounded 用在高频且消费者速度不可控的路径上,否则内存问题很可能比丢失几个元素更难排查。
AsyncStream背压Swift并发修改时间:2026-09-29 09:31:02