【问题标题】:RxJS grouping emitted events nodejsRxJS 对发出的事件进行分组 nodejs
【发布时间】:2016-07-22 12:01:52
【问题描述】:

我正在查询数据库并将结果作为事件“db_row_receieved”的逐行流检索。我正在尝试按公司 ID 对这些结果进行分组,但我在订阅中没有得到任何输出。

db 行格式如下所示。

 // row 1
    {
        companyId: 50,
        value: 200
    }
    // row 2   
    {
        companyId: 50,
        value: 300
    }
    // row 3 
    {
        companyId: 51,
        value: 400
    }

代码:

var source = Rx.Observable.fromEvent(eventEmitter, 'db_row_receieved');
var grouped = source.groupBy((x) => { return x.companyId; });
var selectMany = grouped.selectMany(x => x.reduce((acc, v) => {
                             return acc + v.value;
                          }, 0));

var subscription = selectMany.subscribe(function (obs) {
                        console.log("value: ", obs);
                   }

预期输出:

value: 500    // from the group with companyId 50
value: 400    // from the group with companyId 51

实际输出: 订阅不输出任何内容,但在使用 Rx.Observable.fromArray(someArray) 时有效

谁能告诉我哪里出错了?

【问题讨论】:

  • 您确定eventEmitter 实际上是在发出具有给定名称的事件。其余代码看起来没问题。
  • 你试过 grouped.subscribe(data => { console.log(data); });看看 groupeBy 是否有效
  • @Yury,是的,eventEmitter 确实发出一行。
  • @oliv37 是的,按作品分组。我开始认为这是一个将热门可观察对象分组的问题
  • 你能在你的reduce函数中添加一个console.log吗?例如:{console.log(acc);返回 acc + v.value; }。我认为它不起作用,因为它不知道最终事件是什么。这可以解释为什么它适用于数组

标签: javascript node.js rxjs observable


【解决方案1】:

所以问题是reduce 仅在底层流completed 时才会产生单个值。由于事件发射器是一种无限源,因此它始终处于活动状态。

看看下面的 sn-p - 第一个例子完成,另一个没有。

const data = [
  {k: 'A', v: 1},
  {k: 'B', v: 10},
  {k: 'A', v: 1},
  {k: 'B', v: 10},
  {k: 'A', v: 1},
  {k: 'B', v: 10},
  {k: 'A', v: 1},
  {k: 'A', v: 1},
  {k: 'A', v: 1},
];

Rx.Observable.from(data)
  .concatMap(d => Rx.Observable.of(d).delay(100))
  .groupBy(d => d.k)
  .mergeMap(group => group.reduce((acc, value) => {
    acc.sum += value.v;
    return acc;
  }, {key: group.key, sum: 0}))
  .do(d => console.log('RESULT', d.key, d.sum))
  .subscribe();
  
Rx.Observable.from(data)
  .concatMap(d => Rx.Observable.of(d).delay(100))
  .merge(Rx.Observable.never()) // MERGIN NEVER IN
  // .take(data.length) // UNCOMMENT TO MITIGATE NEVER
  .groupBy(d => d.k)
  .mergeMap(group => group.reduce((acc, value) => {
    acc.sum += value.v;
    return acc;
  }, {key: group.key, sum: 0}))
  .do(d => console.log('RESULT - NEVER - WILL NOT BE PRINTED', d))
  .subscribe();
<script src="https://cdnjs.cloudflare.com/ajax/libs/rxjs/5.0.0-beta.10/Rx.umd.js"></script>

我不知道您的具体用例,但最常见的 2 件事是:

  • 使用scan(可能带有去抖),
  • 如果有指示下级流结束的事件,请使用takeUntil

【讨论】:

    猜你喜欢
    • 2013-08-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-05-25
    • 1970-01-01
    相关资源
    最近更新 更多