【问题标题】:Avoiding overlapping OnNext calls in Rx when using SubscribeOn(Scheduler.TaskPool)使用 SubscribeOn(Scheduler.TaskPool) 时避免在 Rx 中重叠 OnNext 调用
【发布时间】:2012-02-06 12:51:00
【问题描述】:

我有一些使用 Rx 的代码,从多个线程调用:

subject.OnNext(value); // where subject is Subject<T>

我希望在后台处理值,所以我的订阅是

subscription = subject.ObserveOn(Scheduler.TaskPool).Subscribe(value =>
{
    // use value
});

我并不真正关心哪些线程处理来自 Observable 的值,只要将工作放入 TaskPool 并且不阻塞当前线程即可。但是,我在 OnNext 委托中使用“值”不是线程安全的。目前,如果有很多值正在通过 Observable,我会收到对 OnNext 处理程序的重叠调用。

我可以为我的 OnNext 委托添加一个锁,但这不像 Rx 的做事方式。当我有多个线程调用 subject.OnNext(value); 时,确保一次只调用一次 OnNext 处理程序的最佳方法是什么?

【问题讨论】:

    标签: c# system.reactive


    【解决方案1】:

    来自 MSDN 上的Using Subjects

    默认情况下,主题不执行任何同步 线程。 [...] 但是,如果您想 使用调度程序同步对观察者的传出调用,您可以使用 Synchronize 方法来执行此操作。

    因此,正如 Brandon 在 cmets 中所说,您应该同步主题并将其交给您的生产者线程。例如

    var syncSubject = Subject.Synchronize(subject);
    
    // syncSubject.OnNext(value) can be used from multiple threads
    
    subscription = syncSubject.ObserveOn(TaskPoolScheduler.Default).Subscribe(value =>
    {
        // use value
    });
    

    【讨论】:

    • 调用Synchronize 应该在调用ObserveOn 之前,否则你违反了ObserveOn 的并发契约。但实际上,如果用例在多个线程之间共享一个主题,最好的解决方案是同步主题而不是同步主题的订阅者:var syncSubject = Subject.Synchronize(syncSubject); 现在将syncSubject 交给你的生产者线程,他们可以调用@ 987654328@ 不会对原始主题的订阅者造成问题。
    • @Brandon var syncSubject = Subject.Synchronize(syncSubject); ...您在分配之前使用了一个变量,您能澄清一下吗?
    • @Beachwalker 这是一个错字。它应该是 subject 作为参数传递给 Synchronize,如答案所示
    【解决方案2】:

    我认为您正在寻找 .Synchronize() 扩展方法。为了在最近的版本(2011 年末)中获得性能改进,他们 Rx 团队放宽了关于可观察序列生成器的顺序性质的假设。但是,您似乎打破了这些假设(不是一件坏事),但要让 Rx 像用户期望的那样重新播放,您应该同步序列以确保它再次是连续的。

    【讨论】:

      【解决方案3】:

      这里再解释一下为什么要使用Synchronize(第二段)。 另一方面,如果您在代码中积极使用锁定,则同步可能会参与死锁,至少我目睹了这种情况。

      【讨论】:

        【解决方案4】:

        您可以尝试实现自己的 IObserver,它将在其 OnNext 方法中应用锁定。只是一个简单的装饰器:应用锁,调用内部 OnNext,移除锁。

        然后你可以在 IObservable 上实现一个扩展方法,比如 .AsThreadSafe()。

        【讨论】:

        • 不是最好的解决方案,但从我的角度来看不需要被否决,因为它是一个有效的解决方案。所以我 +1 来平衡一下这里的反对意见。
        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2018-02-12
        相关资源
        最近更新 更多