【问题标题】:Rx Java Android : How to convert this callback block to ObserverRx Java Android:如何将此回调块转换为观察者
【发布时间】:2017-04-01 03:44:04
【问题描述】:

我正在尝试通过 Amazon 的 S3 Android SDK 上传文件。我已经使用了一点 RX Java,但我不确定如何将此方法转换为返回 Observable 的方法,因为我想将此方法的结果链接到另一个 Observable 调用。我想这让我很困惑,因为它不会立即返回并且在 OnError 或 OnState 更改之前无法返回。如何以 RX 方式处理这些情况?

public void uploadFile(TransferObserver transferObserver){

    transferObserver.setTransferListener(new TransferListener() {
        @Override
        public void onStateChanged(int id, TransferState state) {

        }

        @Override
        public void onProgressChanged(int id, long bytesCurrent, long bytesTotal) {

        }

        @Override
        public void onError(int id, Exception ex) {

        }
  });

}

如果有人能用 RX Java 2 和 lambdas 来回答,那就太好了,因为我只是在这方面做得不够

【问题讨论】:

标签: callback rx-java rx-android


【解决方案1】:

这通常是在异步/回调工作与反应式之间架起桥梁的正确方法,但现在不鼓励使用 Observable.create(),因为它需要高级知识才能使其正确。
您应该使用更新的创建方法Observable.fromEmitter(),看起来完全一样:

    return Observable.fromEmitter(new Action1<Emitter<Integer>>() {
        @Override
        public void call(Emitter<Integer> emitter) {

            transObs.setTransferListener(new TransferListener() {
                @Override
                public void onStateChanged(int id, TransferState state) {
                    if (state == TransferState.COMPLETED)
                        emitter.onCompleted();
                }

                @Override
                public void onProgressChanged(int id, long bytesCurrent, long bytesTotal) {

                }

                @Override
                public void onError(int id, Exception ex) {
                    emitter.onError(ex);
                }
            });
            emitter.setCancellation(new Cancellable() {
                @Override
                public void cancel() throws Exception {
                    // Deal with unsubscription:
                    // 1. unregister the listener to avoid memory leak
                    // 2. cancel the upload 
                }
            });
        }
    }, Emitter.BackpressureMode.DROP);

这里添加的是:处理取消订阅:取消上传,取消注册以避免内存泄漏,并指定背压策略。
你可以阅读更多here

补充说明:

  • 如果您对进度感兴趣,可以在onProgressChanged() 上调用 onNext() 并将 Observable 转换为 Observable&lt;Integer&gt;
  • 如果没有,您可能需要考虑使用 Completable,它是 Observable,没有 onNext() 排放,但只有 onCompleted(),如果您对进度指示不感兴趣,这可能适合您的情况。

【讨论】:

    【解决方案2】:

    @Yosriz 我无法编译你的代码,但你确实帮了我很多,所以根据你的回答,这就是我现在所拥有的:

    return Observable.fromEmitter(new Action1<AsyncEmitter<Integer>>() {
                @Override
                public void call(AsyncEmitter<Integer> emitter) {
    
                    transObs.setTransferListener(new TransferListener() {
                        @Override
                        public void onStateChanged(int id, TransferState state) {
                            if (state == TransferState.COMPLETED)
                                emitter.onCompleted();
                        }
    
                        @Override
                        public void onProgressChanged(int id, long bytesCurrent, long bytesTotal) {
    
                        }
    
                        @Override
                        public void onError(int id, Exception ex) {
                            emitter.onError(ex);
                        }
                    });
    
                    emitter.setCancellation(new AsyncEmitter.Cancellable() {
                        @Override
                        public void cancel() throws Exception {
    
                            transObs.cleanTransferListener();
                        }
                    });
                }
            }, AsyncEmitter.BackpressureMode.DROP);
    

    【讨论】:

    • 需要在dispose()方法中释放transObs的Listener。而且我想最好使用 Completable,因为您不关心 onNext()s。
    猜你喜欢
    • 1970-01-01
    • 2017-02-01
    • 2010-12-04
    • 1970-01-01
    • 2020-05-31
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-07-10
    相关资源
    最近更新 更多