【问题标题】:RxJS test equality of two streams regardless of order无论顺序如何,RxJS 测试两个流的相等性
【发布时间】:2020-06-18 04:23:03
【问题描述】:

RxJS 提供了 sequenceEqual 运算符来按顺序比较两个流。无论顺序如何,如何测试两个流的相等性?

伪代码:

//how do we implement sequenceEqualUnordered?
from([1,2,3]).pipe(sequenceEqualUnordered(from([3,2,1]))).subscribe((eq) => 
  console.log("Eq should be true because both observables contain the same values")
)

在我的特定用例中,我需要等到发出一组特定的值或出错,但我不在乎它们以什么顺序发出。我只关心每个感兴趣的值都发出一次。

【问题讨论】:

  • 在比较之前订购这两个?
  • 我上面的例子是伪代码。在实际情况下,我不知道来自第一个可观察值的值的顺序,或者不一定知道第二个可观察值中值的顺序。我只关心完成后的值是否相等。

标签: rxjs


【解决方案1】:

这是我的解决方案:

import { Observable, OperatorFunction, Subscription } from 'rxjs';

export function sequenceEqualUnordered<T>(compareTo: Observable<T>, comparator?: (a: T, b: T) => number): OperatorFunction<T, boolean> {
  return (source: Observable<T>) => new Observable<boolean>(observer => {
    const sourceValues: T[] = [];
    const destinationValues: T[] = [];
    let sourceCompleted = false;
    let destinationCompleted = false;

    function onComplete() {
      if (sourceCompleted && destinationCompleted) {
        if (sourceValues.length !== destinationValues.length) {
          emit(false);
          return;
        }

        sourceValues.sort(comparator);
        destinationValues.sort(comparator);

        emit(JSON.stringify(sourceValues) === JSON.stringify(destinationValues));
      }
    }

    function emit(value: boolean) {
      observer.next(value);
      observer.complete();
    }

    const subscriptions = new Subscription();

    subscriptions.add(source.subscribe({
      next: next => sourceValues.push(next),
      error: error => observer.error(error),
      complete: () => {
        sourceCompleted = true;
        onComplete();
      }
    }));
    subscriptions.add(compareTo.subscribe({
      next: next => destinationValues.push(next),
      error: error => observer.error(error),
      complete: () => {
        destinationCompleted = true;
        onComplete();
      }
    }));

    return () => subscriptions.unsubscribe();
  });
}

由于许多 RxJS 操作符都有一些输入参数并且它们都返回函数,sequenceEqualUnordered 也有一些输入参数(大部分与 Rx 的 sequenceEqual 操作符相同)并且它返回一个函数。而这个返回函数有Observable&lt;T&gt;作为source类型,并且有Observable&lt;boolean&gt;作为返回类型。

创建一个将发出boolean 值的新 Observable 正是您所需要的。您基本上希望从source 和compareTo Observables 中收集所有值(并将它们存储到sourceValues 和destinationValues 数组中)。为此,您需要同时订阅 source 和 compareTo Observables。但是,为了能够处理订阅,必须创建一个 subscriptions 对象。在创建source 和compareTo 的新订阅时,只需将add 订阅到subscriptions 对象即可。

订阅其中任何一个时,将发出的值收集到next 处理程序中的适当sourceValues 或destinationValues 数组。如果发生任何错误,请将它们传播到error 处理程序中的observer。在 complete 处理程序中,设置适当的 sourceCompleted 或 destinationCompleted 标志以指示哪个 Observable 已完成。

然后,在onComplete 中检查它们是否都已完成,如果都已完成,则比较发出的值并发出适当的布尔值。如果sourceValues 和destinationValues 数组的长度不同,它们就不能相等,所以发出false。之后,基本上对数组进行排序,然后compare这两个。

发出时,同时发出值和complete 通知。

另外,传递给new Observable&lt;boolean&gt; 的函数的返回值应该是unsubscribe 函数。基本上,当有人取消订阅new Observable&lt;boolean&gt; 时,它也应该取消订阅source 和compareTo Observables,这是通过调用() =&gt; subscriptions.unsubscribe() 来完成的。 subscriptions.unsubscribe() 将取消订阅 add 的所有订阅。

TBH,我还没有为这个运算符编写任何测试,所以我不能完全确定我已经涵盖了所有边缘情况。

【讨论】:

    【解决方案2】:

    我的第一个想法。对两者都使用toArray,然后将它们一起使用zip,最后对结果进行排序和比较?

    【讨论】:

      猜你喜欢
      • 2013-10-12
      • 2019-01-28
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-03-12
      • 2019-06-11
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多