【发布时间】:2016-11-01 16:55:59
【问题描述】:
在使用Rx.Observable.webSocket 公开的Subject 时遇到了一些麻烦。虽然 WebSocket 在complete 之后确实重新连接,但对Subject 的后续订阅也会立即完成,而不是推送通过套接字发送的下一条消息。
我认为我遗漏了一些关于它应该如何工作的基本知识。
这是一个 requirebin/paste,我希望能更好地说明我的意思,以及我所期望的行为。认为这将是我忽略的超级简单的事情。
var Rx = require('rxjs')
var subject = Rx.Observable.webSocket('wss://echo.websocket.org')
subject.next(JSON.stringify('one'))
subject.subscribe(
function (msg) {
console.log('a', msg)
},
null,
function () {
console.log('a complete')
}
)
setTimeout(function () {
subject.complete()
}, 1000)
setTimeout(function () {
subject.next(JSON.stringify('two'))
}, 3000)
setTimeout(function () {
subject.next(JSON.stringify('three'))
subject.subscribe(
function (msg) {
// Was hoping to get 'two' and 'three'
console.log('b', msg)
},
null,
function () {
// Instead, we immediately get here.
console.log('b complete')
}
)
}, 5000)
【问题讨论】:
-
听起来你可能正在处理热与冷的可观察对象。 github.com/Reactive-Extensions/RxJS/blob/master/doc/…