创建Obseravble
该章节展示创建响应源的函数,例如Obseravble.
Outline
Just
Available in: Flowable,Observable,Maybe,Single
ReactiveX doumentation: http://reactivex.io/documentation/operators/just.html
构造一个响应类型通过拿一个预先存在的对象且在订阅时将该对象发射给下游消费者。
Just例子:
String greeting = "Hello world!";
Observable<String> observable = Observable.just(greeting);
observable.subscribe(item -> System.out.println(item));
为了方便,这里存在2-9个重载参数,对象(对于同一种类型)将会被以要求的顺序发射。
Observable<Object> observable = Observable.just("1", "A", "3.2", "def");
observable.subscribe(item -> System.out.print(item), error -> error.printStackTrace,
() -> System.out.println());
From
从一个已经存在的源或者生成器类型构造一个序列。
注意:为了避免重载界定模糊,这些静态函数用后缀命名约定(i.e. 参数类型在方法名中重复出现)
ReactiveX doumentation: http://reactivex.io/documentation/operators/from.html
fromIterable
Available in: Flowable,Observable
给java.lang.Iterable的源的item发信号,接着 完成该序列。
fromIterable 例子
List<Integer> list = new ArrayList<>(Arrays.asList(1, 2, 3, 4, 5, 6, 7, 8));
Observable<Integer> observable = Observable.fromIterable(list);
observable.subscribe(item -> System.out.println(item), error -> error.printStackTrace(),
() -> System.out.println("Done"));
fromArray
Available in: Flowable,Observable
给给定的array的元素发信号接着完成该序列。
fromArray 例子:
Integer[] array = new Integer[10];
for (int i = 0; i < array.length; i++) {
array[i] = i;
}
Observable<Integer> observable = Observable.fromIterable(array);
observable.subscribe(item -> System.out.println(item), error -> error.printStackTrace(),
() -> System.out.println("Done"));
注意:RxJava不支持原始数组,仅引用数组。
fromCallable
Available in:Flowable,Observable,Maybe,Single,Completable
当一个消费者订阅的时候,给定的java.util.concurrent.Callable 被调用 且 它的返回值(或者抛出的异常)被转发给该消费者。
fromCallable example:
Callable<String> callable = () -> {
System.out.println("Hello World!");
return "Hello World!");
}
Observable<String> observable = Observable.fromCallable(callable);
observable.subscribe(item -> System.out.println(item), error -> error.printStackTrace(),
() -> System.out.println("Done"));
注意:在 completable,实际的返回值是被忽略且 completable 简单的完成。
fromAction
Available in: Maybe,Completable
当一个消费者订阅的时候,给定的io.reactivex.function.Action 被调用 且消费者 结束或者 收到 Action 抛出的异常。
fromAction example:
Action action = () -> System.out.println("Hello World!");
Completable completable = Completable.fromAction(action);
completable.subscribe(() -> System.out.println("Done"), error -> error.printStackTrace());
注意:fromAction 和fromRunnable 的区别是 Action 的接口允许抛出一个异常 而java.lang.Runnable 不允许。
fromRunnable
Available in: Maybe,Completable
当一个消费者订阅的时候,给定的io.reactivex.function.Action被调用 且消费者结束或者收到Action抛出的异常
fromRunnable example:
Runnable runnable = () -> System.out.println("Hello World!");
Completable completable = Completable.fromRunnable(runnable);
completable.subscribe(() -> System.out.println("Done"), error -> error.printStackTrace());
注意:fromAction和fromRunnable的区别是Action的接口允许抛出一个异常,而java.lang.Runnable不允许。
fromFuture
Available in: Flowable,Observable,Maybe,Single,Completable
给定一个已有的,已经运行或者完成的java.util.concurrent.Future,等待Future 正常完成 或者以阻塞的形式带一个异常 且 转发产生的值和异常给消费者。
fromFuture example:
ScheduledExecutorService executor = Executors.newSingleThreadedScheduledExecutor();
Future<String> future = executor.schedule(() -> "Hello world!", 1, TimeUnit.SECONDS);
Observable<String> observable = Observable.fromFuture(future);
observable.subscribe(
item -> System.out.println(item),
error -> error.printStackTrace(),
() -> System.out.println("Done"));
executor.shutdown();
from{reactive type}
封装或者转换一个响应类型到目标响应类型。
如下的组合是合理的在不同种类的响应类型当中,这些响应类型带有如下签名模式:targetType.from{sourceType}()
注意:不是所有可能的转换是通过from{reactive type} 函数家庭被实现的,为了更进一步的转换可能性,检查 to{reactive type} 函数家庭。
from{reactive type} example:
<Flux<Integer> reactorFlux = Flux.fromCompletionStage(CompletableFuture.<Integer>completedFuture(1));
Observable<Integer> observable = Observable.fromPublisher(reactorFlux);
observable.subscribe(
item -> System.out.println(item),
error -> error.printStackTrace(),
() -> System.out.println("Done"));
Create
Available in: Flowable,Observable,Maybe,Single,Completable
ReactiveX doumentation: http://reactivex.io/documentation/operators/create.html
当被一个消费者订阅的时候,构造一个安全的响应类型实例 ,执行一个用户提供的函数且提供一个特定类型的Emitter 为这个函数 产生 指定的业务逻辑需求的信号。这个函数允许连接非响应的,通常 listener/callback-style 世界 和响应的世界。
create example:
ScheduledExecutorService executor = Executors.newSingleThreadedScheduledExecutor();
ObservableOnSubscribe<String> handler = emitter -> {
Future<Object> future = executor.schedule(() -> {
emitter.onNext("Hello");
emitter.onNext("World");
emitter.onComplete();
return null;
}, 1, TimeUnit.SECONDS);
emitter.setCancellable(() -> future.cancel(false));
};
Observable<String> observable = Observable.create(handler);
observable.subscribe(item -> System.out.println(item), error -> error.printStackTrace(),
() -> System.out.println("Done"));
Thread.sleep(2000);
executor.shutdown();
注意:为了被应用,Flowable.create()必须指定背压行为,当用户提供的函数生成比下游消费者需求更多item。
defer
Available in: Flowable,Observable,Maybe,Single,Completable
ReactiveX doumentation: http://reactivex.io/documentation/operators/defer.html
当一个消费者订阅该响应类型,调用一个用户提供的 java.util.concurrent.Callable,以便Callable能够生成实际的响应实例 为了转发面向消费者的信号。
defer allows:
将每个消费者状态与产生的反应实例相关联
在实际的/生成的 反应实例被订阅前,允许执行副作用。
主要通过让这些热源直到一个消费者订阅才存在来把热源(i.e.
Subject
s andProcessor
s)转换成冷源
defer example:
Observable<Long> observable = Observable.defer(() -> {
long time = System.currentTimeMillis();
return Observable.just(time);
});
observable.subscribe(time -> System.out.println(time));
Thread.sleep(1000);
observable.subscribe(time -> System.out.println(time));
range
vailable in: Flowable,Observable
ReactiveX doumentation: http://reactivex.io/documentation/operators/range.html
对于每一个独立的消费者,生成一系列值。 Range()函数 生成 Integer ,rangeLong() 生成 Long.
range example:
String greeting = "Hello World!";
Observable<Integer> indexes = Observable.range(0, greeting.length());
Observable<Char> characters = indexes
.map(index -> greeting.charAt(index));
characters.subscribe(character -> System.out.print(character), erro -> error.printStackTrace(),
() -> System.out.println());
interval
Available in: Flowable,Observable
ReactiveX doumentation: http://reactivex.io/documentation/operators/interval.html
周期性生成一个无限的、永远增长的数(长整型). IntervalRange 变体生成一个有限数量的数。
interval example:
Observable<Long> clock = Observable.interval(1, TimeUnit.SECONDS);
clock.subscribe(time -> {
if (time % 2 == 0) {
System.out.println("Tick");
} else {
System.out.println("Tock");
}
});
timer
Available in: Flowable,Observable,Maybe,Single,Completable
ReactiveX doumentation: http://reactivex.io/documentation/operators/timer.html
在特定时间后,这个响应源 发出一个 单一的0L信号(接着结束Flowable和Obseravble)
timer example:
Observable<Long> eggTimer = Observable.timer(5, TimeUnit.MINUTES);
eggTimer.blockingSubscribe(v -> System.out.println("Egg is ready!"));
empty
Available in: Flowable,Observable,Maybe,Completable
ReactiveX doumentation: http://reactivex.io/documentation/operators/empty-never-throw.html
这种类型原在订阅时立马发出完成的信号。
empty example:
Observable<String> empty = Observable.empty();
empty.subscribe(
v -> System.out.println("This should never be printed!"),
error -> System.out.println("Or this!"),
() -> System.out.println("Done will be printed."));
never
Available in: Flowable,Observable,Maybe,Single,Completable
ReactiveX doumentation: http://reactivex.io/documentation/operators/empty-never-throw.html
这种类型原不发出onNext,onSuccess,onError或者onComplete信号。 这种类型的响应源 在测试或者在组合操作符中禁用确切的源非常有用。
never example:
Observable<String> never = Observable.never();
never.subscribe(
v -> System.out.println("This should never be printed!"),
error -> System.out.println("Or this!"),
() -> System.out.println("This neither!"));
error
Available in: Flowable,Observable,Maybe,Single,Completable
ReactiveX doumentation: http://reactivex.io/documentation/operators/empty-never-throw.html
对消费者发出一个错误信号,要么已有,要么通过一个java.util.concurrent.Callable 生成。
error example:
Observable<String> error = Observable.error(new IOException());
error.subscribe(
v -> System.out.println("This should never be printed!"),
error -> error.printStackTrace(),
() -> System.out.println("This neither!"));
一个典型的应用案例 是在链中有条件的映射或者利用 onErrorResumeNext抑制异常
Observable<String> observable = Observable.fromCallable(() -> {
if (Math.random() < 0.5) {
throw new IOException();
}
throw new IllegalArgumentException();
});
Observable<String> result = observable.onErrorResumeNext(error -> {
if (error instanceof IllegalArgumentException) {
return Observable.empty();
}
return Observable.error(error);
});
for (int i = 0; i < 10; i++) {
result.subscribe(
v -> System.out.println("This should never be printed!"),
error -> error.printStackTrace(),
() -> System.out.println("Done"));
}