Code: https://github.com/ZbCiok/RxJava-examples
See RxJava : https://jreact.com/index.php/reactivex/
Project Structure

CompletableClass.java
public class CompletableClass {
public void complete01() throws InterruptedException {
Disposable d = Completable.complete()
.delay(10, TimeUnit.SECONDS, Schedulers.io())
.subscribeWith(new DisposableCompletableObserver() {
@Override
public void onStart() {
System.out.println("Started");
}
@Override
public void onError(Throwable error) {
error.printStackTrace();
}
@Override
public void onComplete() {
System.out.println("Done!");
}
});
Thread.sleep(5000);
d.dispose();
}
public void complete02() {
Completable
.complete()
.subscribe(new DisposableCompletableObserver() {
@Override
public void onComplete() {
System.out.println("Completed!");
}
@Override
public void onError(Throwable e) {
e.printStackTrace();
}
});
}
}
FlowableClass.java
public class FlowableClass {
public void flowable01() {
Flowable<String> flowable = Flowable.create(emitter -> {
emitter.onNext("Hello");
emitter.onNext("World");
emitter.onComplete();
}, BackpressureStrategy.BUFFER);
}
public void flowable02() {
Flowable<Integer> flowable = Flowable.just(1, 2, 3, 4, 5);
}
public void flowable03() {
List<String> names = Arrays.asList("Alice", "Bob", "Charlie");
Flowable<String> flowable = Flowable.fromIterable(names);
}
public void flowable04() {
Flowable<String> flowable = Flowable.fromIterable(Arrays.asList("one", "three", "two"));
flowable.subscribe(new Subscriber<String>() {
@Override
public void onNext(String item) {
// Handle the emitted item
System.out.println(item);
}
@Override
public void onError(Throwable error) {
// Handle errors
error.printStackTrace();
}
@Override
public void onComplete() {
// Handle completion
System.out.println("Flowable completed");
}
@Override
public void onSubscribe(Subscription subscription) {
subscription.request(Long.MAX_VALUE); // Request all items
}
});
}
}
MaybeClass.java
public class MaybeClass {
public void just01() throws InterruptedException {
Disposable d = Maybe.just("Hello World")
.delay(10, TimeUnit.SECONDS, Schedulers.io())
.subscribeWith(new DisposableMaybeObserver<String>() {
@Override
public void onStart() {
System.out.println("Started");
}
@Override
public void onSuccess(String value) {
System.out.println("Success: " + value);
}
@Override
public void onError(Throwable error) {
error.printStackTrace();
}
@Override
public void onComplete() {
System.out.println("Done!");
}
});
Thread.sleep(5000);
d.dispose();
}
}ObservableClass.java
public class ObservableClass {
public void just01() {
Observable<String> observable = Observable.just("Hello, world");
observable.subscribe(System.out::println);
}
public void fromIterable02() {
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"));
}
public void fromArray03() {
Integer[] array = new Integer[10];
for (int i = 0; i < array.length; i++) {
array[i] = i;
}
Observable<Integer> observable = Observable.fromArray(array);
observable.subscribe(item -> System.out.println(item), error -> error.printStackTrace(),
() -> System.out.println("Done"));
}
public void fromCallable04() {
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"));
}
public void fromFuture05() {
ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor();
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();
}
}
SingleClass.java
public class SingleClass {
public void just01() throws InterruptedException {
Disposable d = Single.just("Hello World")
.delay(10, TimeUnit.SECONDS, Schedulers.io())
.subscribeWith(new DisposableSingleObserver<String>() {
@Override
public void onStart() {
System.out.println("Started");
}
@Override
public void onSuccess(String value) {
System.out.println("Success: " + value);
}
@Override
public void onError(Throwable error) {
error.printStackTrace();
}
});
Thread.sleep(5000);
d.dispose();
}
// see tests
// public void just02() {
// TestSubscriber<String> ts = new TestSubscriber<String>();
// Single.just("A")
// .map(new Function<String, String>() {
// @Override
// public String apply(String s) {
// return s + "B";
// }
// })
// .toFlowable().subscribe(ts);
// //ts.assertValueSequence(Arrays.asList("AB"));
// }
}TargetTypeFromSourceType.java
public class TargetTypeFromSourceType {
// Mono
public void FromReactiveType01() {
Mono<Integer> reactorMono = Mono.fromCompletionStage(CompletableFuture.<Integer>completedFuture(1));
Observable<Integer> observable = Observable.fromPublisher(reactorMono);
observable.subscribe(
item -> System.out.println(item),
error -> error.printStackTrace(),
() -> System.out.println("Done"));
}
}