【问题标题】:Combine overlapping observable streams but take most recent values合并重叠的可观察流,但采用最新的值
【发布时间】:2015-01-08 14:13:35
【问题描述】:

我有一个使用 Rx 组合流的用例。我有 3 个流输出值:

S1: ----1----2----3----4----5-----6
S2: ----a---------c---------d------
S3: ---------x---------y----------z

我想合并 (zip) 流 1 & 2 和 1 & 3,逻辑是 2 & 3 永远不能同时输出一个值。

所以想要的输出是:

1,a
2,x
3,c
4,y
5,d
6,z

但实际发生的是:

1,a
1,x
2,c
2,y
3,d
3,z

关于如何获得所需输出的任何线索?

代码:

Subject<int> s1 = new Subject<int>();
Subject<string> s2 = new Subject<string>();
Subject<string> s3 = new Subject<string>();

var zip1 = s1.Zip(s2, (x, y) => {
    return new Tuple<int, string>(x, y);
});

var zip2 = s1.Zip(s3, (x, y) => {
    return new Tuple<int, string>(x, y);
});

zip1.Subscribe(Console.WriteLine, () => Console.WriteLine("Completed"));
zip2.Subscribe(Console.WriteLine, () => Console.WriteLine("Completed"));

s1.OnNext(1);
s2.OnNext("a");

s1.OnNext(2);
s3.OnNext("x");

s1.OnNext(3);
s2.OnNext("c");

s1.OnNext(4);
s3.OnNext("y");

s1.OnNext(5);
s2.OnNext("d");

s1.OnNext(6);
s3.OnNext("z");

【问题讨论】:

  • s2 和 s3 是否总是交替出现,或者一个连续输出多个值而另一个在该时间跨度内不输出任何值?
  • s2和s3可以连续输出多个值任意次数,所以没有模式。
  • 如果 s1 输出多个值而 s2 或 s3 没有输出任何值会发生什么?
  • 我没想到那么远,实际上我会尝试编写代码,这样这种情况就不会发生。

标签: c# system.reactive


【解决方案1】:

对于您拥有的每个流 s2 和 s3,您可以创建一个 Subject,订阅原始流,然后在等待 @987654323 的下一个值后将值传递给新的 Subject @:

var zip1 = new Subject<Tuple<int, string>>();
s2.Subscribe(async value =>
{
    var other = await s1.FirstAsync();
    zip1.OnNext(Tuple.Create(other, value));
});

var zip2 = new Subject<Tuple<int, string>>();
s3.Subscribe(async value =>
{
    var other = await s1.FirstAsync();
    zip2.OnNext(Tuple.Create(other, value));
});

您已声明您将管理流的计时,以便从s1 为其他流的每个值组合产生一个值。如果你能做到这一点,你会没事的。如果从s1 产生的多个值在另一个流中没有任何对应值,它将被忽略,并且如果从其他流组合产生多个值而在s1 中没有对应值,则下一个值产生于s1 将重复。

还有另一种更人为的选项,但在值之间的时间安排上更灵活。这种方法是首先将每个流投影到具有价值但将其与其他流区分开来的东西中。然后可以合并流,使用s1 压缩,然后可以创建多个流来过滤掉每个特定流的值。然后需要将这些值投影回原始值。

var zip = s2.Select(s => Tuple.Create(1, s))
    .Merge(s3.Select(s => Tuple.Create(2, s)))
    .Zip(s1, Tuple.Create);

var zip1 = zip.Where(pair => pair.Item1.Item1 == 1)
    .Select(pair => Tuple.Create(pair.Item2, pair.Item1.Item2));
var zip2 = zip.Where(pair => pair.Item1.Item1 == 2)
    .Select(pair => Tuple.Create(pair.Item2, pair.Item1.Item2));

如果您实际上不需要单独的流,而只想将输出复制为单个流,那么它实际上要简单得多:

var output = s1.Zip(s2.Merge(s3), Tuple.Create);
output.Subscribe(Console.WriteLine, () => Console.WriteLine("Completed"));

【讨论】:

  • 是的,我打算在最后一行 Zip+Merge sn-p 中发布您所拥有的内容。根据弹珠图和线条简单解决问题的描述。
  • @Brandon 假设他不需要单独的流,只需要一个组合流。如果这是他想要的,这当然是最好的解决方案。
【解决方案2】:

这不只是一个例子吗?

var query =
    s1.Zip(s2.Merge(s3),
        (x, y) => new Tuple<int, string>(x, y));

鉴于您在问题中提供的代码,我得到了这个结果:

【讨论】:

  • 绝对更接近我需要的解决方案。问题是 s2 和 s3 不一定相互了解。
  • @David - 那么您的示例代码是一个特例。我认为您需要提供一些更现实的代码来说明您的需求。
猜你喜欢
  • 2019-10-05
  • 2019-10-28
  • 2018-08-07
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-12-30
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多