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

