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

Swift Combine:sink繁忙时收集值,空闲时立即处理且不丢值

Swift Combine实现忙时收集值、闲时立即处理的方案

问题背景

我们有一个快速发布新值的Publisher,但下游的slowProcessing处理速度较慢,导致值堆积。需要实现:

  • 处理忙时自动收集新值,合并为数组后批量处理
  • 处理空闲时立即执行单个值的处理
  • 不能丢失任何已发布的值
  • 无法使用buffer/collect(_:)(需立即启动处理)、debounce/throttle(会丢值)

现有代码:

private let snapshotQueue = DispatchQueue(
    label: "snapshotQueue",
)

QuickUpdating.$currentSnapshot
    .receive(on: snapshotQueue)
    .sink { [weak self] currentSnapshot in
        slowProcessing(currentSnapshot)
        updateUI()

    }
    .store(in: &cancellables)

解决方案

通过scan维护动态缓冲区,结合flatMap限制并发数为1,实现需求:

private let snapshotQueue = DispatchQueue(label: "snapshotQueue")

QuickUpdating.$currentSnapshot
    .receive(on: snapshotQueue)
    // 用scan累积待处理的值,处理完成后重置缓冲区
    .scan([]) { buffer, newValue in
        buffer + [newValue]
    }
    // 限制同时仅一个处理任务执行,确保忙时累积、闲时立即处理
    .flatMap(maxPublishers: .max(1)) { buffer -> AnyPublisher<[Snapshot], Never> in
        Future { promise in
            // 处理所有累积的值
            buffer.forEach { slowProcessing($0) }
            updateUI()
            // 返回空数组,让scan重置缓冲区
            promise(.success([]))
        }
        .eraseToAnyPublisher()
    }
    // 过滤空数组,避免无效回调
    .filter { !$0.isEmpty }
    .store(in: &cancellables)

工作逻辑

  1. scan:持续累积新发布的值到缓冲区数组。当flatMap返回空数组时,缓冲区会被重置为空,准备接收下一批值。
  2. flatMap(maxPublishers: .max(1)):确保同一时间只有一个处理任务在运行。如果当前任务未完成,新值会被scan持续累积;任务完成后,会立即处理所有累积的值。
  3. 无值丢失:所有发布的值都会被scan捕获并累积,直到被处理。
  4. 立即处理:当队列空闲时,新值会被scan包装为单元素数组,flatMap会立即启动处理,无延迟。

内容的提问来源于stack exchange,提问作者förschter

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 12:01:04