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

如何在RxSwift中自定义Sink、Producer?及创建takeWhileInclusive算子的正确方式

在RxSwift中自定义事件处理逻辑及实现takeWhileInclusive算子

嘿,我来帮你理清这个问题!首先得明确一个核心点:RxSwift里的Sink和Producer是内部私有类,框架并没有把它们暴露给外部开发者,所以你没法直接导入使用。这是RxSwift的设计意图——鼓励咱们通过公开API来扩展功能,而不是依赖内部实现细节。

一、如何实现类似Sink/Producer的功能?

如果你想实现类似Producer(生成事件序列)和Sink(处理订阅事件逻辑)的角色,有两种规范的方式:

1. 自定义Observable(模拟Producer)

你可以直接继承Observable类,或者用Observable.create来创建自定义的事件生产者。比如下面这个简单的例子,模拟一个能发送固定元素序列的Observable:

class CustomElementObservable<Element>: Observable<Element> {
    private let elements: [Element]
    
    init(elements: [Element]) {
        self.elements = elements
        super.init()
    }
    
    override func subscribe<Observer: ObserverType>(_ observer: Observer) -> Disposable where Observer.Element == Element {
        // 这里的逻辑就类似Sink的职责:遍历元素并发送事件
        elements.forEach { observer.on(.next($0)) }
        observer.on(.completed)
        return Disposables.create()
    }
}

// 用法示例
CustomElementObservable(elements: [1,2,3])
    .subscribe(onNext: { print($0) })
    .disposed(by: DisposeBag())

2. 自定义算子(结合事件处理逻辑)

如果是想扩展RxSwift的算子(比如你要做的takeWhileInclusive),更推荐参考原生算子的实现思路:通过扩展ObservableType,结合内部的Producer和Sink类(虽然它们是内部的,但在自定义算子时,这种写法是RxSwift官方的标准范式,只要你在项目中正确引入RxSwift,就可以这么写)。


二、实现takeWhileInclusive算子

takeWhileInclusive的需求很明确:当元素满足指定条件时,包含该元素,然后立即终止序列(和原生takeWhile的区别是,takeWhile会在第一个不满足条件的元素处停止且不包含它,而咱们这个算子是在第一个满足条件的元素处停止且包含它)。

方式一:参考原生算子实现(使用内部Producer/Sink)

这种写法和RxSwift原生takeUntil的实现逻辑一致,是最贴近框架原生风格的方式:

extension ObservableType {
    func takeWhileInclusive(_ predicate: @escaping (Element) throws -> Bool) -> Observable<Element> {
        return TakeWhileInclusiveProducer(source: self.asObservable(), predicate: predicate)
    }
}

// 自定义Producer,负责创建Sink并订阅源序列
private final class TakeWhileInclusiveProducer<Element>: Producer<Element> {
    private let source: Observable<Element>
    private let predicate: (Element) throws -> Bool
    
    init(source: Observable<Element>, predicate: @escaping (Element) throws -> Bool) {
        self.source = source
        self.predicate = predicate
        super.init()
    }
    
    override func run<Observer: ObserverType>(_ observer: Observer, cancel: Cancelable) -> (sink: Disposable, subscription: Disposable) where Observer.Element == Element {
        let sink = TakeWhileInclusiveSink(predicate: predicate, observer: observer, cancel: cancel)
        let subscription = source.subscribe(sink)
        return (sink: sink, subscription: subscription)
    }
}

// 自定义Sink,负责处理事件转发和终止逻辑
private final class TakeWhileInclusiveSink<Observer: ObserverType>: Sink<Observer> {
    typealias Element = Observer.Element
    private let predicate: (Element) throws -> Bool
    
    init(predicate: @escaping (Element) throws -> Bool, observer: Observer, cancel: Cancelable) {
        self.predicate = predicate
        super.init(observer: observer, cancel: cancel)
    }
    
    private func handleEvent(_ event: Event<Element>) {
        switch event {
        case .next(let element):
            do {
                let shouldTerminate = try predicate(element)
                // 先发送当前元素
                forwardOn(.next(element))
                // 如果满足终止条件,发送完成事件并终止订阅
                if shouldTerminate {
                    forwardOn(.completed)
                    dispose()
                }
            } catch {
                forwardOn(.error(error))
                dispose()
            }
        case .error(let error):
            forwardOn(.error(error))
            dispose()
        case .completed:
            forwardOn(.completed)
            dispose()
        }
    }
}

// 让Sink符合ObserverType,这样才能订阅源Observable
extension TakeWhileInclusiveSink: ObserverType {
    func on(_ event: Event<Element>) {
        handleEvent(event)
    }
}

使用示例

Observable.of(1,2,3,4,5)
    .takeWhileInclusive { $0 == 3 }
    .subscribe(
        onNext: { print("收到元素:\($0)") },
        onCompleted: { print("序列完成") }
    )
    .disposed(by: DisposeBag())

// 输出:
// 收到元素:1
// 收到元素:2
// 收到元素:3
// 序列完成

方式二:完全使用公开API实现(不依赖内部类)

如果你不想碰框架的内部私有类,也可以用Observable.create来实现,代码更简洁直观:

extension ObservableType {
    func takeWhileInclusive(_ predicate: @escaping (Element) throws -> Bool) -> Observable<Element> {
        return Observable.create { observer in
            var shouldStop = false
            let disposable = self.subscribe(
                onNext: { element in
                    guard !shouldStop else { return }
                    do {
                        shouldStop = try predicate(element)
                        observer.on(.next(element))
                        if shouldStop {
                            observer.on(.completed)
                        }
                    } catch {
                        observer.on(.error(error))
                    }
                },
                onError: { observer.on(.error($0)) },
                onCompleted: { observer.on(.completed) }
            )
            return disposable
        }
    }
}

这种写法完全基于RxSwift的公开API,没有依赖任何内部实现,同样能满足需求,而且更容易理解和维护。


内容的提问来源于stack exchange,提问作者Barış Deniz Sağlam

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 09:17:30