【问题标题】:RxSwift: share() alternative that guarantees single subscription on the upstreamRxSwift:share() 替代方案,保证上游的单一订阅
【发布时间】:2022-07-28 05:16:46
【问题描述】:

我一直认为.share(replay: 1, scope: .forever) 共享一个上游订阅,无论下游有多少订阅者。

但是,我刚刚发现,如果下游订阅的计数降至零,它会停止“共享”并释放上游的订阅(因为在后台使用了 refCount())。所以当一个新的下游订阅发生时,它必须在上游重新订阅。在以下示例中:

let sut = Observable<Int>
    .create { promise in
        print("create")
        promise.onNext(0)
        return Disposables.create()
    }
    .share(replay: 1, scope: .forever)

sut.subscribe().dispose()
sut.subscribe().dispose()

我希望create 只打印一次,但它会打印两次。如果我删除 .dispose() 调用 - 只需一次。

如何设置上游保证最多订阅一次的链?

【问题讨论】:

  • 对我来说看起来像个 bug。可以建议使用deferred 并返回just,就像在.forever 范围的cmets 中一样。
  • promise.onCompleted() 修复了输出。可能是连接到replay: 1:当没有输出且流未完成时,则没有什么可重播。
  • 好吧,我不能在我的代码中使用onCompleted(),因为我正在create 块中开始一个数据库更改观察,这个流没有“完成”
  • >没什么可重播的。

标签: swift rx-swift


【解决方案1】:

您描述的目标意味着您应该使用multicast(或使用它的运算符之一,例如publish()replay(_:)replayAll())而不是share...

let sut = Observable<Int>
    .create { observer in
        print("create")
        observer.onNext(0)
        return Disposables.create()
    }
    .replay(1)

let disposable = sut.connect() // subscription will stay alive until dispose() is called on this disposable...

sut.debug("one").subscribe().dispose()
sut.debug("two").subscribe().dispose()

要了解 .forever 和 .whileConnected 之间的区别,请阅读“ShareReplayScope.swift”文件中的文档。两者都被重新计算,但不同之处在于重新订阅运算符的处理方式。这是一些显示差异的测试代码...

class SandboxTests: XCTestCase {
    var scheduler: TestScheduler!
    var observable: Observable<String>!

    override func setUp() {
        super.setUp()
        scheduler = TestScheduler(initialClock: 0)
        // creates an observable that will error on the first subscription, then call `.onNext("A")` on the second.
        observable = scheduler.createObservable(timeline: "-#-A")
    }

    func testWhileConnected() {
        // this shows that re-subscription gets through the while connected share to the source observable
        let result = scheduler.start { [observable] in
            observable!
                .share(scope: .whileConnected)
                .retry(2)
        }
        XCTAssertEqual(result.events, [
            .next(202, "A")
        ])
    }

    func testForever() {
        // however re-subscription doesn't get through on a forever share
        let result = scheduler.start { [observable] in
            observable!
                .share(scope: .forever)
                .retry(2)
        }
        XCTAssertEqual(result.events, [
            .error(201, NSError(domain: "Test Domain", code: -1))
        ])
    }
}

【讨论】:

    【解决方案2】:

    我不知道为什么.share(replay: 1, scope: .forever) 没有给出你想要的行为(我也认为它应该像你描述的那样工作)但是没有share 的其他方式呢?

    // You will subscribe to this and not directly on sut (maybe hiding Subject interface to avoid onNext calls from observers)
    let subject = ReplaySubject<Int>.create(bufferSize: 1)
    
    let sut = Observable<Int>.create { obs in
      print("Performing work ...")
      obs.onNext(0)
      return Disposables.create()
    }
    
    // This subscription is hidden, happens only once and stays alive forever
    sut.subscribe(subject)
    
    // Observers subscribe to the public stream
    subject.subscribe().dispose()
    subject.subscribe().dispose()
    

    【讨论】:

    • 我现在得到了一个类似的代码,我故意“泄露”订阅。我不喜欢这个,这就是为什么希望有一种更清洁的方法
    【解决方案3】:

    我不喜欢建议的解决方案中泄漏的Disposable,因此提出了以下建议:

    extension ObservableType {
        func shareReplayForever() -> Observable<Element> {
            let relay = BehaviorRelay<Element?>(value: nil)
            let disposeBag = DisposeBag()
            var subscribeOnce: () -> Void = {
                self.bind(to: relay).disposed(by: disposeBag)
            }
            return relay
                .compactMap { $0 }
                .do(onSubscribe: {
                    subscribeOnce()
                    subscribeOnce = { }
                    _ = disposeBag
                })
        }
    }
    

    它通过在闭包中保留disposeBag 将其生命周期绑定到下游的生命周期来处理绑定到中继后留下的disposable。 释放的下游取消了上游的订阅,否则这将是永恒的。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-01-30
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-08-08
      • 1970-01-01
      • 2015-06-13
      相关资源
      最近更新 更多