这是我的解决方案:
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<T>作为source类型,并且有Observable<boolean>作为返回类型。
创建一个将发出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<boolean> 的函数的返回值应该是unsubscribe 函数。基本上,当有人取消订阅new Observable<boolean> 时,它也应该取消订阅source 和compareTo Observables,这是通过调用() => subscriptions.unsubscribe() 来完成的。 subscriptions.unsubscribe() 将取消订阅 add 的所有订阅。
TBH,我还没有为这个运算符编写任何测试,所以我不能完全确定我已经涵盖了所有边缘情况。