【问题标题】:How do I share an observable with publish and connect?如何通过发布和连接共享可观察对象?
【发布时间】:2015-08-27 22:17:29
【问题描述】:

我有一个可观察的数据流,我正在对其应用操作,分成两个单独的流,对两个流中的每一个应用更多(不同的)操作,然后再次合并在一起。我正在尝试使用 PublishConnect 在两个订阅者之间共享 observable,但每个订阅者似乎都在使用单独的流。也就是说,在下面的示例中,我看到 为两个订阅者 为流中的每个项目打印了一次“执行昂贵的操作”。 (想象一下昂贵的操作应该在所有订阅者之间只发生一次,因此我试图重用流。)我使用PublishConnect 尝试与两个订阅者共享合并的 observable,但是好像效果不对。

问题示例:

var foregroundScheduler = new NewThreadScheduler(ts => new Thread(ts) { IsBackground = false });
var timer = Observable.Timer(TimeSpan.Zero, TimeSpan.FromSeconds(10), foregroundScheduler);
var expensive = timer.Select(i =>
{
    // Converting to strings is an expensive operation
    Console.WriteLine("Doing an expensive operation");
    return string.Format("#{0}", i);
});

var a = expensive.Where(s => int.Parse(s.Substring(1)) % 2 == 0).Select(s => new { Source = "A", Value = s });
var b = expensive.Where(s => int.Parse(s.Substring(1)) % 2 != 0).Select(s => new { Source = "B", Value = s });

var connectable = Observable.Merge(a, b).Publish();
connectable.Where(x => x.Source.Equals("A")).Subscribe(s => Console.WriteLine("Subscriber A got: {0}", s));
connectable.Where(x => x.Source.Equals("B")).Subscribe(s => Console.WriteLine("Subscriber B got: {0}", s));
connectable.Connect();

我看到以下输出:

