【问题标题】:RxJava - parallel execution of two API callsRxJava - 两个 API 调用的并行执行
【发布时间】:2017-05-09 01:59:43
【问题描述】:

我有一个 SyncService。

@Override
public int onStartCommand(Intent intent, int flags, final int startId) {
    Timber.i("Starting sync...");
    ...
    RxUtil.unsubscribe(mSubscription);
    mSubscription = mDataManager.syncEvents()
            .subscribeOn(Schedulers.io())
            .subscribe(new Observer<Event>() {
                @Override
                public void onCompleted() {
                    Timber.i("Synced successfully!");
                    stopSelf(startId);
                }

                @Override
                public void onError(Throwable e) {
                    Timber.w(e, "Error syncing.");
                    stopSelf(startId);

                }

                @Override
                public void onNext(Event event) {
                }
            });

    return START_STICKY;
}

Observable events = mDataManager.syncEvents() 是一个 API 调用。

我想做一个并行调用:

单个用户信息 = mDataManager.getUserInfo()

并调用 stopSelf(startId);在这两个调用结束后。

我该怎么做?

我试过RxJava Fetching Observables In Parallel 但这有点不同。

我想我必须使用 .zip 或 .merge 方法。但在我的例子中,一个方法调用返回 Observable(事件列表)和第二个 Single(一个 UserInfo 对象)。

我创建了 z 结果类,它可能是 .zip 方法的结果,但我不知道如何填充它:

public class SyncResponse {
     List<Event> events;
     UserInfo userInfo;
     ...
}

【问题讨论】:

    标签: android rx-java


    【解决方案1】:

    由于您只有两个 observable,您可以使用 android Pair 将它们组合起来以组合结果。

    mDataManager.syncEvents()
        .zipWith(mDataManager.getUserInfo().toObservable(), Pair::create)
        .subscribeOn(Schedulers.io())
        .subscribe(pair -> {
            Event event = pair.first;
            UserInfo userInfo = pair.second;
    
            // your code here
    
            Timber.i("Synced successfully!");
            stopSelf(startId);
        }, throwable -> {
            Timber.w(e, "Error syncing.");
            stopSelf(startId);
        });
    

    没有 java8/retrolmbda 也可以这样做,但为什么

    如果您需要收集所有事件直到完成并将它们与单个用户信息结合起来,这会有点复杂

    Observable<ArrayList<String>> eventsObservable = mDataManager.syncEvents()
        .collect(ArrayList::new, ArrayList::add); 
    
    mDataManager.getUserInfo()
        .zipWith(eventsObservable.toSingle(), Pair::create)
        .subscribeOn(Schedulers.io())
        .subscribe(pair -> {
            UserInfo userInfo = pair.first;
            List<Event> event = pair.second;
    
            // your code here
    
            Timber.i("Synced successfully!");
            stopSelf(startId);
        }, throwable -> {
            Timber.w(e, "Error syncing.");
            stopSelf(startId);
        });
    

    【讨论】:

      【解决方案2】:

      您可以使用 observable.zip 和 single.toObservable 来完成此操作。我相信您还需要在每个 observable 的 zip 调用中分别进行线程调度,以并行执行它们。

      Observable.zip(call1.subscribeOn(Schedulers.io()), call2.toObservable().subscribeOn(Schedulers.io(), zip-function)
      

      因为两个调用都是不同的包装类型。

      zip 函数是您可以使用类来组合结果的地方。

      应该注意,由于您的一个 observables 是一个单一的,这种带有 zip 的方法最多只会产生一个结果,因为单一将在其第一个 onNext 事件之后立即发送一个 onComplete 事件。

      如果您知道您的事件列表将完成,您可以在 call1 上使用 toList 来缓冲它们并将它们作为单个收集的事件列表发出。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2020-05-13
        • 1970-01-01
        • 2018-04-10
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多