【问题标题】:Swift Combine: Unexpected backpressure behaviour with zip operatorSwift Combine:使用 zip 运算符的意外背压行为
【发布时间】:2020-01-10 06:35:44
【问题描述】:

我有一个关于结合背压的组合中的 zip 运算符的问题。

取如下代码sn-p:

let sequencePublisher = Publishers.Sequence<Range<Int>, Never>(sequence: 0..<Int.max)
let subject = PassthroughSubject<String, Never>()

let handle = subject
    .zip(sequencePublisher.print())
    .print()
    .sink { letters, digits in
        print(letters, digits)
    }

subject.send("a")

在操场上执行此操作时,输出如下:

receive subscription: (0..<9223372036854775807)
receive subscription: (Zip)
request unlimited
request unlimited
receive value: (0)
receive value: (1)
receive value: (2)
receive value: (3)
receive value: (4)
receive value: (5)
receive value: (6)
receive value: (7)
...

在 iOS 设备上执行时,由于内存问题,代码在几秒钟后崩溃。

根本原因可以在上面的第四行中看到,其中zipsequencePublisher 请求无限数量的值。由于sequencePublisher 提供了整个Int 值范围,这会导致内存溢出。

我想知道的:

  • zip 等待每个发布者的一个值,然后将它们合并并推送它们
  • 背压用于控制从订阅者到发布者的需求

我的期望是zip 只向每个发布者请求一个值,等待它们到达,并且仅在从每个发布者收到一个值时才请求下一个值。

在这种特殊情况下,我尝试构建一种行为,其中将序列号分配给subject 生成的每个值。但是,我可以想象,当zip 结合来自发布频率非常不同的发布者的值时,这总是一个问题。

zip 运算符中利用背压似乎是解决该问题的完美工具。你知道为什么不是这样吗?这是一个错误还是故意的?如果是故意的,为什么?

谢谢大家

【问题讨论】:

    标签: ios swift reactive-programming combine


    【解决方案1】:

    Sequence 发布者似乎不切实际。它似乎对背压没有反应;它只是一次喷出整个序列,这在发布应该是异步的世界中毫无意义。如果您将Int.max 更改为 3,则没有问题。 :) 我不知道这是一个错误还是仅仅是序列发布者整个概念中的一个缺陷。

    但是,对于您的实际用例来说确实没有问题,因为有一种更好的方法可以为来自主题的每个发射分配一个连续的数字,即scan

    这是一个更现实的方法:

    func delay(_ delay:Double, closure:@escaping ()->()) {
        let when = DispatchTime.now() + delay
        DispatchQueue.main.asyncAfter(deadline: when, execute: closure)
    }
    class ViewController : UIViewController {
        var storage = Set<AnyCancellable>()
        override func viewDidLoad() {
            super.viewDidLoad()
            let subject = PassthroughSubject<String, Never>()
            subject.scan(("",0)) {t,s in (s,t.1+1)}
                .sink { print($0.0, $0.1)
                }.store(in:&storage)
            delay(1) {
                subject.send("a") // a 1
                delay(1) {
                    subject.send("b") // b 2
                }
            }
        }
    }
    

    这假设您有其他原因需要每个连续的枚举通过管道传递。但是,如果您的唯一目标是在每个信号到达 sink 本身时对其进行枚举,那么您可以让 sink 本身维护一个计数器(它可以轻松做到,因为它是关闭):

        var storage = Set<AnyCancellable>()
        let subject = PassthroughSubject<String, Never>()
        override func viewDidLoad() {
            super.viewDidLoad()
            var counter = 1
            subject
                .sink {print($0, counter); counter += 1}
                .store(in:&storage)
            delay(1) {
                self.subject.send("a") // a 1
                self.subject.send("b") // b 2
            }
        }
    
    

    【讨论】:

    • 酷,感谢您的出色回复。我对函数式编程很陌生,因此了解此类食谱非常有用。
    • 其实我在这里使用.scan 可能有点矫枉过正;如果您只需要.sink 中的计数器,则只需将.sink 本身配置为计数即可。我已将这个想法添加到我的答案中。
    【解决方案2】:

    Combine 的 zip 运算符:

    1. 将请求的背压需求从其下游订阅者转发到上游,这对于接收器是无限的
    2. 从第一个上游缓冲整个序列

    除了基于扫描的解决方案外,您还可以通过控制下游的背压需求或使用自定义 zip 运算符来避免该问题。

    我设法开发了一个自定义 zip 运算符,它在两个方面与原来的不同:

    1. 无论它从下游收到什么背压需求,它总是向上游发送一个只需要一个值的需求,然后等待每个上游的响应,然后将结果发送到下游,结束这一轮。如此重复,直到需求耗尽。
    2. 它通过使用所描述的基于“轮”的方法避免缓冲整个上游序列。

    这里包含的代码有点广泛,但请随时在此 repo https://github.com/SergeBouts/XCombine 中查看它

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2022-11-13
      • 1970-01-01
      • 2017-03-05
      • 1970-01-01
      • 2016-08-10
      • 1970-01-01
      相关资源
      最近更新 更多