Doing expensive operation
Doing expensive operation
Subscriber A got: { Source = A, Value = #0 }
Doing expensive operation
Doing expensive operation
Subscriber B got: { Source = B, Value = #1 }

(输出继续,为简洁起见截断。)

如何与两个订阅者共享 observable?

【问题讨论】:

    标签: c# system.reactive observable


    【解决方案1】:

    你发布了错误的 observable。

    您正在合并当前代码,然后像这样Observable.Merge(a, b).Publish(); 发布。现在,由于 ab 是针对 expensive 定义的,因此您仍然可以获得对 expensive 的两个订阅。

    订阅创建这些管道:

    如果您从代码中取出.Publish();,您可以看到这一点。输出变为:

    Doing an expensive operation
    Doing an expensive operation
    Doing an expensive operation
    Doing an expensive operation
    Subscriber A got: { Source = A, Value = #0 }
    Doing an expensive operation
    Doing an expensive operation
    Doing an expensive operation
    Doing an expensive operation
    Subscriber B got: { Source = B, Value = #1 }
    

    这会创建这些管道:

    因此,通过将.Publish() 移回expensive,您可以消除问题。这才是你真正需要它的地方,因为它毕竟是昂贵的操作。

    这是您需要的代码:

    var foregroundScheduler = new NewThreadScheduler(ts => new Thread(ts) { IsBackground = false });
    var timer = Observable.Timer(TimeSpan.Zero, TimeSpan.FromSeconds(10), foregroundScheduler);
    var expensive = timer.Select(i =>
    {
        // Converting to strings is an expensive operation
        Console.WriteLine("Doing an expensive operation");
        return string.Format("#{0}", i);
    });
    
    var connectable = expensive.Publish();
    
    var a = connectable.Where(s => int.Parse(s.Substring(1)) % 2 == 0).Select(s => new { Source = "A", Value = s });
    var b = connectable.Where(s => int.Parse(s.Substring(1)) % 2 != 0).Select(s => new { Source = "B", Value = s });
    
    var merged = Observable.Merge(a, b);
    
    merged.Where(x => x.Source.Equals("A")).Subscribe(s => Console.WriteLine("Subscriber A got: {0}", s));
    merged.Where(x => x.Source.Equals("B")).Subscribe(s => Console.WriteLine("Subscriber B got: {0}", s));
    
    connectable.Connect();
    

    这很好地产生了以下内容:

    Doing an expensive operation
    Subscriber A got: { Source = A, Value = #0 }
    Doing an expensive operation
    Subscriber B got: { Source = B, Value = #1 }
    Doing an expensive operation
    Subscriber A got: { Source = A, Value = #2 }
    Doing an expensive operation
    Subscriber B got: { Source = B, Value = #3 }
    

    这为您提供了这些管道:

    您可以从这张图片中看到仍然存在重复。没关系,因为这些零件不贵。

    复制实际上很重要。管道的共享部分使它们的端点容易受到错误的影响,因此容易被提前终止。共享越少,代码的健壮性就越好。只有当您进行昂贵的操作时,您才应该担心发布。否则,您应该让管道成为它们自己。

    这里有一个例子来展示它。如果您没有发布的源,那么如果一个源产生错误,那么它不会拉下所有管道。

    但是一旦你引入了一个共享的 observable,那么一个错误就会导致所有的管道崩溃。

    【讨论】:

    • 您能否详细说明“由于 a & b 是针对昂贵的定义的,您仍然可以获得两次昂贵的订阅”?我不完全理解那部分。
    • @Whymarrh - 我已经添加了进一步的解释。
    • “重复实际上很重要。管道的共享部分使它们的端点容易受到错误的影响,因此容易被提前终止。共享越少,代码的健壮性就越好。”我不明白这一点。你能解释更多吗?
    • @TimothyShields - 我添加了更多细节。
    • @Enigmativity:我不明白你关于沿着管道传播错误的说法。这是否意味着错误也会向后传播,即如果在图形的特定路径中间发生错误,则该路径上的所有节点都将被终止?我在想错误只会向前传播,即错误生成流和后续流将终止。
    【解决方案2】:

    一种可能的解决方法:

    var foregroundScheduler = new NewThreadScheduler(ts => new Thread(ts) { IsBackground = false });
    var timer = Observable.Timer(TimeSpan.Zero, TimeSpan.FromSeconds(10), foregroundScheduler);
    var expensive = timer.Select(i =>
    {
        // Converting to strings is an expensive operation
        Console.WriteLine("Doing an expensive operation");
        return string.Format("#{0}", i);
    });
    
    var subj = new ReplaySubject<string>();
    expensive.Subscribe(subj);
    
    var a = subj.Where(s => int.Parse(s.Substring(1)) % 2 == 0).Select(s => new { Source = "A", Value = s });
    var b = subj.Where(s => int.Parse(s.Substring(1)) % 2 != 0).Select(s => new { Source = "B", Value = s });
    
    var merged = Observable.Merge(a, b);
    merged.Where(x => x.Source.Equals("A")).Subscribe(s => Console.WriteLine("Subscriber A got: {0}", s));
    merged.Where(x => x.Source.Equals("B")).Subscribe(s => Console.WriteLine("Subscriber B got: {0}", s));
    

    上面的示例本质上创建了一个新的中间可观察对象,它发出昂贵操作的结果。这允许您订阅昂贵操作的结果,而不是应用到计时器的昂贵转换。

    这样你会看到:

    Doing an expensive operation
    Subscriber A got: { Source = A, Value = #0 }
    Doing an expensive operation
    Subscriber B got: { Source = B, Value = #1 }
    

    (输出继续,为简洁起见被截断。)

    或者,您可以将呼叫转移到 PublishConnect

    var foregroundScheduler = new NewThreadScheduler(ts => new Thread(ts) {IsBackground = false});
    var timer = Observable.Timer(TimeSpan.Zero, TimeSpan.FromSeconds(10), foregroundScheduler);
    var expensive = timer.Select(i =>
    {
        // Converting to strings is an expensive operation
        Console.WriteLine("Doing an expensive operation");
        return string.Format("#{0}", i);
    }).Publish();
    
    var a = expensive.Where(s => int.Parse(s.Substring(1)) % 2 == 0).Select(s => new { Source = "A", Value = s });
    var b = expensive.Where(s => int.Parse(s.Substring(1)) % 2 != 0).Select(s => new { Source = "B", Value = s });
    
    var merged = Observable.Merge(a, b);
    merged.Where(x => x.Source.Equals("A")).Subscribe(s => Console.WriteLine("Subscriber A got: {0}", s));
    merged.Where(x => x.Source.Equals("B")).Subscribe(s => Console.WriteLine("Subscriber B got: {0}", s));
    
    expensive.Connect();
    

    为什么是 ReplaySubject,而不仅仅是 Subject 或其他主题?

    .NET Rx 实现中的Subject,默认情况下是the ReactiveX documentation calls a PublishSubject,它只向观察者发出那些在订阅时间之后由源 Observable 发出的项目。另一方面,ReplaySubject 向任何观察者发出源 Observable 发出的所有项目,无论观察者何时订阅。如果我们在第一个示例中使用普通主题,订阅 subj 到计时器将导致订阅 subj 错过从主题订阅昂贵操作到订阅中级主题 (subj)。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2016-08-09
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-07-28
      • 1970-01-01
      • 1970-01-01
      • 2023-01-08
      相关资源
      最近更新 更多