ReactiveX/RxJava
RxJava – Reactive Extensions for the JVM – a library for composing asynchronous and event-based programs using observable sequences for the Java VM.
About ReactiveX/RxJava
ReactiveX/RxJava is an open-source project on GitHub, mainly written in Java. RxJava – Reactive Extensions for the JVM – a library for composing asynchronous and event-based programs using observable sequences for the Java VM. It currently holds 48,258 stars and 7,585 forks with 11 open issues, and was last pushed on 2026-10-05 (repository created 2013-01-08).
Project Overview
Git Homed tracks it on the Today's Trending board, currently at rank #52 with 68 new stars today.
GitHub Repository Details
README
RxJava: Reactive Extensions for the JVM
:information_source: 4.0.0 release: 2026.11.30. Monday. Track
RxJava is a Java VM implementation of Reactive Extensions: a library for composing asynchronous and event-based programs by using observable sequences.
It extends the observer pattern to support sequences of data/events and adds operators that allow you to compose sequences together declaratively while abstracting away concerns about things like low-level threading, synchronization, thread-safety and concurrent data structures.
Version 4.x (Javadoc)
- :+1: Native Java 26 implementation1.
- :+1: No 3rd party library required at runtime.
- :+1: JPMS and :question: OSGi support still intact.
- :+1:
java.util.concurrent.Flow-based implementation. - :+1: Virtual Thread support;
virtualCreate(),virtualTransform(), :eye:Schedulers.virtual(). - :+1: New
Streamablebuilt around Virtual Threads & virtual blocking. ThinkIAsyncEnumerablefor Java. :satellite: in progress. - :+1: Using Java Cleaner API to detect resource leaks and using it for adaptive cleanups.
- :information_source: Reactive Streams Test Compatibility Kit usage; Reactive-Streams.
- :satellite: Reduce overload bloat by using
record-based configurations. - :satellite: Internal optimizations now that I have the master :key:.
- :eye: Possible usages for Scoped variables for context and per-item resource management.
- :eye: Possible inclusion of 2nd and 3rd party operators.
- :question: Android compatibility depends on your API level and what desugaring is available.
- :warning: RxJava 3.x support will be toned down in the coming months, will be offered for +1 year after 4.x official release.
Getting started
Setting up the dependency
The first step is to include RxJava 4 into your project, for example, as a Gradle compile dependency:
implementation "io.reactivex.rxjava4:rxjava:4.x.y"
(Please replace x and y with the latest version numbers:
)
Hello World
The second is to write the Hello World program:
package rxjava.examples;
import io.reactivex.rxjava4.core.*;
public class HelloWorld {
public static void main(String[] args) {
Flowable.just("Hello world").subscribe(System.out::println);
}
}
Note that RxJava 4 components now live under io.reactivex.rxjava4 and the base classes and interfaces live under io.reactivex.rxjava4.core.
Base classes
RxJava 4 features several base classes you can discover operators on:
io.reactivex.rxjava4.core.Flowable: 0 .. N flows, supporting Reactive-Streams and backpressure,io.reactivex.rxjava4.core.Observable: 0 .. N flows, no backpressure,io.reactivex.rxjava4.core.Single: a flow of exactly 1 item or an error,io.reactivex.rxjava4.core.Completable: a flow without items but only a completion or error signal,io.reactivex.rxjava4.core.Maybe: a flow with no items, exactly one item or an error.io.reactivex.rxjava4.core.Streamable: 0 .. N flows based on virtual blocking andCompletionStage-based state machines with native backpressure.
Some terminology
Upstream, downstream
The dataflows in RxJava consist of a source, zero or more intermediate steps followed by a data consumer or combinator step (where the step is responsible to consume the dataflow by some means):
source.operator1().operator2().operator3().subscribe(consumer);
source.flatMap(value -> source.operator1().operator2().operator3());
Here, if we imagine ourselves on operator2, looking to the left towards the source is called the upstream. Looking to the right towards the subscriber/consumer is called the downstream. This is often more apparent when each element is written on a separate line:
source
.operator1()
.operator2()
.operator3()
.subscribe(consumer)
Objects in motion
In RxJava's documentation, emission, emits, item, event, signal, data and message are considered synonyms and represent the object traveling along the dataflow.
Backpressure
When the dataflow runs through asynchronous steps, each step may perform different things with different speed. To avoid overwhelming such steps, which usually would manifest itself as increased memory usage due to temporary buffering or the need for skipping/dropping data, so-called backpressure is applied, which is a form of flow control where the steps can express how many items are they ready to process. This allows constraining the memory usage of the dataflows in situations where there is generally no way for a step to know how many items the upstream will send to it.
In RxJava, the dedicated Flowable class is designated to support backpressure and Observable is dedicated to the non-backpressured operations (short sequences, GUI interactions, etc.). The other types, Single, Maybe and Completable don't support backpressure nor should they; there is always room to store one item temporarily.
Since 4.0.0, the Streamable type gives natural backpressure because producers and consumers have to wait for
each other to hand over data. Since waiting is blocking, the type natively works with Virtual Threaded ExecutorServices
and the new Schedulers.virtual() Scheduler.
Assembly time
The preparation of dataflows by applying various intermediate operators happens in the so-called assembly time:
Flowable flow = Flowable.range(1, 5)
.map(v -> v * v)
.filter(v -> v % 3 == 0)
;
At this point, the data is not flowing yet and no side effects are happening.
Subscription time
This is a temporary state when subscribe() is called on a flow that establishes the chain of processing steps internally:
flow.subscribe(System.out::println)
`
This is when the subscription side effects are triggered (see doOnSubscribe). Some sources block or start emitting items right away in this state.
Runtime
This is the state when the flows are actively emitting items, errors or completion signals:
Observable.create(emitter -> {
while (!emitter.isDisposed()) {
long time = System.currentTimeMillis();
emitter.onNext(time);
if (time % 2 != 0) {
emitter.onError(new IllegalStateException("Odd millisecond!"));
break;
}
}
})
.subscribe(System.out::println, Throwable::printStackTrace);
Practically, this is when the body of the given example above executes.
Simple background computation
One of the common use cases for RxJava is to run some computation, network request on a background thread and show the results (or error) on the UI thread:
import io.reactivex.rxjava4.schedulers.Schedulers;
Flowable.fromCallable(() -> {
Thread.sleep(1000); // imitate expensive computation
return "Done";
})
.subscribeOn(Schedulers.io())
.observeOn(Schedulers.single())
.subscribe(System.out::println, Throwable::printStackTrace);
Thread.sleep(2000); // <--- wait for the flow to finish
This style of chaining methods is called a fluent API which resembles the builder pattern. However, RxJava's reactive types are immutable; each of the method calls returns a new Flowable with added behavior. To illustrate, the example can be rewritten as follows:
Flowable source = Flowable.fromCallable(() -> {
Thread.sleep(1000); // imitate expensive computation
return "Done";
});
Flowable runBackground = source.subscribeOn(Schedulers.io());
Flowable showForeground = runBackground.observeOn(Schedulers.single());
showForeground.subscribe(System.out::println, Throwable::printStackTrace);
Thread.sleep(2000);
Typically, you can move computations or blocking IO to some other thread via subscribeOn. Once the data is ready, you can make sure they get processed on the foreground or GUI thread via observeOn.
Schedulers
RxJava operators don't work with Threads or ExecutorServices directly but with so-called Schedulers that abstract away sources of concurrency behind a uniform API. RxJava 4 features several standard schedulers accessible via Schedulers utility class.
Schedulers.computation(): Run computation intensive work on a fixed number of dedicated threads in the background. Most asynchronous operators use this as their defaultScheduler.Schedulers.cached(): Run I/O-like or blocking operations on a dynamically changing set of threads backed by native OS threads. :warning: Can exhaust system resources!Schedulers.virtual(): Run I/O-like or blocking scatter-gather operations in a sequential manner on threads with virtualized stacks attached and detached to native OS threads on demand. :information_source: Helps with the issues around unboundedness ofSchedulers.cached().Schedulers.single(): Run work on a single thread in a sequential and FIFO manner.Schedulers.trampoline(): Run work in a sequential and FIFO manner in one of the participating threads, usually for testing purposes.Schedulers.createParallel(): Allows creating aSchedulerwith an user-configurable worker pool size and other parameters to contrastcomputation()which is always set toavailableProcessors()/configured amount globally.Schedulers.createBlocking(): Allows creating an event-loop styleSchedulerwhich runs tasks and blocks on the thread callingexecute(). Can be used to pull tasks onto a specific thread or have it itself run in a virtual threaded executor for maximum efficiency.Scheduler.shared(): EverySchedulerorWorkercan now be shared and act like its own fullSchedulerwith lifecycle tracking and dispose support. I.e., share one worker ofcached()like it is some kind ofSchedulers.single().
Schedulers.io() has been API deprecated and delegates to Schedulers.cached() for compatibility reasons. It is recommended you decide at these deprecated code locations which standard (or custom) scheduler
to use: cached() like before or the new virtual() for more efficient system resource usages.
These are available on all JVM platforms but some specific platforms, such as Android, have their own typical Schedulers defined: AndroidSchedulers.mainThread(), SwingScheduler.instance() or JavaFXScheduler.platform().
In addition, there is an option to wrap an existing Executor (and its subtypes such as ExecutorService) into a Scheduler via Schedulers.from(Executor). This can be used, for example, to have a larger but still fixed pool of threads (unlike computation() and io() respectively).
The Thread.sleep(2000); at the end is no accident. In RxJava the default Schedulers run on daemon threads, which means once the Java main thread exits, they all get stopped and background computations may never happen. Sleeping for some time in this example situations lets you see the output of the flow on the console with time to spare.
Concurrency within a flow
Flows in RxJava are sequential in nature split into processing stages that may run concurrently with each other:
Flowable.range(1, 10)
.observeOn(Schedulers.computation())
.map(v -> v * v)
.blockingSubscribe(System.out::println);
This example flow squares the numbers from 1 to 10 on the computation Scheduler and consumes the results on the "main" thread (more precisely, the caller thread of blockingSubscribe). However, the lambda v -> v * v doesn't run in parallel for this flow; it receives the values 1 to 10 on the same computation thread one after the other.
Parallel processing
Processing the numbers 1 to 10 in parallel is a bit more involved:
Flowable.range(1, 10)
.flatMap(v ->
Flowable.just(v)
.subscribeOn(Schedulers.computation())
.map(w -> w * w)
)
.blockingSubscribe(System.out::println);
Practically, parallelism in RxJava means running independent flows and merging their results back into a single flow. The operator flatMap does this by first mapping each number from 1 to 10 into its own individual Flowable, runs them and merges the computed squares.
Note, however, that flatMap doesn't guarantee any order and the items from the inner flows may end up interleaved. There are alternative operators:
concatMapthat maps and runs one inner flow at a time andconcatMapEagerwhich runs all inner flows "at once" but the output flow will be in the order those inner flows were created.
Flowable.parallel() operator and the ParallelFlowable type help achieve the same parallel processing pattern:
Flowable.range(1, 10)
.parallel()
.runOn(Schedulers.computation())
.map(v -> v * v)
.sequential()
.blockingSubscribe(System.out::println);
Dependent sub-flows
flatMap is a powerful operator and helps in a lot of situations. For example, given a service that returns a Flowable, we'd like to call another service with values emitted by the first service:
Flowable inventorySource = warehouse.getInventoryAsync();
inventorySource
.flatMap(inventoryItem -> erp.getDemandAsync(inventoryItem.getId())
.map(demand -> "Item " + inventoryItem.getName() + " has demand " + demand))
.subscribe(System.out::println);
Continuations
Sometimes, when an item has become available, one would like to perform some dependent computations on it. This is sometimes called continuations and, depending on what should happen and what types are involved, may involve various operators to accomplish.
Dependent
The most typical scenario is given a value, invoke another service, await and continue with its result:
service.apiCall()
.flatMap(value -> service.anotherApiCall(value))
.flatMap(next -> service.finalCall(next))
It is often the case also that later sequences would require values from earlier mappings. This can be achieved by moving the outer flatMap into the inner parts of the previous flatMap for example:
service.apiCall()
.flatMap(value ->
service.anotherApiCall(value)
.flatMap(next -> service.finalCallBoth(value, next))
)
Here, the original value will be available inside the inner flatMap, courtesy of lambda variable capture.
Non-dependent
In other scenarios, the result(s) of the first source/dataflow is irrelevant and one would like to continue with a quasi independent another source. Here, flatMap works as well:
Observable continued = sourceObservable.flatMapSingle(ignored -> someSingleSource)
continued.map(v -> v.toString())
.subscribe(System.out::println, Throwable::printStackTrace);
however, the continuation in this case stays Observable instead of the likely more appropriate Single. (This is understandable because
from the perspective of flatMapSingle, sourceObservable is a multivalued source and thus the mapping may result in multiple values as well).
Often though there is a way that is somewhat more expressive (and also lower overhead) by using Completable as the mediator and its operator andThen to resume with something else:
sourceObservable
.ignoreElements() // returns Completable
.andThen(someSingleSource)
.map(v -> v.toString())
The only dependency between the sourceObservable and the someSingleSource is that the former should complete normally in order for the latter to be consumed.
Deferred-dependent
Sometimes, there is an implicit data dependency between the previous sequence and the new sequence that, for some reason, was not flowing through the "regular channels". One would be inclined to write such continuations as follows:
AtomicInteger count = new AtomicInteger();
Observable.range(1, 10)
.doOnNext(ignored -> count.incrementAndGet())
.ignoreElements()
.andThen(Single.just(count.get()))
.subscribe(System.out::println);
Unfortunately, this prints 0 because Single.just(count.get()) is evaluated at assembly time when the dataflow hasn't even run yet. We need something that defers the evaluation of this Single source until runtime when the main source completes:
AtomicInteger count = new AtomicInteger();
Observable.range(1, 10)
.doOnNext(ignored -> count.incrementAndGet())
.ignoreElements()
.andThen(Single.defer(() -> Single.just(count.get())))
.subscribe(System.out::println);
or
AtomicInteger count = new AtomicInteger();
Observable.range(1, 10)
.doOnNext(ignored -> count.incrementAndGet())
.ignoreElements()
.andThen(Single.fromCallable(() -> count.get()))
.subscribe(System.out::println);
Type conversions
Sometimes, a source or service returns a different type than the flow that is supposed to work with it. For example, in the inventory example above, getDemandAsync could return a Single. If the code example is left unchanged, this will result in a compile-time error (however, often with a misleading error message about lack of overload).
In such situations, there are usually two options to fix the transformation: 1) convert to the desired type or 2) find and use an overload of the specific operator supporting the different type.
Converting to the desired type
Each reactive base class features operators that can perform such conversions, including the protocol conversions, to match some other type. The following matrix shows the available conversion options:
| | Flowable | Observable | Single | Maybe | Completable | Streamable |
|----------|----------|------------|--------|-------|-------------|------------|
|Flowable | | toObservable | first, firstOrError, single, singleOrError, last, lastOrError1 | firstElement, singleElement, lastElement | ignoreElements | toStreamable |
|Observable| toFlowable2 | | first, firstOrError, single, singleOrError, last, lastOrError1 | firstElement, singleElement, lastElement | ignoreElements | toStreamable |
|Single | toFlowable3 | toObservable | | toMaybe | ignoreElement | toStreamable |
|Maybe | toFlowable3 | toObservable | toSingle | | ignoreElement | toStreamable |
|Completable | toFlowable | toObservable | toSingle | toMaybe | | toStreamable |
|Streamable | toFlowable | toObservable | TBD | TBD | TBD | |
1: When turning a multivalued source into a single-valued source, one should decide which of the many source values should be considered as the result.
2: Turning an Observable into Flowable requires an additional decision: what to do with the potential unconstrained flow
of the source Observable? There are several strategies available (such as buffering, dropping, keeping the latest) via the BackpressureStrategy parameter or via standard Flowable operators such as onBackpressureBuffer, onBackpressureDrop, onBackpressureLatest which also
allow further customization of the backpressure behavior.
3: When there is only (at most) one source item, there is no problem with backpressure as it can be always stored until the downstream is ready to consume.
Using an overload with the desired type
Many frequently used operator has overloads that can deal with the other types. These are usually named with the suffix of the target type:
| Operator | Overloads |
|----------|-----------|
| flatMap | flatMapSingle, flatMapMaybe, flatMapCompletable, flatMapIterable |
| concatMap | concatMapSingle, concatMapMaybe, concatMapCompletable, concatMapIterable |
| switchMap | switchMapSingle, switchMapMaybe, switchMapCompletable |
The reason these operators have a suffix instead of simply having the same name with different signature is type erasure. Java doesn't consider signatures such as operator(Function>) and operator(Function>) different (unlike C#) and due to erasure, the two operators would end up as duplicate methods with the same signature.
Operator naming conventions
Naming in programming is one of the hardest things as names are expected to be not long, expressive, capturing and easily memorable. Unfortunately, the target language (and pre-existing conventions) may not give too much help in this regard (unusable keywords, type erasure, type ambiguities, etc.).
Unusable keywords
In the original Rx.NET, the operator that emits a single item and then completes is called Return(T). Since the Java convention is to have a lowercase letter start a method name, this would have been return(T) which is a keyword in Java and thus not available. Therefore, RxJava chose to name this operator just(T). The same limitation exists for the operator Switch, which had to be named switchOnNext. Yet another example is Catch which was named onErrorResumeNext.
Type erasure
Many operators that expect the user to provide some function returning a reactive type can't be overloaded because the type erasure around a Function turns such method signatures into duplicates. RxJava chose to name such operators by appending the type as suffix as well:
Flowable flatMap(Function<? super T, ? extends Publisher<? extends R>> mapper)
Flowable flatMapMaybe(Function<? super T, ? extends MaybeSource<? extends R>> mapper)
Type ambiguities
Even though certain operators have no problems from type erasure, their signature may turn up being ambiguous, especially if one uses Java 8 and lambdas. For example, there are several overloads of concatWith taking the various other reactive base types as arguments (for providing convenience and performance benefits in the underlying implementation):
Flowable concatWith(Publisher<? extends T> other);
Flowable concatWith(SingleSource<? extends T> other);
Both Publisher and SingleSource appear as functional interfaces (types with one abstrac