You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PassthroughSubject的AsyncPublisher values属性未输出全部值问题排查

问题原因

你的问题出在getValues()方法的实现上:你用Task包裹for await in subject.values来接收事件,而PassthroughSubject本身没有缓冲机制。当你连续同步发送三个事件时,Task里的异步循环无法及时跟上事件发送的速度,后续的两个事件因为没有订阅者处于就绪状态接收,直接被丢弃了。

解决方案

替换Task+for await的实现方式,改用PassthroughSubject的sink方法来订阅事件,这样能同步捕获每个发送的事件并传递给AsyncStream的continuation,而AsyncStream会自动缓冲这些事件,确保for await循环能逐个处理所有值。

修改后的Channel类代码:

protocol Channelable {}

final class Channel {
    private let subject = PassthroughSubject<Channelable, Never>()
    
    func send(_ value: Channelable) {
        subject.send(value)
    }
    
    func getValues() -> AsyncStream<Channelable> {
        AsyncStream { continuation in
            // 使用sink同步接收事件
            let cancellable = subject.sink { value in
                continuation.yield(value)
            }
            
            continuation.onTermination = { _ in
                cancellable.cancel()
            }
        }
    }
}
补充说明
  • PassthroughSubject是热信号,没有内置缓冲,事件发送时如果没有订阅者就绪接收就会丢失;
  • AsyncStream本身具备缓冲能力,通过sink将事件同步传递给continuation,能确保所有事件被缓冲并等待for await循环处理;
  • 原实现中Task的异步循环会导致事件接收的时序滞后,无法匹配同步发送的多事件场景。

内容的提问来源于stack exchange,提问作者hydro1337x

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.27 19:32:56