Angular2合并可观察量

Hec*_*tor 8 rxjs typescript angular

我很难找到一些可观察的东西.我似乎无法让两个可观察者放在一起.他们自己工作得很好,但我需要两个值.

db.glass.subscribe( ( glass: GlassData[] ): void => {
  console.log( glass ); // This prints
});

db.cassette_designs.subscribe( ( cassettes: CassetteData[] ): void => {
  console.log( cassettes ): // This prints
});
Run Code Online (Sandbox Code Playgroud)

不是那么熟悉observables,我尝试的第一件事就是将一个嵌套在另一个内部,但内部的东西似乎没有做任何事情.

db.glass.subscribe( ( glass: GlassData[] ): void => {
  console.log( glass ); // This prints

  db.cassette_designs.subscribe( ( cassettes: CassetteData[] ): void => {
    console.log( cassettes ): // This doesn't
  });
});
Run Code Online (Sandbox Code Playgroud)

这似乎有点傻,所以我在Google上搜索是否有更好的方法来组合observables,结果发现有一些.我试过 zip,forkJoin因为他们看起来最喜欢我想做的事,但也没有什么能和他们一起工作.

Observable.zip( db.cassette_designs, db.glass, ( cassettes: CassetteData[], glass: GlassData[] ): void => {
  console.log( cassettes ); // doesn't print
  console.log( glass     ); // doesn't print
});

Observable.forkJoin( [ db.cassette_designs, db.glass ] ).subscribe( ( data: any ): void => {
  console.log( data ); // doesn't print
});
Run Code Online (Sandbox Code Playgroud)

它可能是一件简单的事情,因为我没有正确地调用函数,但我想我会在某些时候得到某种警告或错误.tsc 代码没有任何问题,我在Chrome或Firefox上的开发者控制台中都没有收到任何消息.

更新

我试过了combineLatest,但它仍然没有在控制台中显示任何内容.我一定错过了什么,但我不确定是什么.他们个人工作.

Observable.combineLatest( db.cassette_designs, db.glass, ( cassettes: CassetteData[], glass: GlassData[] ): void => {
  console.log( cassettes ); // doesn't print
  console.log( glass     ); // deson't print
});
Run Code Online (Sandbox Code Playgroud)

可观察量以以下方式创建:

...

public Listen( event: string ): Observable<Response>
{
  return new Observable<Response>( ( subscriber: Subscriber<Response> ): Subscription => {
    const listen_func = ( res: Response ): void => subscriber.next( res );

    this._socket.on( event, listen_func );

    return new Subscription( (): void =>
      this._socket.removeListener( event, listen_func ) );
  });
}

...
Run Code Online (Sandbox Code Playgroud)

然后,为了实际获得可观察量,我发送了一个关于相关事件的回复.例如

...

public cassette_designs: Observable<CassetteData[]>;

...

this.cassette_designs = _socket.Listen( "get_cassette_designs" )
    .map( ( res: Response ) => res.data.data );
Run Code Online (Sandbox Code Playgroud)

Hec*_*tor 9

我设法combineLatest通过实际订阅生成的observable来开始工作.

最初我这样做:

Observable.combineLatest( db.cassette_designs, db.glass, ( cassettes: CassetteData[], glass: GlassData[] ): void => {
  console.log( cassettes );
  console.log( glass     );
});
Run Code Online (Sandbox Code Playgroud)

现在我这样做:

Observable.combineLatest( db.cassette_designs, db.glass ).subscribe( ( data: any[] ): void => {
  console.log( data );
  // cassettes - data[0]
  // glass     - data[1]
});
Run Code Online (Sandbox Code Playgroud)


pau*_*els 6

跟进你的发现:

1)An Observable是一种惰性执行数据类型,这意味着它在订阅之前不会在管道中执行任何操作.这也适用于组合运算符.zip,forkJoin,combineLatest,并且withLatestFrom将所有递归订阅Observables您传递他们只后,他们自己也有订阅.

因此:

var output = Observable.combinelatest(stream1, stream2, (x, y) => ({x, y}));
Run Code Online (Sandbox Code Playgroud)

直到调用实际上并不会做任何事情output.subscribe(),在这一点subscribe也被称为上stream1和stream2,你会得到Rx的所有魔法.

2)更小的一点,但是无论何时开始使用自己的创建方法,首先要查看文档,看它是否已经存在.有用于静态创建方法Arrays,Promises中,节点式的回调和肯定,甚至标准的事件模式.

因此,您的Listen方法可以成为:

public Listen<R>(event: string): Observable<R> {
  return Observable.fromEvent(this._socket, event);
}
Run Code Online (Sandbox Code Playgroud)