如何在两个Flow间共享Combine Publisher避免重复请求
问题
我正在使用一个Publisher,当它输出值时需要执行两个基于Publisher的不同流程。如何共享第一个Publisher的结果来完成以下两个任务,且不重复调用请求?
相关定义
// 1. 获取凭证的Publisher func getCredentialsPublisher(userInfo: String)-> AnyPublisher<CredentialsResponse, Error> // 任务A - 创建套接字 // A1. 获取会话 func getSession(token: String) -> AnyPublisher<Session, Error> // A2. 设置套接字 func setupSockets(sessionID: String) -> SocketManager? // 任务B - 从其他服务请求展示信息 // B1. 获取外部对象 func getExternalObject(token: String) -> AnyPublisher<Object, Error> // B2. 获取视图信息 func getViewInfo(object: Object) -> AnyPublisher<ViewInfo, Error>
当前已实现的任务A流程
let pipeline = self.getCredentialsPublisher(userInfo: "howdy") .flatMap { credentialsResponse -> AnyPublisher<Session, Error> in // 从credentialsResponse中提取token let token = credentialsResponse.token // 替换示例值为实际提取逻辑 return self.getSession(token: token) } .sink(receiveCompletion: { completion in // 错误处理逻辑 switch completion { case .failure(let error): print("任务A出错:\(error)") case .finished: break } }, receiveValue: { session in // 从session中提取sessionID let sessionID = session.id // 替换示例值为实际提取逻辑 if let manager = self.setupSockets(sessionID: sessionID) { // 处理SocketManager的逻辑 } })
重复请求的问题
如果单独编写任务B的流程,会重复调用getCredentialsPublisher:
let pipeline2 = self.getCredentialsPublisher(userInfo: "howdy") .flatMap { credentialsResponse -> AnyPublisher<Object, Error> in let token = credentialsResponse.token // 替换示例值为实际提取逻辑 return getExternalObject(token: token) } .flatMap { object -> AnyPublisher<ViewInfo, Error> in return getViewInfo(object: object) } .sink(receiveCompletion: { completion in // 错误处理逻辑 switch completion { case .failure(let error): print("任务B出错:\(error)") case .finished: break } }, receiveValue: { viewInfo in // 处理viewInfo的逻辑 })
解决方案
要共享第一个Publisher的结果且避免重复请求,核心是用share()操作符让多个订阅者共享同一个上游Publisher的输出,同时可以根据需求选择独立订阅或合并流程。
方法一:独立订阅(推荐,任务逻辑分离)
先创建共享的凭证Publisher,再分别为两个任务搭建流程:
// 1. 创建共享的凭证Publisher,确保只发起一次请求 let sharedCredentials = self.getCredentialsPublisher(userInfo: "howdy") .share() .eraseToAnyPublisher() // 2. 任务A流程 let taskAPipeline = sharedCredentials .flatMap { credentials in self.getSession(token: credentials.token) } .sink(receiveCompletion: { completion in if case .failure(let error) = completion { print("任务A失败:\(error)") } }, receiveValue: { session in if let manager = self.setupSockets(sessionID: session.id) { // 处理SocketManager的业务逻辑 } }) // 3. 任务B流程 let taskBPipeline = sharedCredentials .flatMap { credentials in getExternalObject(token: credentials.token) } .flatMap { object in getViewInfo(object: object) } .sink(receiveCompletion: { completion in if case .failure(let error) = completion { print("任务B失败:\(error)") } }, receiveValue: { viewInfo in // 处理viewInfo的业务逻辑 })
方法二:合并流程(统一处理结果)
如果希望两个任务并行执行,且统一处理完成状态,可以用Publishers.Zip合并两个任务的Publisher:
let combinedPipeline = self.getCredentialsPublisher(userInfo: "howdy") .flatMap { credentials -> AnyPublisher<(Session?, ViewInfo?), Error> in // 任务A的流程:捕获错误避免单个任务失败导致整体终止 let taskA = self.getSession(token: credentials.token) .map { $0 as Session? } .catch { _ in Just(nil) } // 任务B的流程:同样捕获错误 let taskB = getExternalObject(token: credentials.token) .flatMap { getViewInfo(object: $0) } .map { $0 as ViewInfo? } .catch { _ in Just(nil) } // 并行执行两个任务,合并结果 return Publishers.Zip(taskA, taskB) .eraseToAnyPublisher() } .sink(receiveCompletion: { completion in if case .failure(let error) = completion { print("整体流程出错:\(error)") } }, receiveValue: { session, viewInfo in // 处理任务A结果 if let session = session, let manager = self.setupSockets(sessionID: session.id) { // 处理SocketManager逻辑 } // 处理任务B结果 if let viewInfo = viewInfo { // 处理viewInfo逻辑 } })
关键说明
share():核心操作符,确保上游的getCredentialsPublisher只会被订阅一次,所有下游订阅者共享其输出,彻底避免重复请求。- 流程选择:如果两个任务的错误处理、业务逻辑相对独立,优先用方法一;如果需要统一监控两个任务的完成状态,再考虑方法二。
内容的提问来源于stack exchange,提问作者Maria
相关产品推荐
相关产品推荐

