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

如何在两个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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 10:40:58