RxJava Subject - Publish, Replay, Behavior, and Async

Authors
RxJava Subject - Publish, Replay, Behavior, and Async

A Subject in RxJava acts both as an Observer and as an Observable. Its four types differ in what a new subscriber gets: Publish gives the items after subscribing, Replay gives all the items, Behavior gives the most recent item and after, and Async gives only the last item after completion.

In this blog, we will learn about the RxJava Subject - Publish, Replay, Behavior, and Async.

I am Amit Shekhar, Founder @ Outcome School, I have taught and mentored many developers, and their efforts landed them high-paying tech jobs, helped many tech companies in solving their unique problems, and created many open-source libraries being used by top companies. I am passionate about sharing knowledge through open-source, blogs, and videos.

I teach AI and Machine Learning, and Android at Outcome School.

Join Outcome School and get a high-paying tech job:

Let's get started.

This article is all about the Subject available in RxJava.

  • Publish Subject
  • Replay Subject
  • Behavior Subject
  • Async Subject

As we already have the sample project based on RxJava2 to learn RxJava (many developers have learned from this sample project), So I have included the Subject examples in the same project. Fork, clone, build, run, and learn RxJava.

Project Link: RxJava2-Android-Samples

What is Subject?

From the official documentation:

A Subject is a sort of bridge or proxy that is available in some implementations of ReactiveX that acts both as an observer and as an Observable. Because it is an observer, it can subscribe to one or more Observables, and because it is an Observable, it can pass through the items it observes by re-emitting them, and it can also emit new items.

I believe: learning by examples is the best way to learn.

Let's learn it through the Professor-Student analogy.

Observable: Assume that a professor is an observable. The professor teaches a topic to the students.

Observer: Assume that a student is an observer. The student observes the topic that is being taught by the professor.

Publish Subject

It emits all the subsequent items of the source Observable at the time of subscription.

Here, if a student entered late into the classroom, he just wants to listen from that point of time when he entered the classroom. So, Publish will be the best for this use case.

See the below example:

PublishSubject<Integer> source = PublishSubject.create();

// It will get 1, 2, 3, 4 and onComplete
source.subscribe(getFirstObserver());

source.onNext(1);
source.onNext(2);
source.onNext(3);

// It will get 4 and onComplete for second observer also.
source.subscribe(getSecondObserver());

source.onNext(4);
source.onComplete();

A quick note for you

No matter which tech domain you work in, get familiar with these topics:

  • LLM
  • RAG
  • MCP
  • Agent
  • Fine-tuning
  • Quantization

We put it all together in one video:

AI Engineering Explained: LLM, RAG, MCP, Agent, Fine-Tuning, and Quantization

No need to stop reading - bookmark it and watch later when you get time. Future you will thank you.

Now, let's get back to the topic.

Replay Subject

It emits all the items of the source Observable, regardless of when the subscriber subscribes.

Here, if a student entered late into the classroom, he wants to listen from the beginning. So, here we will use Replay to achieve this.

See the below example:

ReplaySubject<Integer> source = ReplaySubject.create();

// It will get 1, 2, 3, 4
source.subscribe(getFirstObserver());

source.onNext(1);
source.onNext(2);
source.onNext(3);
source.onNext(4);
source.onComplete();

// It will also get 1, 2, 3, 4 as we have used replay Subject
source.subscribe(getSecondObserver());

Behavior Subject

It emits the most recently emitted item and all the subsequent items of the source Observable when an observer subscribes to it.

Here, if a student entered late into the classroom, he wants to listen to the most recent things(not from the beginning) being taught by the professor so that he gets the idea of the context. So, here we will use Behavior.

See the below example:

BehaviorSubject<Integer> source = BehaviorSubject.create();

// It will get 1, 2, 3, 4 and onComplete
source.subscribe(getFirstObserver());

source.onNext(1);
source.onNext(2);
source.onNext(3);

// It will get 3(last emitted)and 4(subsequent item) and onComplete
source.subscribe(getSecondObserver());

source.onNext(4);
source.onComplete();

Async Subject

It only emits the last value of the source Observable(and only the last value) only after that source Observable completes.

Here, if a student entered the classroom at any point in time, and wants to listen only to the last thing(and only the last thing) being taught after class is over. So, here we will use Async.

See the below example:

AsyncSubject<Integer> source = AsyncSubject.create();

// It will get only 4 and onComplete
source.subscribe(getFirstObserver());

source.onNext(1);
source.onNext(2);
source.onNext(3);

// It will also get only get 4 and onComplete
source.subscribe(getSecondObserver());

source.onNext(4);
source.onComplete();

So, whenever you are stuck with these types of cases, the RxJava Subject will be your best friend.

Stay updated: Subscribe to our newsletter to get our latest AI and Machine Learning blogs straight to your inbox.

Frequently Asked Questions

What is the difference between PublishSubject and BehaviorSubject?

PublishSubject gives a new subscriber only the items emitted after it subscribes. BehaviorSubject also gives the most recently emitted item, followed by all the subsequent items. So a late subscriber to a BehaviorSubject gets some context, while a late subscriber to a PublishSubject does not.

Which Subject should we use if a late subscriber needs all the items?

We should use ReplaySubject. It emits all the items of the source Observable, regardless of when the subscriber subscribes. Even an observer that subscribes after onComplete still gets every item from the beginning.

Does AsyncSubject emit anything before the source completes?

No. AsyncSubject emits only the last value of the source Observable, and only after the source completes. Every observer, whether it subscribed early or late, gets just that last value and onComplete.

Can a Subject subscribe to other Observables?

Yes. Because a Subject is an observer, it can subscribe to one or more Observables. And because it is also an Observable, it can pass through the items it observes by re-emitting them, and it can also emit new items.

That's it for now.

Thanks

Amit Shekhar
Founder @ Outcome School

You can connect with me on:

Follow Outcome School on:

Read all of our high-quality blogs here.