【问题标题】:Create an Rx.Subject using Subject.create that allows onNext without subscription使用 Subject.create 创建一个 Rx.Subject 允许 onNext 没有订阅
【发布时间】: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


【解决方案1】:

基本解决方案,在订阅前使用 ReplaySubject 缓存 nexts:

我认为你想要做的就是使用ReplaySubject 作为你的观察者。

const { Observable, Subject, ReplaySubject } = Rx;

const replay = new ReplaySubject();

const observable = Observable.create(observer => {
  replay.subscribe(observer);
});

const mySubject = Subject.create(replay, observable);


mySubject.onNext(1);
mySubject.onNext(2);
mySubject.onNext(3);

mySubject.subscribe(x => console.log(x));

mySubject.onNext(4);
mySubject.onNext(5);

结果:

1
2
3
4
5

套接字实现(示例,不要使用)

...但是如果您正在考虑执行 Socket 实现,它会变得更加复杂。这是一个有效的套接字实现,但我不建议您使用它。相反,我建议您使用rxjs-dom(如果您是 RxJS 4 或更低版本)或RxJS 5 中的社区支持的实现之一,这两个我都参与过。

function createSocketSubject(url) {
  let replay = new ReplaySubject();
  let socket;

  const observable = Observable.create(observer => {
    socket = new WebSocket(url);

    socket.onmessage = (e) => {
      observer.onNext(e);
    };

    socket.onerror = (e) => {
      observer.onError(e);
    };

    socket.onclose = (e) => {
      if (e.wasClean) {
        observer.onCompleted();
      } else {
        observer.onError(e);
      }
    }

    let sub;
    socket.onopen = () => {
      sub = replay.subscribe(x => socket.send(x));      
    };
    return () => {
      socket && socket.readyState === 1 && socket.close();
      sub && sub.dispose();
    }
  });

  return Subject.create(replay, observable);
}

const socket = createSocketSubject('ws://echo.websocket.org');

socket.onNext('one');
socket.onNext('two');
socket.subscribe(x => console.log('response: ' + x.data));
socket.onNext('three');
socket.onNext('four');

Here's the obligatory JsBin

【讨论】:

  • 爱你的工作本!我很期待 RxJS 5。这实际上是一个 ZeroMQ 套接字,服务器端(它的接口与 TCP 套接字大致相同)。我从这里删除了其中一些细节。实际代码库:github.com/nteract/enchannel-zmq-backend
  • 谢谢! :) 所以是的,同样的基本原则也适用,真的。但是要记住的事情:如果您希望您的主题是“可重用”或“可重新订阅”,您需要保护 replay 主题免受 onCompleteonError 调用,或者您需要回收它在那些事件中。您很可能想要保护它。取消订阅主题后,它就完成了,您需要重新创建它。每当您准备好将其移至 RxJS 5 时,请在问题中联系我们。
  • 这个项目是上周开始的(好吧,我们使用 zmq 已经有一段时间了——observables 对我们来说是新的),所以我们可以合理地切换到 RxJS 5。想我们会等到你发布,虽然你在这里的回复 + youtube.com/watch?v=QhjALubBQPg 诱惑了我。我担心如果太早采用,我最终会调试新模块而不是我们的应用程序。
  • 现在人们经常会遇到这个问题(大约 1000 次浏览),我很想用我们 now 的方式来更新它,即包装一旦订阅了 observable 就发送:github.com/nteract/nteract/blob/… 对于大多数用例来说,这感觉是一种更好的模式(即使有时我们希望能够在没有显式订阅的情况下发送)。
猜你喜欢
  • 2019-02-28
  • 2017-03-02
  • 2018-07-04
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-09-19
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多