RXJava技术问题:当某一Observable有数据时永久切换Observable
Hey there! Let's tackle your problem of switching permanently to an Observable once it becomes available in RxJava. Based on your sample code, I'll cover two common scenarios you might be aiming for:
Scenario 1: Start with initial values, switch to new values once they're available
If you want to receive emissions from initialValues first, but immediately switch to newValues as soon as it emits any data (and stop listening to initialValues), you can combine takeUntil and concat:
import org.junit.Test; import io.reactivex.Observable; import com.google.common.util.concurrent.SettableFuture; import java.util.concurrent.TimeUnit; @Test public void testSwitchToNewValuesWhenAvailable() throws InterruptedException { final Observable<Long> initialValues = Observable.fromArray(100L, 200L, 300L) // Simulate real-world streaming delays between emissions .concatMap(val -> Observable.just(val).delay(200, TimeUnit.MILLISECONDS)); final SettableFuture<Long> future = SettableFuture.create(); // Simulate new values becoming available after 1 second new Thread(() -> { try { Thread.sleep(1000); future.set(400L); // If newValues was an infinite stream, you could emit more values here // e.g., future.set(Observable.interval(300, TimeUnit.MILLISECONDS).map(i -> 400L + i)); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }).start(); final Observable<Long> newValues = Observable.fromFuture(future); // Combine observables: take initial values until newValues emits, then switch to newValues Observable.concat( initialValues.takeUntil(newValues), newValues ).subscribe( val -> System.out.println("Received: " + val), err -> err.printStackTrace(), () -> System.out.println("Completed") ); // Wait long enough to see all emissions Thread.sleep(2000); }
How this works:
initialValues.takeUntil(newValues)emits all values frominitialValuesuntilnewValuessends its first emission. At that point, it completes the initial stream.Observable.concatthen subscribes tonewValuesand emits all its values (including the one that triggered the switch).
Scenario 2: Choose whichever Observable emits first and stick with it
If you want to ignore all other Observables once the first one emits any data (regardless of which one it is), use ambWith (short for "ambiguous"):
import org.junit.Test; import io.reactivex.Observable; import com.google.common.util.concurrent.SettableFuture; import java.util.concurrent.TimeUnit; @Test public void testPickFirstEmittingObservable() throws InterruptedException { final Observable<Long> initialValues = Observable.fromArray(100L, 200L, 300L); final SettableFuture<Long> future = SettableFuture.create(); new Thread(() -> { try { Thread.sleep(500); // Make newValues emit later than initialValues future.set(400L); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }).start(); final Observable<Long> newValues = Observable.fromFuture(future); // Subscribe to the first Observable that emits, ignore the other initialValues.ambWith(newValues) .subscribe( val -> System.out.println("Received: " + val), err -> err.printStackTrace(), () -> System.out.println("Completed") ); Thread.sleep(1000); }
How this works:
ambWithsubscribes to both Observables simultaneously, but only forwards emissions from the first one that sends data. The other Observable is immediately disposed of and ignored.
In your original code, since initialValues emits immediately, this would only receive the initial 100, 200, 300 values and ignore newValues entirely. If you wanted newValues to be picked first, adjust the delay so it emits before initialValues.
Both approaches let you "permanently switch" to an Observable once it becomes available—choose the one that fits your exact use case!
内容的提问来源于stack exchange,提问作者nucandrei

