重新连接Angular和rxjs中的websocket?

stw*_*sel 12 websocket rxjs typescript ngrx angular

我有一个基于ngrx/store(v2.2.2)和rxjs(v5.1.0)的应用程序,它使用observable监听Web套接字的传入数据.当我启动应用程序时,我可以完美地接收传入的数据.

但是过了一段时间(更新很少发生)连接似乎丢失了,我不再收到传入的数据.我的代码:

服务

import { Injectable, OnInit } from '@angular/core';
import { Observable } from 'rxjs';

@Injectable()
export class MemberService implements OnInit {

  private websocket: any;
  private destination: string = "wss://notessensei.mybluemix.net/ws/time";

  constructor() { }

  ngOnInit() { }

  listenToTheSocket(): Observable<any> {

    this.websocket = new WebSocket(this.destination);

    this.websocket.onopen = () => {
      console.log("WebService Connected to " + this.destination);
    }

    return Observable.create(observer => {
      this.websocket.onmessage = (evt) => {
        observer.next(evt);
      };
    })
      .map(res => res.data)
      .share();
  }
}
Run Code Online (Sandbox Code Playgroud)

订户

  export class AppComponent implements OnInit {

  constructor(/*private store: Store<fromRoot.State>,*/ private memberService: MemberService) {}

  ngOnInit() {
    this.memberService.listenToTheSocket().subscribe((result: any) => {
      try {
        console.log(result);
        // const member: any = JSON.parse(result);
        // this.store.dispatch(new MemberActions.AddMember(member));
      } catch (ex) {
        console.log(JSON.stringify(ex));
      }
    })
  }
}
Run Code Online (Sandbox Code Playgroud)

在超时时重新连接Web套接字需要做什么,所以observable会继续发出传入值?

我在这里,这里这里看了一些问答,似乎没有解决这个问题(以我能理解的方式).

注意:websocket at wss://notessensei.mybluemix.net/ws/time是live并且每分钟发出一次时间戳(如果有人想测试那个).

建议非常感谢!

Her*_*man 25

实际上现在rxjs中有一个WebsocketSubject!

 import { webSocket } from 'rxjs/webSocket' // for RxJS 6, for v5 use Observable.webSocket

 let subject = webSocket('ws://localhost:8081');
 subject.subscribe(
    (msg) => console.log('message received: ' + msg),
    (err) => console.log(err),
    () => console.log('complete')
  );
 subject.next(JSON.stringify({ op: 'hello' }));
Run Code Online (Sandbox Code Playgroud)

当您重新订阅断开的连接时,它会处理重新连接.所以例如写这个重新连接:

subject.retry().subscribe(...)
Run Code Online (Sandbox Code Playgroud)

有关详细信息,请参阅文档.不幸的是,搜索框没有显示该方法,但您可以在此处找到它:

http://reactivex.io/rxjs/class/es6/Observable.js~Observable.html#static-method-webSocket

#-navigation在我的浏览器中无效,因此在该页面上搜索"webSocket".

资料来源:http://reactivex.io/rxjs/file/es6/observable/dom/WebSocketSubject.js.html#lineNumber15

  • 要在 rxjs 6 中导入它,请使用“import {webSocket} from 'rxjs/webSocket'” (2认同)
  • 类型“WebSocketSubject”上不存在属性“重试”? (2认同)

Joh*_*uez 5

对于 rxjs 6 实现

import { webSocket } from 'rxjs/webSocket'
import { retry, RetryConfig } from "rxjs/operators";

const retryConfig: RetryConfig = {
  delay: 3000,
};

let subject = webSocket('ws://localhost:8081');
subject.pipe(
   retry(retryConfig) //support auto reconnect
).subscribe(...)
Run Code Online (Sandbox Code Playgroud)