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
相关产品推荐
相关产品推荐

