如何在RxJava中合并多Observable时忽略空源以实现实时发射
Great question—this is a common gotcha with combineLatest! The root issue is that this operator requires every source Observable to emit at least one item before it produces its first result. When obs4 emits an empty list, your flatMap { Observable.fromIterable(it) } call creates an empty stream—so combineLatest waits forever for that first emission from obs4.
Here are two clean, idiomatic approaches to fix this:
Approach 1: Wrap items in Optional with a default empty value
We can convert each source Observable to emit Optional<MyItem> values, and use startWithItem to inject an empty optional as the initial emission. This guarantees every source has at least one emission, so combineLatest will start producing results immediately.
// Helper function to wrap any Observable<MyItem> with an initial empty Optional fun Observable<MyItem>.withEmptyFallback(): Observable<Optional<MyItem>> { return this.map { Optional.of(it) } .startWithItem(Optional.empty()) } // Process obs4: flatMap to individual items, then apply the fallback val processedObs4 = obs4.flatMap { Observable.fromIterable(it) } .withEmptyFallback() // Wrap all other observables val observables = arrayOf( obs0.withEmptyFallback(), obs1.withEmptyFallback(), obs2.withEmptyFallback(), obs3.withEmptyFallback(), processedObs4, obs5.withEmptyFallback(), obs6.withEmptyFallback(), obs7.withEmptyFallback(), obs8.withEmptyFallback(), obs9.withEmptyFallback(), obs10.withEmptyFallback() ) // Combine and filter out empty optionals Observable.combineLatestDelayError(observables) { items -> items.filterIsInstance<Optional<MyItem>>() .filter { it.isPresent } .map { it.get() } .filter { ... } // Your existing filter logic .map { ... } // Your existing map logic }
Approach 2: Use BehaviorSubject with a null default
BehaviorSubject holds the latest value emitted by the source, and emits it to new subscribers immediately. By initializing each subject with a null default, we ensure combineLatest always has a value to work with from every source.
// Helper function to convert an Observable<MyItem> to a BehaviorSubject with null default fun Observable<MyItem>.toBehaviorSubject(): BehaviorSubject<MyItem?> { val subject = BehaviorSubject.createDefault<MyItem?>(null) this.subscribe(subject) return subject } // Process obs4 first val processedObs4 = obs4.flatMap { Observable.fromIterable(it) } .toBehaviorSubject() // Wrap all other observables val subjects = arrayOf( obs0.toBehaviorSubject(), obs1.toBehaviorSubject(), obs2.toBehaviorSubject(), obs3.toBehaviorSubject(), processedObs4, obs5.toBehaviorSubject(), obs6.toBehaviorSubject(), obs7.toBehaviorSubject(), obs8.toBehaviorSubject(), obs9.toBehaviorSubject(), obs10.toBehaviorSubject() ) // Combine and filter out null values Observable.combineLatestDelayError(subjects) { items -> items.filterIsInstance<MyItem>() // Nulls are excluded here .filter { ... } // Your existing filter logic .map { ... } // Your existing map logic }
Key Notes
- Both approaches ensure
combineLatestreceives at least one emission from every source, so it will start producing results even ifobs4emits an empty list. - If you don't want to use
Optionalornull, you can define a dedicated "empty" marker instance (e.g.,object : MyItem()), then filter that out during processing instead.
内容的提问来源于stack exchange,提问作者Ilya E

