如何在单条RxJava链中将Observable值分发至3个不同对象?
Merging Multiple RxJava Processing Paths into One Chain
Hey there! Let's tackle this RxJava problem together. The key here is to use multicasting to share the source Observable's emissions across all three processing branches—this way we don't re-emit the source data multiple times, and all branches get the exact same sequence of values.
Core Approach: Use publish() for Shared Emissions
The publish() operator converts a regular Observable into a ConnectableObservable, which lets multiple subscribers share the same stream of events. We'll define all three processing branches within publish()'s lambda, then merge them into a single Completable to coordinate completion and error handling.
Complete Code Example
// Define your source Observable once Observable<KV> source = Observable.just( new KV("1", "Alex"), new KV("2", "LeBron"), new KV("3", "Kobe"), new KV("1", "Kevin") ); // Merge all processing logic into one chain source.publish(sharedObservable -> { // Branch 1: Regular processing of each KV Completable regularProcessing = sharedObservable .flatMapCompletable(keyValue -> doSth(keyValue)); // Branch 2: Group by key, collect values, then write result Completable groupedProcessing1 = sharedObservable .groupBy(KV::getKey) .collect(Result::new, (resultAccumulator, keyGroup) -> { // Collect all values in this key group into the accumulator keyGroup.map(KV::getVal).subscribe(resultAccumulator::add); }) .flatMapCompletable(result -> writeResultSomeWhere1(result)); // Branch 3: Second grouped collection/processing path Completable groupedProcessing2 = sharedObservable .groupBy(KV::getKey) .collect(Result::new, (resultAccumulator, keyGroup) -> { keyGroup.map(KV::getVal).subscribe(resultAccumulator::add); }) .flatMapCompletable(result -> writeResultSomeWhere2(result)); // Merge all branches into a single Completable to track overall completion return Completable.mergeArray(regularProcessing, groupedProcessing1, groupedProcessing2); }) .subscribe( () -> System.out.println("All processing tasks finished successfully!"), throwable -> System.err.println("Processing failed: " + throwable.getMessage()) );
Key Details Explained
publish(): Ensures all three branches receive the exact same emissions from the source Observable—no duplicate data emission, which is efficient and consistent.- Grouped Collection: In each grouped branch,
groupBy()gives us an Observable per key. We subscribe to each group's Observable to collect all its values into theResultaccumulator. - Completable Merge:
Completable.mergeArray()combines all three processing branches into one. This means the entire chain will only triggeronCompleteonce all three branches have finished, and any error in any branch will immediately trigger the error handler. - Alternative:
share(): If you don't need strict control over when the source starts emitting, you can replacepublish()withshare()(a shorthand forpublish().refCount()). This automatically connects when the first subscriber joins and disconnects when the last leaves.
内容的提问来源于stack exchange,提问作者Alex Kokorin
相关产品推荐
相关产品推荐

