RxJava2 / RxAndroid2的concat拼接多个Observable

简介: RxJava2 / RxAndroid2的concat拼接多个Observable concat操作符和merge类似,把多个Observable拼接成一个可以观察的输出,例如代码: package zhangphil.

RxJava2 / RxAndroid2的concat拼接多个Observable

 

concat操作符和merge类似,把多个Observable拼接成一个可以观察的输出,例如代码:

 

package zhangphil.app;

import android.os.Bundle;
import android.support.annotation.NonNull;
import android.support.annotation.Nullable;
import android.support.v7.app.AppCompatActivity;
import android.util.Log;

import java.util.concurrent.Callable;

import io.reactivex.Observable;
import io.reactivex.android.schedulers.AndroidSchedulers;
import io.reactivex.disposables.CompositeDisposable;
import io.reactivex.functions.BiFunction;
import io.reactivex.observers.DisposableObserver;
import io.reactivex.schedulers.Schedulers;

public class MainActivity extends AppCompatActivity {
    private final String TAG = getClass().getSimpleName();
    private CompositeDisposable mCompositeDisposable = new CompositeDisposable();

    @Override
    public void onCreate(@Nullable Bundle savedInstanceState) {
        super.onCreate(savedInstanceState);

        test();
    }

    private void test() {

        DisposableObserver disposableObserver = new DisposableObserver<String>() {
            @Override
            public void onNext(String s) {
                Log.d(TAG, "#####开始#####");
                Log.d(TAG + "数据", String.valueOf(s));
                Log.d(TAG, "#####结束#####");
            }

            @Override
            public void onComplete() {
                Log.d(TAG, "onComplete");
            }

            @Override
            public void onError(Throwable e) {
                Log.e(TAG, e.toString(), e);
            }
        };

        mCompositeDisposable.add(
                Observable.concat(
                        getObservableA(null),
                        getObservableB(null),
                        getObservableA(null),
                        getObservableB(null))
                        .subscribeOn(Schedulers.io())
                        .observeOn(AndroidSchedulers.mainThread())
                        .subscribeWith(disposableObserver));
    }

    @Override
    protected void onDestroy() {
        super.onDestroy();

        // 如果退出程序,就清除后台任务
        mCompositeDisposable.clear();
    }

    private Observable<String> getObservableA(Object o) {
        return Observable.fromCallable(new Callable<String>() {
            @Override
            public String call() throws Exception {
                try {
                    Thread.sleep(500); // 假设此处是耗时操作
                } catch (Exception e) {
                    e.printStackTrace();
                }

                return "A";
            }
        });
    }

    private Observable<String> getObservableB(Object o) {
        return Observable.fromCallable(new Callable<String>() {
            @Override
            public String call() throws Exception {
                try {
                    Thread.sleep(1000); // 假设此处是耗时操作
                } catch (Exception e) {
                    e.printStackTrace();
                }

                return "B";
            }
        });
    }
}
AI 代码解读

 

输出:

 

05-15 14:39:18.667 14456-14456/zhangphil.app D/MainActivity: #####开始#####
05-15 14:39:18.667 14456-14456/zhangphil.app D/MainActivity数据: A
05-15 14:39:18.667 14456-14456/zhangphil.app D/MainActivity: #####结束#####
05-15 14:39:19.669 14456-14456/zhangphil.app D/MainActivity: #####开始#####
05-15 14:39:19.669 14456-14456/zhangphil.app D/MainActivity数据: B
05-15 14:39:19.669 14456-14456/zhangphil.app D/MainActivity: #####结束#####
05-15 14:39:20.170 14456-14456/zhangphil.app D/MainActivity: #####开始#####
05-15 14:39:20.170 14456-14456/zhangphil.app D/MainActivity数据: A
05-15 14:39:20.170 14456-14456/zhangphil.app D/MainActivity: #####结束#####
05-15 14:39:21.171 14456-14456/zhangphil.app D/MainActivity: #####开始#####
05-15 14:39:21.172 14456-14456/zhangphil.app D/MainActivity数据: B
05-15 14:39:21.172 14456-14456/zhangphil.app D/MainActivity: #####结束#####
05-15 14:39:21.172 14456-14456/zhangphil.app D/MainActivity: onComplete
AI 代码解读

 
目录
打赏
0
0
0
0
15
分享
相关文章
RxJava之Hello, World
RxJava之Hello, World
6. Observable 和 数组的区别
Observable 和 数组都有filter, map 等运算操作operators,具体的区别是什么? 主要是两点: 延迟运算 渐进式取值 延迟运算 延迟运算很好理解,所有 Observable 一定会等到订阅后才开始对元素做运算,如果没有订阅就不会有运算的行为 var source = Rx.
873 0
RxJava2-map操作符源码解析
RxJava2的map操作符用于对输入对象进行转换。 map操作图 下图所示为将String的输出转化为Integer的场景。 String转Integer Map的源码解析如下,首先涉及到以下几个类: 1、Observable:被观察者,通过Observable.create创建一个被观察者,即观察者模式里面的主题Subject对象。
1099 0
RxJava操作符大全
RxJava操作符大全 创建操作 以下操作符用于创建Observable。 create: 使用OnSubscribe从头创建一个Observable,这种方法比较简单。需要注意的是,使用该方法创建时,建议在OnSubscribe#call方法中检查订阅状态,以便及时停止发射数据或者运算。
1513 0
AI助理

你好,我是AI助理

可以解答问题、推荐解决方案等