【问题标题】:Convert existing WebSocket implementation into Reactive WebSocket in Android Java将现有的 WebSocket 实现转换为 Android Java 中的 Reactive WebSocket
【发布时间】:2020-08-23 17:26:07
【问题描述】:

我在 Android 中使用 OKHttpClient 和 WebSockets。 我想将其转换为反应式编程。 如何实现这一点。 现在我正在为 WebSocket 连接做这个。

    // OkHttp Client
    OkHttpClient httpClient = new OkHttpClient.Builder().build();

    // WebSocket Object
    Request request = new Request.Builder().url(url).build();
    mWebSocket = httpClient.newWebSocket(request, new WebSocketListener() {
        // Override methods - OnOpen, OnMessage, OnClosing, OnFailure
    }

    // Then Calling this to make websocket request
    mWebSocket.send(message);

我发现这个库使用 Reactive WebSocket https://github.com/jacek-marchwicki/JavaWebsocketClient

但是没有回调监听器作为“WebSocketListener”,所以我可以处理消息。

在这个方向上的任何帮助将不胜感激。谢谢

【问题讨论】:

标签: android websocket rx-java reactive-programming rx-android


【解决方案1】:

Tinder 有改造灵感的 websocket 客户端 https://github.com/Tinder/Scarlet

这里除了发送和接收常规消息之外,您还可以创建一个响应流来接收 Websocket.Event 并过滤传入的事件,如 onOpen、onMessage、onFailed 等。

以下是有关如何执行此操作的示例。转到上面提供的链接以获取详细示例

//Service declaration similar to retrofit
interface MyWebsocketService{
    @Receive
    Flowable<WebSocket.Event> observeWebSocketEvent();
}

Scarlet scarletInstance = new Scarlet.Builder()
    .webSocketFactory(OkHttpClientUtils.newWebSocketFactory(okhttpClient,"websocket-server-url")
    .addStreamAdapterFactory(new RxJava2StreamAdapterFactory())
    .build();

MyWebsocketService myWebsocketService = scarletInstance.create<MyWebsocketService>();

myWebsocketService.observeWebsocketEvent()
                  .subscribe(event -> {
                      if(event instanceof Websocket.Event.OnConnectionOpened){
                         //do something here
                      }else if(event instanceof Websocket.Event.OnConnectionClosed){
                        //do something here
                      }
                  });

【讨论】:

    【解决方案2】:

    我根本没有使用 OkHttp,但简短的回答是在您的 OnMessage 方法中输入 PublishSubject。 Subject 会将此消息提供给您的订阅者,您也可以在订阅前创建一个链。

    // OkHttp Client
    OkHttpClient httpClient = new OkHttpClient.Builder().build();
    
    // Creating a Subject
    // You can use it as Observer or Subscriber
    PublishSubject<String> subject = PublishSubject.create();
    
    // WebSocket Object
    Request request = new Request.Builder().url(url).build();
    mWebSocket = httpClient.newWebSocket(request, new WebSocketListener() {
        @Override
        public void onMessage(String text, ...) {
            subject.onNext(text);
        }
    }
    
    // I hope that is an queue option
    mWebSocket.send(message);
    
    // Do what you want with your message
    Disposable subscription = subject.flatMap(...)
                                     .observeOn(AndroidSchedulers.mainThread())
                                     .subscribe(...);
    

    总之,这是一个复杂的问题,因为您应该处理连接错误、数据压力、配置更改并具有生命周期意识,但它也与非反应式方式(例如回调地狱)相关,因此 PublishSubject 是一个起点。

    【讨论】:

      猜你喜欢
      • 2013-11-28
      • 1970-01-01
      • 1970-01-01
      • 2016-05-27
      • 2012-10-04
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多