【问题标题】:Combine - merging multiple shared filters合并 - 合并多个共享过滤器
【发布时间】:2021-12-24 19:24:16
【问题描述】:

我使用 RxSwift 已经有一段时间了,刚刚切换到 Combine,我正试图围绕这个特定的 .filter 行为。这是一个简短的游乐场示例:

import Combine

let publisher = [1, 2, 3, 4, 5]
    .publisher
    .share()

let filter1 = publisher
    .filter { $0 == 1 }
    .print("filter1")

let filter2 = publisher
    .filter { $0 == 2 }
    .print("filter2")

Publishers
    .Merge(filter1, filter2)
    .sink {
        print("Result is: \($0)")
    }

输出是

filter1: receive subscription: (Multicast)
filter1: request unlimited
filter1: receive value: (1)
Result is: 1
filter1: receive finished
filter2: receive subscription: (Multicast)
filter2: request unlimited
filter2: receive finished

令我惊讶的是,Result is: 2 从未被调用,因为流结束了。我可以删除 .share() 运算符,这将导致接收到我期望的两个值

filter1: receive subscription: ([1])
filter1: request unlimited
filter1: receive value: (1)
Result is: 1
filter1: receive finished
filter2: receive subscription: ([2])
filter2: request unlimited
filter2: receive value: (2)
Result is: 2
filter2: receive finished

但是如果我的发布者是一个 API 调用并且我不想创建重复的网络请求怎么办?这正是我现在要处理的情况,这也是我需要使用.share() 运算符的原因。

有什么更好的解释为什么会发生这种情况以及如何处理您想要过滤流、在每个流中执行单独的逻辑然后将结果重新合并在一起的情况?

【问题讨论】:

    标签: ios asynchronous rx-swift combine


    【解决方案1】:

    所以这里发生了一些不同的事情。

    首先[1, 2, 3].publisher 的工作方式与 Observable.from([1, 2, 3]) 不同。后者每个周期发出一次值,而前者则背靠背发出所有值。 Publisher 示例在 Rx 中的工作方式更像这样:

    Observable<Int>.create { observer in
        [1, 2, 3, 4, 5].forEach {
            observer.onNext($0)
        }
        observer.onCompleted()
        return Disposables.create()
    }
    

    因此,在 Observable.from 的情况下,在订阅 filter2 可观察对象时,排放完成。因此,即使您省略了share(),“Result is: 1”和“Result is: 2”都会被发出。

    第二share() 运算符的工作方式也不同。默认情况下,一旦所有订阅都处理完毕,RxSwift 共享操作符将重置 Observable(这是一个引用计数共享)。在组合案例中,共享运算符使发布者可连接,然后连接到它。本质上,它与 RxSwift 中的 .share(replay: 0, scope: .forever) 运算符相同(这是我在 Rx BTW 中从未需要的)。

    所以与你发布的Combine代码等效的Rx代码实际上是这样的:

    let observable = emitSequence([1, 2, 3, 4, 5])
        .share(replay: 0, scope: .forever)
    
    let filter1ʹ = observable
        .filter { $0 == 1 }
        .debug("filterʹ1")
    
    let filter2ʹ = observable
        .filter { $0 == 2 }
        .debug("filterʹ2")
    
    Observable.merge(filter1ʹ, filter2ʹ)
        .subscribe(onNext: {
            print("Resultʹ is: \($0)")
        })
    
    func emitSequence<S>(_ sequence: S) -> Observable<S.Element> where S: Sequence {
        Observable.create { observer in
            sequence.forEach {
                observer.onNext($0)
            }
            observer.onCompleted()
            return Disposables.create()
        }
    }
    

    所有这些都表明处理 API 调用的实际方面很好。在这种情况下,假设调用不会立即返回(至少需要一个周期)并且因为它是一次性的,只要您确保没有重新订阅 Observable,@ 987654329@不重置不是问题。

    【讨论】:

    • 有些人可能将Combine 中发生的事情称为错误,其他人可能将其称为哲学差异。就个人而言,我认为使用经过 6 年以上实战测试的 API 与使用只有几年历史的 API 是不同的。
    猜你喜欢
    • 2010-10-29
    • 1970-01-01
    • 1970-01-01
    • 2022-09-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多