You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在单条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 the Result accumulator.
  • Completable Merge: Completable.mergeArray() combines all three processing branches into one. This means the entire chain will only trigger onComplete once 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 replace publish() with share() (a shorthand for publish().refCount()). This automatically connects when the first subscriber joins and disconnects when the last leaves.

内容的提问来源于stack exchange,提问作者Alex Kokorin

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 06:51:49