【发布时间】:2015-08-27 22:17:29
【问题描述】:
我有一个可观察的数据流,我正在对其应用操作,分成两个单独的流,对两个流中的每一个应用更多(不同的)操作,然后再次合并在一起。我正在尝试使用 Publish 和 Connect 在两个订阅者之间共享 observable,但每个订阅者似乎都在使用单独的流。也就是说,在下面的示例中,我看到 为两个订阅者 为流中的每个项目打印了一次“执行昂贵的操作”。 (想象一下昂贵的操作应该在所有订阅者之间只发生一次,因此我试图重用流。)我使用Publish 和Connect 尝试与两个订阅者共享合并的 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