【问题标题】:onNext of the Subscribe method not emitting items after using the ZIP WITH operator in RxJava?在 RxJava 中使用 ZIP WITH 运算符后,订阅方法的 onNext 不发出项目?
【发布时间】:2018-09-18 15:32:45
【问题描述】:

主要 POJO:

class VideoResponse{
    List<VideoFiles> videosFiles;
}

我有以下情况,我将两个数据库操作的结果合并并返回为 Observable(List(VideoResponse))

##Update##

mDbHelper
  ===/* getVideoCategory() returns Observable<List<VideoResponse>> */=========                 
    .getVideoCategory()  
    .flatMapIterable(videoResponses -> videoResponses)
    .doOnNext(videoResponse -> {
              Timber.d("Flatmap Iterable Thread == > %s", Thread.currentThread().getName());})
    .concatMap((Function<VideoResponse, ObservableSource<VideoResponse>>) videoResponse -> {
                Integer videoId = videoResponse.getId();
                return Observable.fromCallable(() -> videoResponse)
                               .doOnNext(videoResponse1 -> {

                            Timber.d("Thread before ZIP WITH ===>  
                                  %s", Thread.currentThread().getName());
                                })

     ===/* getVideoFileList(int videoId) returns Observable<List<VideoFiles>> */====

                              .zipWith(getVideoFilesList(videoId)),
                                        (videoResponse1, videoFiles) -> {
                                            videoResponse1.setFiles(videoFiles);
                                            return videoResponse1;
                                        })
                              .doOnNext(vResponse -> {

                                 Timber.d("Thread After ZIP WITH ===>  
                                  %s",Thread.currentThread().getName());
                                })
           ======= /*This Gets printed*/ ======================
                           .doOnComplete(()->{
                                Timber.d(" OnComplete Thread for Video Files ===>  %s ",Thread.currentThread().getName());
                            });

                     })
                    .toList()
                    .toObservable()

     ===/* Below Print statement is not getting Executed */=================              
                     .doOnComplete(()->{
                    Timber.d(" Thread doOnComplete");
                })
                    .doOnNext(videoResponses -> {
                        Timber.d("Thread after loading from the LOCAL DB ====> %s", Thread.currentThread().getName());
            }); 

下面是正在执行的调度线程:

 Flatmap Iterable Thread == >  RxCachedThreadScheduler-1
 Thread before ZIP WITH  ===>  RxCachedThreadScheduler-1
 Flatmap Iterable Thread == >  RxCachedThreadScheduler-1
 Thread After ZIP WITH   ===>  RxCachedThreadScheduler-2
 Thread before ZIP WITH  ===>  RxCachedThreadScheduler-2
 Thread After ZIP WITH   ===>  RxCachedThreadScheduler-2

最后的 onNext 永远不会被执行。我需要在 OnNext 中返回 List。 我已将 observeOn 放在不同的位置,但似乎没有任何效果..!!任何建议..

##更新## 使用 SqlBrite,

  @Override
public Observable<List<VideoResponse>> getVideoCategory() {
        return mDBHelper
                .createQuery(VideoEntry.TABLE_NAME,
                        DbUtils.getSelectAllQuery(VideoEntry.TABLE_NAME))
                .mapToOne(DbUtils::videosFromCursor);



  @Override
    public Observable<List<VideoFiles>> getVideoFilesList(int videoId) {
        return mDBHelper.createQuery(VideoDetailsEntry.TABLE_NAME,
                     DbUtils.getSelectFromId(VideoDetailsEntry.TABLE_NAME,VideoDetailsEntry.COLUMN_VIDEO_ID),
                String.valueOf(videoId))
                .mapToOne(DbUtils::videoDetailsFromCursor);
    }

【问题讨论】:

  • 如果上游观察者链永远不会完成,toList() 运算符永远不会完成。您可以使用doOnComplete() 运算符在观察者链中插入日志语句,以查看何时或是否有任何阶段完成
  • @BobDalgleish ... 将 doOnComplete 放在 zipWith 运算符之后(在 getVideoFile(id) 内)打印线程名称..而放置 doOnCompletetoList 无法执行之前...我使用运算符的方式有问题吗..?
  • @BobDalgleish..在将值从一个可观察对象传递到另一个对象后,还有其他方法可以组合可观察对象吗!!
  • 你的例子的逻辑隐藏在细节中。 videoResponse 是一个永远不会完成的可观察对象,或者getVideoFilesList() 永远不会完成,或者两者兼而有之。
  • 如果您只想要来自getVideoCategory() 的第一个值,您可以使用take(1) 运算符——它会在发出第一个值后完成。

标签: android rx-java rx-java2 android-mvp sqlbrite


【解决方案1】:

正如@BobDalgleish 所暗示的, OnComplete 在 toList 之前而不是在 ZipWith 之后被调用。现在我做了以下更改,我得到了完整列表从分贝。我已使用 concatMap 来保留订单并等待完成。

mLocalDataSource
  .getVideoCategory()
  .compose(RxUtils.applySchedulers())
  .flatMap(new Function<List<VideoResponse>, 
              ObservableSource<List<VideoResponse>>>() {
                    @Override
                     public ObservableSource<List<VideoResponse>> apply(List<VideoResponse> videoResponses) throws Exception {
                             return Observable.just(videoResponses)
                                    .concatMap(videoResponses1 -> Observable.fromIterable(videoResponses1)
                                    .concatMap(videoResponse -> Observable.just(videoResponse)
                                    .concatMap(videoResponse1 -> {
                                                Integer videoId = videoResponse1.getId();
                                                return Observable.just(videoResponse1)
                                                      .zipWith(getVideoFilesList(videoId), new BiFunction<VideoResponse, List<VideoFiles>, VideoResponse>() {
                                                                @Override
                                                                public VideoResponse apply(VideoResponse videoResponse1, List<VideoFiles> videoFiles) throws Exception {
                                                                        videoResponse1.setFiles(videoFiles);
                                                                        Timber.d("Video Responses == >  %s",videoResponse1);
                                                                        return videoResponse1;
                                                                        }
                                                                    });
                                                     })))
                                      .toList()
                                      .toObservable()
                                     .observeOn(AndroidSchedulers.mainThread());
                            }
                        }); 

我知道这看起来有点乱!!任何建议或优化,请发布..!!

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多