【发布时间】:2019-02-28 04:51:38
【问题描述】:
我尝试使用 Rx.Subject 创建一个可被多个组件订阅的 observable。
我有我的 WebsocketService :
export class WebsocketService {
private ws = null;
private subject: Rx.Subject<{}>;
constructor() {
}
public connect(url): Rx.Subject<{}> {
if (!this.subject) {
this.subject = this.create(url);
console.log('Successfully connected: ' + url);
} else {
console.log('Already connected: ' + url);
}
return this.subject;
}
private create(url): Rx.Subject<MessageEvent> {
if (this.ws === null) {
console.log('Connected to WS');
this.ws = new WebSocket(url);
} else {
console.log('Already connected to WS');
}
const observable = Rx.Observable.create(
(obs: Rx.Observer<MessageEvent>) => {
this.ws.onmessage = obs.next.bind(obs);
this.ws.onerror = obs.error.bind(obs);
this.ws.onclose = obs.complete.bind(obs);
return this.ws.close.bind(this.ws);
});
return Rx.Subject.create({}, observable);
}
}
还有两个以相同方式订阅 int 的组件:
// ...
constructor(private stationService: StationService, private websocketService: WebsocketService) {
websocketService.connect('ws://localhost:8080/ws')
.subscribe(msg => {
console.log(msg);
console.log('[Dashboard] Response from websocket: ' + msg.data);
});
}
// ...
第二个:
// ...
constructor(private http: HttpClient, private websocketService: WebsocketService) {
websocketService.connect('ws://localhost:8080/ws')
.subscribe(msg => {
console.log(msg);
console.log('[Station] Response from websocket: ' + msg.data);
});
}
// ...
我们刷新,这两个组件调用我的服务:
连接到 WS websocket.service.ts:26:12 连接成功:ws://localhost:8080/ws websocket.service.ts:17:12 已经连接:ws://localhost:8080/ws
但是当我发送一个套接字时,我只有一个响应的订阅者:
[仪表板] 来自 websocket 的响应:qwerty
有人打电话帮我找出我的错误吗?
谢谢,
【问题讨论】:
-
在为您的应用提供服务后,您的两个组件是否都已正确初始化,以便注册这些订阅?还是只是您在仪表板页面上,因此在您检查控制台时只有仪表板组件被初始化?
-
是的,两者都已初始化,在我的控制台中,我看到“成功连接”和“已连接”,这意味着 WebsocketService.connetc 被调用了 2 次。如果我评论第一个订阅,第二个有效。很奇怪
-
它看起来不错,因为它在这里:stackblitz.com/edit/typescript-omdef8 它应该可以工作。您可以尝试摆脱像这样的自定义可观察对象 const observable = new Subject
(); this.ws.onmessage = (msg) => observable.next(msg); this.ws.onerror = (err) => observable.error(err); this.ws.onclose = () => observable.complete(); });返回可观察的
标签: angular websocket rxjs angular6 rxjs6