【问题标题】:How to create an observable from pubnub subscribe如何从 pubnub subscribe 创建一个 observable
【发布时间】:2015-04-24 23:02:39
【问题描述】:

我正在努力解决如何使用 rxjs 库将以下内容转换为可观察对象。

var client = PUBNUB.init({
  publish_key: 'pubkey',
  subscribe_key: 'subkey'
});

client.subscribe({
  channel: 'admin:posts',
  message: function(message, env, channel){console.log("Message");},
  connect: function(){console.log("Connected");},
  disconnect: function(){console.log("Disconnected");},
  reconnect: function(){console.log("Reconnected");},
  error: function(){console.log("Network Error");}, 
 });

我希望将消息回调与其他回调一起转换为 observable。

关于如何做到这一点的任何想法?

谢谢

更新 - 2015 年 4 月 29 日

这是我最终用来完成这项工作的。我应该补充一点,只有在用户登录并清理注销后,我才需要订阅 pubnub。

如果这是一个好方法,请告诉我:

var self = this;

var logins = Rx.Observable.create(function (obs) {
  //Using a session manager in ember.
  self.get('session').on('sessionAuthenticationSucceeded', function(e){
    var data = {
      token:this.content.secure.token,
      email:this.content.secure.email
    }

    obs.onNext(data);
  });
  self.get('session').on('sessionAuthenticationFailed', function(e){obs.onError(e)});
  return function(){
    self.get('session').off('sessionAuthenticationSucceeded', function(e){obs.onNext(e)});
    self.get('session').off('sessionAuthenticationFailed', function(e){obs.onError(e)});
  }
});

var logouts = Rx.Observable.create(function (obs) {
  self.get('session').on('sessionInvalidationSucceeded', function(){obs.onNext()});
  self.get('session').on('sessionInvalidationFailed', function(e){obs.onError(e)});
  return function(){
    self.get('session').off('sessionInvalidationSucceeded', function(e){obs.onNext(e)});
    self.get('session').off('sessionInvalidationFailed', function(e){obs.onError(e)});
  }
});

var dataStream = logins
  .map(function(credentials){
    return PUBNUB.init({
      publish_key: 'pub_key',
      subscribe_key: 'sub_key',
      auth_key: credentials.token,
      uuid: credentials.email
    });
  })
  .scan(function(prev, current){
    prev.unsubscribe({
      channel_group:'admin:data'
    });

    return current;
  })
  .concatMap(function(client){
    return Rx.Observable.create(function (observer) {

      client.subscribe({
        channel_group:'admin:data',
        message:function(message, env, channel){
          observer.onNext({message:message, env:env, channel:channel, error:null});
        },
        error:function(error){
          observer.onNext({error:error})
        }
      });

      return function(){
        client.unsubscribe({channel_group:'admin:data'});
      }
    });
  })
  .takeUntil(logouts)
  .publish()
  .refCount();

  var sub1 = dataStream.subscribe(function(data){
    console.log('sub1', data);
  });

  var sub2 = dataStream.subscribe(function(data){
    console.log('sub2', data);
  });

【问题讨论】:

    标签: pubnub rxjs


    【解决方案1】:

    当然,您可以像这样使用 .create 创建源:

    // your client
    var client = PUBNUB.init({
        publish_key: 'pubkey',
        subscribe_key: 'subkey'
    });
    
    var source = Rx.Observable.create(function (observer) {
        client.subscribe({
             channel: 'admin:posts',
             message: function(message){ observer.onNext(message); }, // next item
             error: function(err){ observer.onError(message); },
             // other functions based on your logic, might want handling
        });
        return function(){ // dispose
            // nothing to do here, yet, you might want to define completion too
        }
    });
    

    【讨论】:

    • +1 这真的很接近。可能想在一次性函数中调用 unsubscribe ,并在末尾添加 .singleInstance() 以使 Observable “冷”直到订阅一次,然后“热”直到所有订阅者都取消订阅。此外,可能希望将 PUBNUB.init 延迟到第一次订阅,因为它有点贵,而且现代应用程序通常有足够多的引导程序。
    • 哦,嗨,没想到会在 SO 上见到你 :) 我留下了懒惰的电话并取消订阅 OP 来指定,因为我不确定他的实际用途 - 在我们的代码中我取​​消订阅在 dispose 和处理重新连接/断开连接 - 我只是觉得这是对他的代码的一个强有力的假设,但对于典型的用例来说,这些绝对是好主意。
    • 谢谢大家,我刚刚添加了一个更新,任何关于它的反馈将不胜感激。
    • 顺便说一句,很想看到你所指的完整代码。
    • 我会检查是否有一种方法可以使用自动转换为 observables 并合并它们,而不是使用显式 onNext 和 onError 因为你正在使用事件——除此之外它还可以,可能是 blesh有更多的 cmets :)
    【解决方案2】:

    我没有使用 PUBNUB 的经验,但根据您提供的信息,我可能会这样实现它:

    // your client
    var client = PUBNUB.init({
        publish_key: 'pubkey',
        subscribe_key: 'subkey'
    });
    
    var source = Rx.Observable.defer(function() {
        return Rx.Observable.create(function (observer) {
            client.subscribe({
                channel: 'admin:posts',
                message: function(message){ observer.onNext(message); },
                connect: function(){console.log("Connected");},
                disconnect: function(){console.log("Disconnected");},
                reconnect: function(){console.log("Reconnected");},
                error: function(err){ observer.onError(err); }
            });
    
            return Rx.Disposable.create(function() { // dispose
                client.unsubscribe();
            });
        });
    })
    .publish()
    .refCount();
    

    【讨论】:

      猜你喜欢
      • 2020-10-10
      • 2023-04-03
      • 1970-01-01
      • 2016-12-29
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-10-16
      相关资源
      最近更新 更多