【问题标题】:Preserve Sorting While Merging Two Observables合并两个 Observable 时保留排序
【发布时间】:2019-07-28 05:35:15
【问题描述】:

我想合并 2 个 observables 并保持顺序(可能基于选择器)。我还想对 observable 的来源施加背压。

所以选择器会选择其中一个项目通过 observable 推进,而另一个项目也会等待另一个项目来比较。

Src1、Src2 和 Result 都是 IObservable<T> 类型。

Src1: { 1,3,6,8,9,10 }
Src2: { 2,4,5,7,11,12 }
Result: 1,2,3,4,5,6,7,8,9,10,11,12

Timeline:
Src1:    -1---3----6------8----9-10
Src2:    --2-----4---5-7----11---------12
Result:  --1--2--3-4-5-6--7-8--9-10-11-12
  1. 在上面的示例中,src1 发出 '1' 并被阻塞,直到 src2 发出 这是第一项,“2”。
  2. 应用选择最小项目的选择器 从 src1 中选择项目。
  3. Src2 现在等待下一项(来自 src1)与其比较 当前项目 ('2')。
  4. 当 src1 发出下一项“3”时,再次运行选择,这 从 src2 中选择项目的时间。
  5. 重复此过程,直到其中一个可观察对象完成。然后,剩余的 observable 推送项目直到完成。

这可以通过现有的 .net Rx 方法实现吗?

编辑:注意 2 个源 observables 保证是有序的。

示例测试:

var source1 = new List<int>() { 1, 4, 6, 7, 8, 10, 14 }.AsEnumerable();
var source2 = new List<int>() { 2, 3, 5, 9, 11, 12, 13, 15 }.AsEnumerable();

var src1 = source1.ToObservable();
var src2 = source2.ToObservable();

var res = src1.SortedMerge(src2, (a, b) =>
    {
       if (a <= b)
           return a;
       else
           return b;
    });

res.Subscribe((x) => Console.Write($"{x}, "));

DesiredResult: 1,2,3,4,5,6,7,8,9,10,11,12,13,14,15

【问题讨论】:

  • 你是如何合并 observables 的?请给出你的代码。因为Merge 运算符不会按顺序等待源
  • @MohammadOmidvar - 这就是他的要求。
  • 是的,我没有正确理解。那么,“现有的 .net Rx 方法”是指内置运算符还是完全有可能? @msauce4。因为有可能写出能够做到这一点的主题。
  • @MohammadOmidvar 我希望使用内置运算符的任何组合来实现,但如果这不可能,我对编写主题感到好奇。假设我们忽略了背压问题,那能做到吗?有点类似这个,stackoverflow.com/questions/50298555/…,问题。

标签: c# .net .net-core system.reactive


【解决方案1】:

这很有趣。不得不稍微调整一下算法。可以进一步改进。

假设:

  1. 有两个流,streamAstreamB 的普通类型T
  2. 两个流分别排序为streamA[i] &lt; streamA[i+1]streamB[i] &lt; stream[i+1]
  3. 您不能假设streamA[i]streamB[i] 之间存在任何关系。
  4. 流 A 和 B 是谨慎的:不会从两者发出相同的元素。如果发生这种情况,我会抛出NotImplementedException这个案子很容易处理,但我想避免歧义。
  5. 有一个函数min 用于类型T
  6. 没有对两个流的相对速度做出任何假设,但如果一个始终比另一个快,则背压将是一个问题。

这是我使用的算法:

  • 假设有两个队列,qAqB
  • 当您从streamA 获取项目时,将其排入qA
  • 当您从streamB 获得项目时,将其排入qB
  • 虽然both qAqB 中有一个项目,但比较两个队列的顶部项目。删除并发出这两者的最小值。如果两个队列仍然非空,请重复。
  • 如果streamAstreamB 完成,则转储队列的内容并终止。 注意这无疑是懒惰的,可能应该改为转储,然后继续返回未完成的可观察对象

