【发布时间】:2016-04-23 02:17:56
【问题描述】:
当使用Subject.create(observer, observable) 创建Rx.Subject 时,Subject 太懒了。当我尝试在没有订阅的情况下使用subject.onNext 时,它不会传递消息。如果我先subject.subscribe(),我可以在之后立即使用onNext。
假设我有一个Observer,这样创建:
function createObserver(socket) {
return Observer.create(msg => {
socket.send(msg);
}, err => {
console.error(err);
}, () => {
socket.removeAllListeners();
socket.close();
});
}
然后,我创建一个接受消息的 Observable:
function createObservable(socket) {
return Observable.fromEvent(socket, 'message')
.map(msg => {
// Trim out unnecessary data for subscribers
delete msg.blobs;
// Deep freeze the message
Object.freeze(msg);
return msg;
})
.publish()
.refCount();
}
主题是使用这两个函数创建的。
observer = createObserver(socket);
observable = createObservable(socket);
subject = Subject.create(observer, observable);
使用此设置,我无法立即subject.onNext(即使我不关心订阅)。这是设计使然吗?有什么好的解决方法?
这些实际上是 TCP 套接字,这就是为什么我没有依赖超级光滑的 websocket 主题。
【问题讨论】:
-
你能不能描述得更详细一些,也许有一个代码示例?不将消息传递到哪里?你想用
Subject.create()方法完成什么?
标签: rxjs