【问题标题】:How to implement RxJS flatMapLatestTwo如何实现 RxJS flatMapLatestTwo
【发布时间】:2015-08-27 22:31:47
【问题描述】:

RxJS 的 flatMapLatest 扁平化了最新的(只有一个)嵌套的 Observable。我有一个用例,我不想要 flatMap(将过去所有嵌套的 Observables 展平),并且我不想要 flatMapWithConcurrency(因为它有利于旧的 Observables,而不是最新的 Observables),所以我想要的是 flatMapLatestTwo 或某些版本的 flatMapLatest,您可以在其中指定并发嵌套 Observable 的最大数量,例如flatMapLatest(2, selectorFn).

这是我想要的输出(_X 指的是嵌套的 Observable XeX 指的是它的第 X 个 onNext 事件):

_0e0
_0e1
 _1e0
_0e2
 _1e1
  _2e0
 _1e2
  _2e1
   _3e0
  _2e2
   _3e1
    _4e0
   _3e2
    _4e1
   _3e3
    _4e2
    _4e3

这是 flatMapLatest 产生的:

_0e0
_0e1
 _1e0
 _1e1
  _2e0
  _2e1
   _3e0
   _3e1
    _4e0
    _4e1
    _4e2
    _4e3

我更喜欢使用现有运算符的解决方案,而不是实现这个低级。

【问题讨论】:

  • 我试过 map 到 Observable(这给了我一个 Observable-of-Observables),然后是 bufferWithCount(2,1),然后是 flatMapLatest,但我得到了一些奇怪的重复和遗漏.

标签: javascript rxjs


【解决方案1】:

这看起来很幼稚。我正在寻找改进的方法,但这里是:

Rx.Observable.prototype.flatMapLatestN = function (count, transform) {

  let queue = [];

  return this.flatMap(x => {
    return Rx.Observable.create(observer => {

      let disposable;

      if (queue.length < count) {
        disposable = transform(x).subscribe(observer);
        queue.push(observer);
      }
      else {
        let earliestObserver = queue[0];
        if (earliestObserver) {
          earliestObserver.onCompleted();
        }

        disposable = transform(x).subscribe(observer);
        queue.push(observer);
      }

      return () => {
        disposable.dispose();
        let i = queue.indexOf(observer);
        queue.splice(i, 1);
      };
    });
  });
};

测试:

function space(n) {
  return Array(n+1).join(' ');
}

Rx.Observable
  .interval(1000)
  .take(6)
  .flatMapLatestN(2, (x) => {
    return Rx.Observable
      .interval(300)
      .take(10)
      .map(n => `${space(x*4)}${x}-${n}`);
  })
  .subscribe(console.log.bind(console));

它会输出:

0-1
0-2
0-3
    1-0
0-4
    1-1
0-5
    1-2
    1-3
        2-0
    1-4
        2-1
    1-5
        2-2
        2-3
            3-0
        2-4
            3-1
        2-5
            3-2
            3-3
                4-0
            3-4
                4-1
            3-5
                4-2
                4-3
                    5-0
                4-4
                    5-1
                4-5
                    5-2
                4-6
                    5-3
                4-7
                    5-4
                4-8
                    5-5
                4-9
                    5-6
                    5-7
                    5-8
                    5-9

【讨论】:

  • 是的,我考虑了几个小时并得出结论,如果没有自定义运算符,可能没有任何方法可以做到这一点。感谢这个运算符的实现。
【解决方案2】:

这是一个使用内置运算符的解决方案。首先,我们将 Observable 拆分为 N 个 observable,其中每个 observable 具有序列中对应的第 N 个最新项目。然后我们flatMapLatest每一个并合并它们。

Rx.Observable.prototype.flatMapLatestN = function(N, selector, thisArg) {
    var self = this;
    return Rx.Observable.merge(Rx.Observable.range(0, N).flatMap(function(n) {
        return self.filter(function(x, i) {
            return i % N === n;
        }).flatMapLatest(selector, thisArg);
    }));
}

或者在 ES2015 中:

Rx.Observable.prototype.flatMapLatestN = function(N, selector, thisArg) {
    const {merge, range} = Rx.Observable;
    return merge(
        range(0, N)
            .flatMap(n => 
                this.filter((x, i) => i % N === n).flatMapLatest(selector, thisArg))
    );
}

使用与戴维相同的测试:

N=1 的输出(与 flatMapLatest 相同):

