RxJava技术问询:合并两个Single<MyData>并处理失败与成功场景
Solution for Combining Two
Single<MyData> with Custom Error Logic Hey there! Let's break down how to implement your desired logic for merging two Single<MyData> streams. The key here is to handle success/failure combinations explicitly without letting a single failure immediately terminate the whole stream.
Core Approach
We'll use RxJava's materialize() operator to wrap all stream events (success, error) into Notification objects, then collect these events to decide how to proceed based on your rules. Here's the step-by-step implementation:
1. Define Your Reporting Method
First, let's assume you have a method to report the error and successful data when one stream fails:
private void reportErrorAndSuccess(Throwable failureError, MyData successData) { // Your custom reporting logic here (e.g., log to analytics, send to server) System.out.println("Reported failure: " + failureError.getMessage()); System.out.println("Reported success data: " + successData.toString()); }
2. Implement the Merging Logic
Now, let's combine the two Single streams and apply your rules:
Single<MyData> single1 = ...; // Your first Single source Single<MyData> single2 = ...; // Your second Single source Observable<MyData> combinedObservable = Observable.merge( single1.toObservable().materialize(), single2.toObservable().materialize() ) .toList() .flatMapObservable(notifications -> { List<MyData> successfulResults = new ArrayList<>(); List<Throwable> failures = new ArrayList<>(); // Classify each event as success or failure for (Notification<MyData> notification : notifications) { if (notification.isOnNext()) { successfulResults.add(notification.getValue()); } else if (notification.isOnError()) { failures.add(notification.getError()); } } // Apply your business rules if (failures.size() == 2) { // Both failed: emit the first error (or combine errors as needed) return Observable.error(failures.get(0)); } else if (failures.size() == 1) { // One failed, one succeeded: report and emit the successful data reportErrorAndSuccess(failures.get(0), successfulResults.get(0)); return Observable.just(successfulResults.get(0)); } else { // Both succeeded: emit all successful data return Observable.fromIterable(successfulResults); } });
How This Works
materialize(): Converts each stream's events (successonNext/onComplete, failureonError) intoNotificationobjects. This lets us capture failures instead of letting them terminate the stream early.merge(): Combines the two wrapped streams into one, so we can process both events together.toList(): Collects bothNotificationobjects into a list (since we have exactly twoSinglesources, this will always collect two items).- Post-collection logic: We sort the events into successes and failures, then apply your three rules:
- Both failed → emit an error.
- One failed, one succeeded → report the error and success data, then emit the successful result.
- Both succeeded → emit all successful results and complete.
Optional Adjustments
- If you need the final result to be a
Singleinstead of anObservable, just add.singleOrError()at the end of the chain (note: this will emit an error if both streams succeed, sincesingleOrErrorexpects exactly one item—use.firstOrError()if you want to take the first success instead). - If you need to run the reporting logic on a specific thread (e.g., main thread for UI, background thread for network calls), add
observeOn(Schedulers.io())orobserveOn(AndroidSchedulers.mainThread())before theflatMapObservableoperator.
内容的提问来源于stack exchange,提问作者IS1_SO
相关产品推荐
相关产品推荐

