【问题标题】:RxJs split stream into multiple streamsRxJs 将流拆分为多个流
【发布时间】:2017-10-19 18:40:51
【问题描述】:

如何根据分组方法将一个永不结束的流拆分为多个结束流?

--a--a-a-a-a-b---b-b--b-c-c---c-c-d-d-d-e...>

进入这些可观察对象

--a--a-a-a-a-|
             b---b-b--b-|
                        c-c---c-c-|
                                  d-d-d-|
                                        e...>

如您所见,a 位于开头,收到b 后,我将不再收到a,因此应该结束。这就是为什么普通的groupBy不好。

【问题讨论】:

  • 为了清楚起见,您想要一个在每次序列更改时发出 observable 的 observable,还是想要一个发出连续项目数组的 observable?
  • 第一个。我想要在原始发射时发射值的 observables。如果原来的已经改变了,(从ab)那么a的observable需要停止

标签: javascript rxjs system.reactive


【解决方案1】:

你可以使用 windowshare 源 Observable。 bufferCount(2, 1)还有一个小技巧:

const str = 'a-a-a-a-a-b-b-b-b-c-c-c-c-d-d-d-e';
const source = Observable.from(str.split('-'), Rx.Scheduler.async).share();

source
    .bufferCount(2, 1) // delay emission by one item
    .map(arr => arr[0])
    .window(source
        .bufferCount(2, 1) // keep the previous and current item
        .filter(([oldValue, newValue]) => oldValue !== newValue)
    )
    .concatMap(obs => obs.toArray())
    .subscribe(console.log);

这个打印(因为toArray()):

[ 'a', 'a', 'a', 'a', 'a' ]
[ 'b', 'b', 'b', 'b' ]
[ 'c', 'c', 'c', 'c' ]
[ 'd', 'd', 'd' ]
[ 'e' ]

这个解决方案的问题是订阅source 的顺序。我们需要window 通知程序在第一个bufferCount 之前订阅。否则,一个项目首先被进一步推送,然后检查它是否与 .filter(([oldValue, newValue]) ...) 上一个项目不同。

这意味着需要在window 之前将发射延迟一个(即第一个.bufferCount(2, 1).map(arr => arr[0])

或者用publish()自己控制订阅顺序可能更容易:

const str = 'a-a-a-a-a-b-b-b-b-c-c-c-c-d-d-d-e';
const source = Observable.from(str.split('-'), Rx.Scheduler.async).share();

const connectable = source.publish();

connectable
    .window(source
        .bufferCount(2, 1) // keep the previous and current item
        .filter(([oldValue, newValue]) => oldValue !== newValue)
    )
    .concatMap(obs => obs.toArray())
    .subscribe(console.log);

connectable.connect();

输出是一样的。

【讨论】:

  • 是的,我们似乎失去了非常有用的 publish(selector: Observable<T> => Observable<R>) : Observable<R> 重载发布以简化订阅共享。
  • @martin 如果我错了,请纠正我,但在您的第一个解决方案中,您延迟了可观察到的源比实际需要的更多。而不是 bufferCount(2, 1).map(arr => arr[0]),我认为你可以使用 .delay(0)
  • 另外,为了保留前一项和当前项,您可以将 bufferCount(2, 1) 替换为更优雅的 pairwise() 运算符。
【解决方案2】:

也许有人可以想出一些更简单的方法,但这很有效(小提琴:https://fiddle.jshell.net/uk01njgc/)...

let counter = 0;

let items = Rx.Observable.interval(1000)
.map(value => Math.floor(value / 3))
.publish();

let distinct = items.distinctUntilChanged()
.publish();

distinct
.map(value => {
  return items
  .startWith(value)
  .takeUntil(distinct);
})
.subscribe(obs => {
  let obsIndex = counter++;
  console.log('New observable');
  obs.subscribe(
    value => {
      console.log(obsIndex.toString() + ': ' + value.toString());
    },
    err => console.log(err),
    () => console.log('Completed observable')
  );
});

distinct.connect();
items.connect();

【讨论】:

  • 另外,您可以使用share 而不是publish 为您节省连接调用的麻烦,但我没有尝试。
【解决方案3】:

这是一个包含所有订阅共享的变体...

const stream = ...;

// an Observable<Observable<T>>
// each inner observable completes when the value changes
const split = Observable
  .create(o => {
    const connected = stream.publish();

    // signals each time the values change (ignore the initial value)
    const newWindowSignal = connected.distinctUntilChanged().skip(1);

    // send the observables to our observer
    connected.window(newWindowSignal).subscribe(o);

    // now "start"
    return connected.connect();
  });

【讨论】:

  • 它让我感到[a,a,a,a,b], [b,b,b,c]
【解决方案4】:
import { from  } from 'rxjs'; 
import { share, window, map, publish, switchMap, 
         skip, toArray, distinct, bufferCount  } 
from 'rxjs/operators';


export function splitToArray(){
  const str = 'a-a-a-a-a-b-b-b-b-c-c-c-c-d-d-d-e-e';

  const lettters = from(str.split('-'))
    .pipe(
      share(),
      publish()
    );

  lettters
  .pipe(
    bufferCount(2, 1), 
    window(
      lettters
        .pipe(distinct(),
              skip(1)
            ),
    ),
      switchMap(obs => obs.pipe(
      map(([val1, val2]) => {return val1;}),
      toArray()))
  )
  .subscribe(console.log);

  lettters.connect();
}

我对这个问题的看法。使用 rxjs6 管道函数。

这里的技巧是按照@martin 的建议延迟元素 使用缓冲区计数(2,1)。所以每个元素都会发射两次。除了第一个。 大红色箭头是我们时间轴上创建新窗口时的时刻。之后,将数组映射到获取双精度数组的第一个元素就很简单了。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-01-08
    • 2013-06-12
    • 2023-03-03
    • 2016-06-29
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多