【问题标题】:RxJS: Working on groupBy and Observable.fromEvents in NodeJSRxJS:在 NodeJS 中处理 groupBy 和 Observable.fromEvents
【发布时间】:2016-10-26 09:26:40
【问题描述】:

我对 RxJS 印象深刻,并开始着手研究。但是,我的以下 nodejs 代码至少对我来说不能按预期工作。

let events = new EventEmitter();
let source = Rx.Observable.fromEvent( events, 'data' );

source
    .groupBy( event => event.type )
    .flatMap( group => group.reduce( ( acc, cur ) => _.merge( acc, cur ), [] ) )
    .subscribe( ( data ) => {
        console.log( data );
    } );


events.emit( 'data', { 'type': 1, msg: 'Test 1' } );
events.emit( 'data', { 'type': 1, msg: 'Test 2' } );
events.emit( 'data', { 'type': 2, msg: 'Test 3' } );

我希望subscribe 产生一些输出

【问题讨论】:

  • 这应该产生什么输出?
  • 它实际上什么也没产生。我期望根据groupBy 分组的数组数组。 [ [ { 'type': 1, msg: 'Test 1' }, { 'type': 1, msg: 'Test 2' ], [ { 'type': 2, msg: 'Test 3' } ] }.
  • 您从哪里获得 EventEmitter?它是starndard rxjs 库的一部分还是ng2 类?你有没有机会把它放在一个正在运行的 jsbin 中?
  • @Meir 这是一个标准的node.js 类nodejs.org/api/events.html#events_class_eventemitter
  • 我会准备一个在线运行版本。

标签: node.js reactive-programming rxjs


【解决方案1】:

这是我在 JSBin http://jsbin.com/cufana/edit?js,console 中创建的稍有改动的版本。

我在您的代码中看到的问题是您有一个未完成的可观察对象。如果您正在执行 groupBy,则只要源 observable 未完成,内部保存的结果(即分组结果)就不会被推送。

let events = new Rx.Subject();

events
    .groupBy( event => event.type)
    .flatMap(group => group.reduce((acc, curr) => [...acc, curr], []))
    .subscribe( ( data ) => {
        console.log( data );
    } );


events.next({ 'type': 1, msg: 'Test 1' } );
events.next({ 'type': 1, msg: 'Test 2' } );
events.next({ 'type': 2, msg: 'Test 3' } );
events.complete();

在这里你可以看到我已经将 eventemitter 更改为一个主题以获得一个没有 angular2 依赖的工作 jsbin。我正在完成这个主题,所以我可以从 groupBy 观察到的来源已经完成。这将推动结果通过。

其余代码非常正确。

如果 EventEmitter 确实来自 Angular2,我猜你在完成这个时会遇到问题。这可以从子组件中完成吗?

【讨论】:

    【解决方案2】:

    正如其他人所建议的,Observable 永远不会完成,因此任何需要前面的 Observable 完成的运算符都不会发出任何东西:

    groupBy() 运算符会为每个组发出一个GroupedObservable 实例,以便您订阅它。我知道这会产生与您预期不同的结果,但也许您可以使用它:

    const Rx = require('rxjs/Rx');
    const EventEmitter = require('events');
    
    let events = new EventEmitter();
    let source = Rx.Observable.fromEvent(events, 'data');
    
    source
        .groupBy( event => event.type )
        .subscribe( ( groupedObservable ) => {
            groupedObservable.subscribe(val => {
                console.log(groupedObservable.key, val);
            });
        } );
    
    
    events.emit( 'data', { 'type': 1, msg: 'Test 1' } );
    events.emit( 'data', { 'type': 1, msg: 'Test 2' } );
    events.emit( 'data', { 'type': 5, msg: 'Test 3' } );
    

    每个组都有自己的GroupedObservable

    这会打印到控制台:

    1 { type: 1, msg: 'Test 1' }
    1 { type: 1, msg: 'Test 2' }
    5 { type: 5, msg: 'Test 3' }
    

    【讨论】:

      【解决方案3】:

      问题是您的序列没有终止。 groupBy 仅在终止时触发。见jsbin example with both cases

      当您执行Rx.Observable.from(...) 时,您将获得一个完整的序列,否则,您需要手动终止。 groupBy documentation 上的示例证明了这一点。

      【讨论】:

      • 是否有助于在 groupBy 之前进行一些缓冲?
      • Nope :-) 因为 groupBy 仅在序列状态完成时触发。缓冲并没有完成它。我认为这与 groupBy 的目的非常一致。在实时序列上对数据进行分组没有意义,只有在所有数据都推送完(因此完成)后,您才能进行分组并发出结果
      猜你喜欢
      • 1970-01-01
      • 2021-11-30
      • 1970-01-01
      • 2018-10-24
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-02-18
      相关资源
      最近更新 更多