【问题标题】:RxJS #zip groups created using #groupBy使用 #groupBy 创建的 RxJS #zip 组
【发布时间】:2016-01-16 13:02:51
【问题描述】:

我需要压缩分组的 observables(以形成相关组的笛卡尔积,但这与问题无关)。

当运行以下代码时,只有子可观察组实际上在#zip 中发出值 - 为什么会这样?

https://jsbin.com/coqeqaxoci/edit?js,console

var parent = Rx.Observable.from([1,2,3]).publish();
var child = parent.map(x => x).publish();
var groupedParent = parent.groupBy(x => x);
var groupedChild = child.groupBy(x => x);

Rx.Observable.zip([groupedChild, groupedParent])
  .map(groups => {
    groups[0].subscribe(x => console.log('zipped child ' + x)); // -> emitting
    groups[1].subscribe(x => console.log('zipped parent ' + x)); // -> not emitting
  })
  .subscribe();

groupedChild.subscribe(group => {
  group.subscribe(value => console.log('child ' + value)); // -> emitting
});

groupedParent.subscribe(group => {
  group.subscribe(value => console.log('parent ' + value)); // -> emitting
});

child.connect();
parent.connect();

编辑: 正如 user3743222 的回答中所解释的,groupBy 发出的组是hot,并且对父组 (groups[1]) 的订阅发生在第一个值已经发出之后。这发生在 #zip 等待 groupedChild 和 groupedParent 发出,后者发出更快(意味着它的组在 #zip 函数运行之前发出值)。

