RxJava中Observable与Observer的概念、工作原理及触发问题咨询
Hey there! Let’s clear up your confusion about RxJava’s Observable and Observer, and figure out why modifying your list isn’t triggering the Observer like you expected.
First, let’s get the basics straight—these two have very distinct roles:
- Observable: Think of this as a data source that emits a sequence of items (or events) over time. Once it finishes emitting all its items, throws an error, or signals completion, that’s the end of its lifecycle for that subscription. It doesn’t "watch" for changes to the data it originally emitted.
- Observer: This is the consumer that listens to the Observable’s emissions. It reacts to three key events:
onNext()(receives an item),onError()(handles a failure), andonComplete()(gets notified when the Observable has no more items to emit).
When you create an Observable from a List (using something like Observable.fromIterable(list)), here’s exactly what happens:
- The Observable takes a snapshot of the list’s current state at the time you subscribe.
- It immediately emits every item in that snapshot one by one.
- Once all items are emitted, it calls
onComplete()to let the Observer know it’s done.
Modifying the list later (like adding a new string) doesn’t do anything because the Observable has already finished its emission sequence. It’s not monitoring the list for updates—it just did its job and moved on.
Here’s a quick code example to illustrate this:
List<String> list = new ArrayList<>(Arrays.asList("A", "B")); Observable<String> observable = Observable.fromIterable(list); observable.subscribe(new Observer<String>() { @Override public void onSubscribe(Disposable d) {} @Override public void onNext(String s) { System.out.println("Received: " + s); } @Override public void onError(Throwable e) {} @Override public void onComplete() { System.out.println("Observable finished emitting items"); } }); // This won't trigger the Observer at all list.add("C");
When you run this, you’ll see "Received: A", "Received: B", then "Observable finished emitting items"—but never "Received: C".
If you want your Observer to react when the list is modified, you need a way to emit new events whenever the list updates. Here are two common, practical approaches:
1. Use a Subject (e.g., PublishSubject)
Subjects act as both an Observable and an Observer—they let you manually emit new items whenever you want. This is the simplest way to handle dynamic list updates.
Option 1: Emit the entire updated list
// Create a Subject that emits List<String> objects PublishSubject<List<String>> listSubject = PublishSubject.create(); // Subscribe to the Subject listSubject.subscribe(new Observer<List<String>>() { @Override public void onSubscribe(Disposable d) {} @Override public void onNext(List<String> updatedList) { System.out.println("Updated list: " + updatedList); } @Override public void onError(Throwable e) {} @Override public void onComplete() {} }); // Initial list emission List<String> list = new ArrayList<>(Arrays.asList("A", "B")); listSubject.onNext(list); // Observer sees [A, B] // After modifying, emit the updated list again list.add("C"); listSubject.onNext(list); // Observer sees [A, B, C]
Option 2: Emit individual new items
If you don’t want to send the whole list every time, just emit the new item directly:
PublishSubject<String> itemSubject = PublishSubject.create(); // Subscribe to receive individual items itemSubject.subscribe(item -> System.out.println("New item received: " + item)); List<String> list = new ArrayList<>(Arrays.asList("A", "B")); // Emit initial items list.forEach(itemSubject::onNext); // Observer gets A, then B // Add a new item and emit it String newItem = "C"; list.add(newItem); itemSubject.onNext(newItem); // Observer gets C
2. Create a Custom Observable List
For a more integrated solution, you can wrap your list in a custom class that notifies listeners when items are added/removed. Then, create an Observable that hooks into those notifications:
// A custom ArrayList that triggers listeners when items are added class ObservableList<T> extends ArrayList<T> { private final List<Consumer<T>> updateListeners = new ArrayList<>(); public void addUpdateListener(Consumer<T> listener) { updateListeners.add(listener); } @Override public boolean add(T item) { boolean added = super.add(item); // Notify all listeners about the new item updateListeners.forEach(listener -> listener.accept(item)); return added; } } // Create our observable list ObservableList<String> observableList = new ObservableList<>(); // Build an Observable that emits items when the list updates Observable<String> listObservable = Observable.create(emitter -> { // Emit all existing items first observableList.forEach(emitter::onNext); // Add a listener to emit new items as they're added observableList.addUpdateListener(emitter::onNext); // Clean up the listener if the subscription is disposed emitter.setCancellable(() -> observableList.updateListeners.clear()); }); // Subscribe to the Observable listObservable.subscribe(item -> System.out.println("Received: " + item)); // Add initial items observableList.add("A"); // Observer gets A observableList.add("B"); // Observer gets B // Later addition triggers the Observer automatically observableList.add("C"); // Observer gets C
- RxJava Observables don’t track changes to the collections they’re initialized with—they emit a fixed sequence based on the collection’s state at the time of emission.
- To react to dynamic list changes, use a Subject to manually emit updates, or build a custom observable list that notifies the Observer when modifications happen.
内容的提问来源于stack exchange,提问作者Mohammed Raheez