代码如下:

public static IObservable<T> SortedMerge<T>(this IObservable<T> source, IObservable<T> other)
{
    return SortedMerge(source, other, (a, b) => Enumerable.Min(new[] { a, b}));
}

public static IObservable<T> SortedMerge<T>(this IObservable<T> source, IObservable<T> other, Func<T, T, T> min)
{
    return source
        .Select(i => (key: 1, value: i)).Materialize()
        .Merge(other.Select(i => (key: 2, value: i)).Materialize())
        .Scan((qA: ImmutableQueue<T>.Empty, qB: ImmutableQueue<T>.Empty, exception: (Exception)null, outputMessages: new List<T>()), 
            (state, message) =>
        {
            if (message.Kind == NotificationKind.OnNext)
            {
                var key = message.Value.key;
                var value = message.Value.value;
                var qA = state.qA;
                var qB = state.qB;
                if (key == 1)
                    qA = qA.Enqueue(value);
                else
                    qB = qB.Enqueue(value);
                var output = new List<T>();
                while(!qA.IsEmpty && !qB.IsEmpty)
                {
                    var aVal = qA.Peek();
                    var bVal = qB.Peek();
                    var minVal = min(aVal, bVal);
                    if(aVal.Equals(minVal) && bVal.Equals(minVal))
                        throw new NotImplementedException();

                    if(aVal.Equals(minVal))
                    {
                        output.Add(aVal);
                        qA = qA.Dequeue();
                    }
                    else
                    {
                        output.Add(bVal);
                        qB = qB.Dequeue();
                    }
                }
                return (qA, qB, null, output);
            }
            else if (message.Kind == NotificationKind.OnError)
            {
                return (state.qA, state.qB, message.Exception, new List<T>());
            }
            else //message.Kind == NotificationKind.OnCompleted
            {
                var output = state.qA.Concat(state.qB).ToList();
                return (ImmutableQueue<T>.Empty, ImmutableQueue<T>.Empty, null, output);
            }
        })
        .Publish(tuples => Observable.Merge(
            tuples
                .Where(t => t.outputMessages.Any() && (!t.qA.IsEmpty || !t.qB.IsEmpty))
                .SelectMany(t => t.outputMessages
                    .Select(v => Notification.CreateOnNext<T>(v))
                    .ToObservable()
            ),
            tuples
                .Where(t => t.outputMessages.Any() && t.qA.IsEmpty && t.qB.IsEmpty)
                .SelectMany(t => t.outputMessages
                    .Select(v => Notification.CreateOnNext<T>(v))
                    .ToObservable()
                    .Concat(Observable.Return(Notification.CreateOnCompleted<T>()))
            ),
            tuples
                .Where(t => t.exception != null)
                .Select(t => Notification.CreateOnError<T>(t.exception))
        ))
        .Dematerialize();

ImmutableQueue 来自System.Collections.Immutable。需要Scan 来跟踪状态。由于OnCompleted 处理,需要实现。诚然,这是一个复杂的解决方案,但我不确定是否有更清洁的以 Rx 为中心的方式。

如果您有任何需要澄清的问题,请告诉我。

【讨论】:

  • 我在问题中添加了一个示例测试。运行您的方法会产生以下结果:“1, 2, 3, 4, 5, 6, 8, 7, 9, 10, 11, 13, 12, 14”。 8 和 13 放错了地方,而 15 似乎没有被抛弃。知道为什么吗?
  • 另外,是什么驱动了 Materialize() 调用的必要性?
  • 将 observable 形式更改为 ToObservable(NewThreadScheduler.Default) 会导致测试返回预期结果。
猜你喜欢
  • 2021-06-26
  • 2017-03-13
  • 1970-01-01
  • 2015-11-26
  • 1970-01-01
  • 2014-05-08
  • 1970-01-01
  • 2013-02-28
相关资源
最近更新 更多