0-0
0-1
0-2
    1-0
    1-1
    1-2
        2-0
        2-1
        2-2
            3-0
            3-1
            3-2
                4-0
                4-1
                4-2
                    5-0
                    5-1
                    5-2
                    5-3
                    5-4
                    5-5
                    5-6
                    5-7
                    5-8
                    5-9

N=2 的输出:

0-0
0-1
0-2
0-3
    1-0
0-4
    1-1
0-5
    1-2
    1-3
        2-0
    1-4
        2-1
    1-5
        2-2
        2-3
            3-0
        2-4
            3-1
        2-5
            3-2
            3-3
                4-0
            3-4
                4-1
            3-5
                4-2
                4-3
                    5-0
                4-4
                    5-1
                4-5
                    5-2
                4-6
                    5-3
                4-7
                    5-4
                4-8
                    5-5
                4-9
                    5-6
                    5-7
                    5-8
                    5-9

N=3 的输出:

0-0
0-1
0-2
0-3
    1-0
0-4
    1-1
0-5
    1-2
0-6
    1-3
        2-0
0-7
    1-4
        2-1
0-8
    1-5
        2-2
    1-6
        2-3
            3-0
    1-7
        2-4
            3-1
    1-8
        2-5
            3-2
        2-6
            3-3
                4-0
        2-7
            3-4
                4-1
        2-8
            3-5
                4-2
            3-6
                4-3
                    5-0
            3-7
                4-4
                    5-1
            3-8
                4-5
                    5-2
            3-9
                4-6
                    5-3
                4-7
                    5-4
                4-8
                    5-5
                4-9
                    5-6
                    5-7
                    5-8
                    5-9

【讨论】:

  • 这是一个可爱的解决方案。
  • 另外,这不是一张相当科学准确的文档图片,狗没有几个鼻子,不像这张照片。请解决这个问题。谢谢。
  • ES6 答案中的合并似乎与 ES5 不同。在 ES6 中,它缺少“this”。
  • @Daiwei 感谢 ES6 中的箭头函数,this.filter 中的 this 的上下文是正确的上下文。
【解决方案3】:

我的答案是 C#。对不起。

你没有指定你的 observables 是热的还是冷的。可能是您的奇怪数字来自这样一个事实,即您的窗口对您的“内部”可观察对象进行了新订阅,因为它从窗口中的第一个被推到第二个。我的第一次尝试也是这样:

var q = Observable.Interval(TimeSpan.FromSeconds(1))
        .Select(i => Observable.Interval(TimeSpan.FromMilliseconds(100))
        .Select(x => $"_{i}e{x}"));

w = q.Zip(q.Skip(1), (prev, curr)=> prev.Merge(curr)).Switch(); 

我认为你将努力做任何事情来避免这种情况,除非你创建一个运营商,因为有状态参与管理它。 (显然有人会在这里证明我错了!!)

这是我的运算符方法,它也恰好支持您要求的参数化。

public static class Ex
{
    public static IObservable<T> SelectManyLatest<T>(this IObservable<IObservable<T>> source, int latest)
    {
        return Observable.Create<T>(o => 
        {
            var d = new Queue<IDisposable>();

            source.Subscribe(os => 
            {
                if(d.Count == latest)
                    d.Dequeue().Dispose();

                d.Enqueue(os.Subscribe(o.OnNext, o.OnError, () => {}));

            }, o.OnError, o.OnCompleted);

            return Disposable.Create(()=>new CompositeDisposable(d).Dispose());
        });     
    }
}

再次为 C# 感到抱歉

【讨论】:

  • 是的,如果没有新的操作员,我看不出你怎么能做到。但是,如果第二个内部可观察对象在第三个内部可观察对象到达之前和第一个完成之前完成,则您的运算符将无法正常工作:如果发生这种情况,我相信您的运算符将取消订阅第一个内部可观察对象...我认为您需要删除如果完成则从队列中订阅
  • 感谢您富有洞察力的回答,但我会接受更适合该问题的 JavaScript 回答。
  • @Brandon 是的,这有点快。我看到了这个问题,但解决它涉及定义一些问题中不存在的政策,所以我忽略了它。您的政策似乎是明智的,但这取决于流所代表的内容,第二个的沉默和第三个的值实际上可能是要求。
猜你喜欢
  • 1970-01-01
  • 2017-10-12
  • 2018-07-12
  • 1970-01-01
  • 2021-09-08
  • 2020-07-30
  • 2018-09-08
  • 1970-01-01
  • 2017-06-12
相关资源
最近更新 更多