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

Swift合并多Publisher/AsyncStream为单一数据流的实现问题

问题

数据库设置了两个监听,数据更新时触发:

func listenToAltersForUser(userId: String) async -> PassthroughSubject<[AlterModel], Never> { ... }

func listenForUpdates() -> PassthroughSubject<Bool, Never> { ... }

需要将这两个发布者合并为AnyPublisher<[AlterUIModel], Never>类型的发布者。

当前Combine实现存在两个问题:异步处理逻辑错误,且只有一个发布者能触发更新:

func invoke(userId: String, lastAlterId: String?) async -> AnyPublisher<[AlterUIModel], Never> {
        let subject = PassthroughSubject<[AlterUIModel], Never>()
        
        let alterStream = await alterRepository.listenToAltersForUser(userId: userId)
        let frontRecordUpdates = frontRepository.listenForUpdates()
        let altersPublisher = alterStream.eraseToAnyPublisher()
        let frontRecordsPublisher = frontRecordUpdates.eraseToAnyPublisher()
        
        return altersPublisher.combineLatest(frontRecordsPublisher) { alters, frontUpdates in
            Task.detached {
                let altersUIModels = await alters.asyncMap { alter in
                    let frontRecord = await self.frontRepository.getLastFrontRecordForAlter(alterId: alter.id)
                    let isFronting = frontRecord != nil && frontRecord?.endTime == nil
                    return AlterUIModel(
                        fromAlterModel: alter,
                        isFronting: isFronting,
                        frontingDate: frontRecord?.endTime ?? frontRecord?.startTime
                    )
                }
            }
        }.eraseToAnyPublisher()
    } 

切换到AsyncStream后出现类型错误:

func invoke(userId: String, lastAlterId: String?) async {
        AsyncStream([AlterUIModel].self) { cont in
            let alterStream = await alterRepository.listenToAltersForUser(userId: userId)
            let frontRecordStream = frontRepository.listenForUpdates()
            for await (alters, _) in combineLatest(alterStream, frontRecordStream) {
                let alters = await alters.asyncMap { alter in
                    let frontRecord = await frontRepository.getLastFrontRecordForAlter(alterId: alter.id)
                    let isFronting = frontRecord != nil && frontRecord?.endTime == nil
                    return AlterUIModel(
                        fromAlterModel: alter,
                        isFronting: isFronting,
                        frontingDate: frontRecord?.endTime ?? frontRecord?.startTime
                    )
                }
                return alters
            }
        }
    } 

错误信息:

Cannot pass function of type '(AsyncStream<[AlterUIModel]>.Continuation) async -> Void' to parameter expecting synchronous function type

解决方案

方案1:用Combine实现正确的异步合并

Combine的combineLatest闭包要求返回同步值,不能直接返回Task。需要通过flatMap拆分异步逻辑,用Future包装异步操作,最终合并为目标数组:

func invoke(userId: String, lastAlterId: String?) async -> AnyPublisher<[AlterUIModel], Never> {
    let alterStream = await alterRepository.listenToAltersForUser(userId: userId)
    let frontRecordUpdates = frontRepository.listenForUpdates()
    
    return alterStream
        .combineLatest(frontRecordUpdates)
        // 将AlterModel数组拆分为单个元素的流,逐个处理
        .flatMap { alters, _ in
            Publishers.Sequence(sequence: alters)
        }
        // 用Future包装异步查询,转换为Combine Publisher
        .flatMap { alter -> AnyPublisher<AlterUIModel, Never> in
            Future { promise in
                Task {
                    let frontRecord = await self.frontRepository.getLastFrontRecordForAlter(alterId: alter.id)
                    let isFronting = frontRecord != nil && frontRecord?.endTime == nil
                    let uiModel = AlterUIModel(
                        fromAlterModel: alter,
                        isFronting: isFronting,
                        frontingDate: frontRecord?.endTime ?? frontRecord?.startTime
                    )
                    promise(.success(uiModel))
                }
            }
            .eraseToAnyPublisher()
        }
        // 将单个UI模型收集为数组,匹配输出类型
        .collect()
        .eraseToAnyPublisher()
}

方案2:用AsyncStream实现正确的异步流合并

AsyncStream初始化闭包是同步的,需将异步逻辑放到独立Task中执行,同时要先把PassthroughSubject转换为AsyncStream:

func invoke(userId: String, lastAlterId: String?) -> AsyncStream<[AlterUIModel]> {
    AsyncStream { continuation in
        // 启动异步Task处理流合并逻辑
        Task {
            // 将PassthroughSubject转换为AsyncStream
            let alterAsyncStream = AsyncStream<[AlterModel]> { cont in
                let alterStream = await alterRepository.listenToAltersForUser(userId: userId)
                let cancellable = alterStream.sink { alters in
                    cont.yield(alters)
                }
                cont.onTermination = { _ in
                    cancellable.cancel()
                }
            }
            
            let frontAsyncStream = AsyncStream<Bool> { cont in
                let frontStream = frontRepository.listenForUpdates()
                let cancellable = frontStream.sink { update in
                    cont.yield(update)
                }
                cont.onTermination = { _ in
                    cancellable.cancel()
                }
            }
            
            // 合并两个AsyncStream
            for await (alters, _) in combineLatest(alterAsyncStream, frontAsyncStream) {
                let uiModels = await alters.asyncMap { alter in
                    let frontRecord = await self.frontRepository.getLastFrontRecordForAlter(alterId: alter.id)
                    let isFronting = frontRecord != nil && frontRecord?.endTime == nil
                    return AlterUIModel(
                        fromAlterModel: alter,
                        isFronting: isFronting,
                        frontingDate: frontRecord?.endTime ?? frontRecord?.startTime
                    )
                }
                continuation.yield(uiModels)
            }
            
            continuation.finish()
        }
    }
}

// 辅助函数:合并两个AsyncStream的最新值
func combineLatest<T, U>(_ stream1: AsyncStream<T>, _ stream2: AsyncStream<U>) -> AsyncStream<(T, U)> {
    AsyncStream { continuation in
        var latestT: T?
        var latestU: U?
        var isStream1Finished = false
        var isStream2Finished = false
        
        let task1 = Task {
            for await value in stream1 {
                latestT = value
                if let u = latestU {
                    continuation.yield((value, u))
                }
            }
            isStream1Finished = true
            if isStream2Finished {
                continuation.finish()
            }
        }
        
        let task2 = Task {
            for await value in stream2 {
                latestU = value
                if let t = latestT {
                    continuation.yield((t, value))
                }
            }
            isStream2Finished = true
            if isStream1Finished {
                continuation.finish()
            }
        }
        
        continuation.onTermination = { _ in
            task1.cancel()
            task2.cancel()
        }
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 08:10:43