【问题讨论】:

    标签: javascript stream rxjs cartesian-product


    【解决方案1】:

    我修改了你的代码如下:

    var countChild = 0, countParent = 0;
    function emits ( who ) {
      return function ( x ) {console.log(who + " emits : " + x);};
    }
    function checkCount ( who ) {
      return function ( ) {
        if (who === "parent") {
          countParent++;
        }
        else {
          countChild++;
        }
        console.log("Check : Parent groups = " + countParent + ", Child groups = " + countChild );
      };
    }
    function check ( who, where ) {
      return function ( x ) {
        console.log("Check : " + who + " : " + where + " :" + x);
      };
    }
    function completed ( who ) {
      return function () { console.log(who + " completed!");};
    }
    function zipped ( who ) {
      return function ( x ) { console.log('zipped ' + who + ' ' + x); };
    }
    function plus1 ( x ) {
      return x + 1;
    }
    function err () {
      console.log('error');
    }
    
    var parent = Rx.Observable.from([1, 2, 3, 4, 5, 6])
        .do(emits("parent"))
        .publish();
    var child = parent
        .map(function ( x ) {return x;})
        .do(emits("child"))
    //    .publish();
    
    var groupedParent = parent
        .groupBy(function ( x ) { return x % 2;}, function ( x ) {return "P" + x;})
        .do(checkCount("parent"))
        .share();
    
    var groupedChild = child
        .groupBy(function ( x ) { return x % 3;}, function (x) {return "C" + x;})
        .do(checkCount("child"))
        .share();
    
    Rx.Observable.zip([groupedChild, groupedParent])
    //    .do(function ( x ) { console.log("zip args : " + x);})
        .subscribe(function ( groups ) {
                     groups[0]
                         .do(function ( x ) { console.log("Child group observable emits : " + x);})
                         .subscribe(zipped('child'), err, completed('Child Group Observable'));
                     groups[1]
                         .do(function ( x ) { console.log("Parent group observable emits : " + x);})
                         .subscribe(zipped('parent'), err, completed('Parent Group Observable'));
                   }, err, completed('zip'));
    
    //child.connect();
    parent.connect();
    

    这是输出:

    "parent emits : 1"
    "child emits : 1"
    "Check : Parent groups = 0, Child groups = 1"
    "Check : Parent groups = 1, Child groups = 1"
    "Parent group observable emits : P1"
    "zipped parent P1"
    "parent emits : 2"
    "child emits : 2"
    "Check : Parent groups = 1, Child groups = 2"
    "Check : Parent groups = 2, Child groups = 2"
    "Parent group observable emits : P2"
    "zipped parent P2"
    "parent emits : 3"
    "child emits : 3"
    "Check : Parent groups = 2, Child groups = 3"
    "Parent group observable emits : P3"
    "zipped parent P3"
    "parent emits : 4"
    "child emits : 4"
    "Child group observable emits : C4"
    "zipped child C4"
    "Parent group observable emits : P4"
    "zipped parent P4"
    "parent emits : 5"
    "child emits : 5"
    "Child group observable emits : C5"
    "zipped child C5"
    "Parent group observable emits : P5"
    "zipped parent P5"
    "parent emits : 6"
    "child emits : 6"
    "Parent group observable emits : P6"
    "zipped parent P6"
    "Child Group Observable completed!"
    "Child Group Observable completed!"
    "Parent Group Observable completed!"
    "Parent Group Observable completed!"
    "zip completed!"
    

    这里要说明两点:

    1. zip 和 group by 与订阅时刻的行为

      • groupBy 在父子节点和子节点中按预期创建可观察对象

      使用这些值,您可以在日志中查看 Child 创建三个组,Parent 创建两个组

      • Zip 将等待您作为参数传递的每个源中都有一个值。在您的情况下,这意味着您将订阅按可观察对象分组的子对象和父对象,当它们都已发布时。在日志中,只有在匹配"Check : Parent groups = 1, Child groups = 1" 上的数字后,您才会看到"Parent group observable emits : P1"

      • 然后您订阅这两个分组的 observables,并记录从那里出来的任何内容。这里的问题是父 grouped-by observable 有一个值要传递,但是子 'group-by' observable 之前创建并且已经传递了它的值,所以当你在事后订阅时,你看不到那个值- 但你会看到下一个。

      • 因此,[1-3] 中的值将生成 3 个新的子可观察对象分组,并且您将看不到任何这些,因为您订阅得太晚了。但是您会在[4-6] 中看到值。您可以查看日志:"zipped child C4" 等。

      • 您将看到父级可观察对象中的所有值,因为您在创建它们后立即订阅它们。

    2. 连接和发布

      • 我对连接和发布没有完全清楚的了解,但由于您的孩子有父母作为来源,您不需要延迟连接到它。如果您连接到父级,子级将自动开始发出其值。因此我对您的代码进行了修改。

      • 这应该回答您的直接问题,但不是您最初的笛卡尔积目标。也许您应该将其表述为一个问题,然后看看人们能给出什么答案。

    【讨论】:

    • 谢谢!你能解释为什么在 #zip 中订阅时,子 observable 已经传递了它的值吗?我不明白这是怎么发生的...关于 2.: 在我给出的示例中不需要发布,我只是将它包含在完整的脚本中。
    • groupBy 链接到 groupByUntil : https://github.com/Reactive-Extensions/RxJS/blob/master/src/core/linq/observable/groupbyuntil.js 每次有一个新的键(即组)时,都会创建一个可观察对象。但是那个 observable 是一个Rx.Subject,它也是一个热源:它立即发出它的值(参见 L35)。该 observable 被打包并发送给观察者(L41,L49)。然后,该键(组)的值(立即)通过主题(L75)发出。简而言之,如果您在传递该 observable 时不立即订阅它,您将始终丢失第一个值。
    • 如果你不熟悉冷热可观察对象,这里有两个资源处理这个问题:github.com/Reactive-Extensions/RxJS/blob/master/doc/…; jaredforsyth.com/2015/03/06/…
    • 再次感谢 :) 在这种情况下,我看到的唯一方法是使用 .toArray() 将组转换为数组(因为它仅在底层组完成后才发出)-否则我将永远丢失值组合组时,对吗?我也试过 combineLatest...
    • 我认为你应该解释你试图解决的问题(笛卡尔积)而不是试图解释解决方案。例如:CARTESIAN_PRODUCT :: INPUT1 -> INPUT2 -> OUTPUT,描述您的输入和预期输出,使用一些示例输入和输出更好。也就是说,如果你找到了一个可行的解决方案,当然没必要。
    猜你喜欢
    • 1970-01-01
    • 2023-04-09
    • 2018-06-19
    • 1970-01-01
    • 2021-12-11
    • 1970-01-01
    • 2016-12-28
    • 1970-01-01
    • 2018-10-24
    相关资源
    最近更新 更多