当内部发布者使用 subscribe(on:) 时,CombineLatest 运算符不会发出

Jos*_*reu 3 ios swift combinelatest combine

我观察到有关 JointLatest 的意外行为,如果内部发布者有subscribe(on:),则 JointLatest 流不会发出任何值。

笔记:

  • Zip 操作员正在工作
  • 将 subscribe(on:) / receive(on:) 移动到combineLatest 流也可以。但在这个特定的用例中,内部发布者正在定义他们的订阅/接收,因为在其他地方(重新)使用。
  • 仅将 subscribe(on:)/receive(on:) 添加到其中一个内部发布者也可以工作,因此问题就在于两者都拥有它时。
    func makePublisher() -> AnyPublisher<Int, Never> {
        Deferred {
            Future { promise in
                DispatchQueue.global(qos: .background).asyncAfter(deadline: .now() + 3) {
                    promise(.success(Int.random(in: 0...3)))
                }
            }
        }
        .subscribe(on: DispatchQueue.global())
        .receive(on: DispatchQueue.main)
        .eraseToAnyPublisher()
    }
    
    var cancellables = Set<AnyCancellable>()
    Publishers.CombineLatest(
        makePublisher(),
        makePublisher()
    )
    .sink { completion in
        print(completion)
    } receiveValue: { (a, b) in
        print(a, b)
    }.store(in: &cancellables)
Run Code Online (Sandbox Code Playgroud)

这是组合错误还是​​预期行为?您是否知道如何设置这种内部可以定义自己的订阅调度程序的流?

rob*_*off 5

是的,这是一个错误。我们可以将测试用例简化为:

import Combine
import Dispatch

let pub = Just("x")
    .subscribe(on: DispatchQueue.main)

let ticket = pub.combineLatest(pub)
    .sink(
        receiveCompletion: { print($0) },
        receiveValue: { print($0) })
Run Code Online (Sandbox Code Playgroud)

这永远不会打印任何东西。但如果您注释掉该subscribe(on:)运算符,它会打印出预期的内容。如果您离开subscribe(on:),但插入一些print()操作员,您将看到CombineLatest操作员永远不会向上游发送任何需求。

我建议您复制CombineX 重新实现CombineLatest以及它需要编译的实用程序(我认为是Lock和 的CombineX 实现LockedAtomic)。我也不知道CombineX版本是否有效,但如果它有问题,至少你有源代码并且可以尝试修复它。