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

如何在Combine中创建可发送多个值的Publisher并封装多次触发的回调

Combine 封装多次调用回调为Publisher的实现方案

直接上两种常用实现方式,不需要自定义复杂Publisher即可满足需求:


方案1:全系统版本兼容实现

兼容所有支持Combine的系统版本(iOS13+/macOS10.15+),和你熟悉的RxSwift Observable.create、ReactiveSwift SignalProducer语义完全一致:

import Combine

func webSocketMessagePublisher() -> AnyPublisher<String, Error> {
    Deferred {
        // 用于转发事件的Subject,不需要缓存历史值用PassthroughSubject即可
        let subject = PassthroughSubject<String, Error>()
        
        // 注册原有的消息接收回调
        webSocketClient.receiveMessage { message, error in
            // 错误场景发送失败结束事件
            if let error = error {
                subject.send(completion: .failure(error))
                return
            }
            // 收到消息向下游发送值
            if let message = message {
                subject.send(message)
            }
        }
        
        return subject
            .handleEvents(receiveCancel: {
                // 订阅取消时执行资源清理,比如取消回调注册、断开连接等
                // webSocketClient.stopReceivingMessage()
            })
    }
    .eraseToAnyPublisher()
}

逻辑说明:

  • Deferred 确保每次有新订阅者时才执行内部的回调注册逻辑,避免提前触发资源加载
  • PassthroughSubject 作为事件中转,接收回调的参数转化为Publisher事件发给下游
  • handleEvents 的receiveCancel闭包处理订阅销毁时的资源清理,避免内存泄漏

方案2:iOS15+/macOS12+ 简化实现

如果你的项目不需要兼容旧系统,可以用AsyncThrowingStream桥接实现,代码更简洁不需要手动管理Subject:

@available(iOS 15.0, macOS 12.0, *)
func webSocketMessagePublisher() -> AnyPublisher<String, Error> {
    AsyncThrowingStream<String, Error> { continuation in
        // 注册回调
        webSocketClient.receiveMessage { message, error in
            if let error = error {
                continuation.finish(throwing: error)
                return
            }
            if let message = message {
                continuation.yield(message)
            }
        }
        // 订阅取消时的清理逻辑
        continuation.onTermination = { @Sendable _ in
            // webSocketClient.stopReceivingMessage()
        }
    }
    .eraseToAnyPublisher()
}

注意事项

  • 如果回调闭包会被原对象(比如webSocketClient)强持有,注意在闭包内避免循环引用,必要时添加[weak ...]修饰符
  • 你提到的TCA的Effect.run本质上就是对上述逻辑的通用封装,普通场景不需要引入TCA即可自己实现,没有额外依赖

内容的提问来源于stack exchange,提问作者Yoshikuni Kato

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 02:18:00