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

如何用RxJava2/RxAndroid按序批量处理多报告的链式接口请求?

嘿,这个场景我之前做项目的时候刚好碰到过,用RxJava的几个操作符组合起来就能完美满足你的需求,咱们一步步来拆解实现思路和代码:

核心思路拆解

你的需求核心有三个关键点:

  1. 按顺序处理每个Report:必须等前一个Report的三个API都处理完(不管成功失败),再开始下一个
  2. 单个Report的API链式调用:API1→API2→API3,每个API的结果作为下一个的输入
  3. 错误隔离:单个Report的任何API失败,终止该Report的后续请求,但不影响其他Report的处理

针对这几点,我们需要用到concatMap和onErrorResumeNext这两个核心操作符,再配合基础的flatMap来做API链式调用。

具体代码实现

首先,先模拟一下你的三个API请求(实际项目中可以替换成Retrofit返回的Observable):

// 假设你的Report实体类
public class Report {
    private String id;
    // getter、setter省略
}

// 模拟每个API返回的结果类
public class Result1 { private String reportId; /* getter */ }
public class Result2 { private String reportId; /* getter */ }
public class Result3 { private String reportId; /* getter */ }

// 模拟API1:传入Report,发起请求并返回Result1的Observable
private Observable<Result1> callApi1(Report report) {
    return Observable.just(report)
            .delay(1, TimeUnit.SECONDS) // 模拟网络请求延迟
            .map(report -> new Result1(report.getId()));
}

// 模拟API2:传入API1的结果,发起请求并返回Result2的Observable
private Observable<Result2> callApi2(Result1 result1) {
    return Observable.just(result1)
            .delay(1, TimeUnit.SECONDS)
            .map(result1 -> new Result2(result1.getReportId()));
}

// 模拟API3:传入API2的结果,发起请求并返回Result3的Observable
private Observable<Result3> callApi3(Result2 result2) {
    return Observable.just(result2)
            .delay(1, TimeUnit.SECONDS)
            .map(result2 -> new Result3(result2.getReportId()));
}

接下来是核心的处理逻辑:

// 假设你已经获取到了List<Report> reportsList
Observable.fromIterable(reportsList)
        // 用concatMap保证串行处理每个Report——前一个处理完才会触发下一个
        .concatMap(report -> {
            // 单个Report的API链式调用:API1 → API2 → API3
            return callApi1(report)
                    .flatMap(this::callApi2) // 用flatMap把API1的结果传给API2
                    .flatMap(this::callApi3) // 把API2的结果传给API3
                    // 关键:捕获当前Report链式请求中的任何错误,返回empty()终止该Report的后续请求,但不中断整个流
                    .onErrorResumeNext(throwable -> {
                        // 这里可以做错误日志、UI提示等操作
                        Log.e("RxJava", "报告[" + report.getId() + "]处理失败:" + throwable.getMessage());
                        return Observable.empty();
                    });
        })
        // 线程调度:根据实际需求添加,比如IO线程发起请求,主线程处理结果
        .subscribeOn(Schedulers.io())
        .observeOn(AndroidSchedulers.mainThread())
        .subscribe(
                // 单个Report处理成功的回调
                result3 -> Log.d("RxJava", "报告[" + result3.getReportId() + "]处理完成"),
                // 全局错误回调:理论上这里不会触发,因为每个Report的错误已经被onErrorResumeNext捕获
                throwable -> Log.e("RxJava", "全局异常:" + throwable.getMessage()),
                // 所有Report处理完毕的回调
                () -> Log.d("RxJava", "所有报告处理完成!")
        );

关键操作符解释

  • concatMap:这是实现串行处理的核心。和flatMap不同,concatMap会严格按照原始Observable的顺序来处理每个item,必须等前一个item的Observable完全结束(不管成功还是失败),才会订阅下一个item的Observable,完美匹配你“按序处理”的需求。
  • onErrorResumeNext:这个操作符用来捕获单个Report链式请求中的错误。当某个API请求失败时,它会替代当前的错误流,返回一个Observable.empty(),这样当前Report的流就正常结束,不会把错误传递到全局的subscribe的onError方法,从而保证整个列表的处理可以继续进行下一个Report。
  • flatMap:用来实现API的链式调用,把前一个API的结果作为下一个API的输入,保证请求顺序执行。

额外注意事项

  1. 如果你的API是用Retrofit实现的,那它返回的Observable本身就是冷Observable,每次订阅都会发起新的请求,完全适配这个逻辑。
  2. 如果需要在某个API失败时做更复杂的处理(比如重试几次再放弃),可以在单个API的Observable里添加retry()操作符,再结合onErrorResumeNext。
  3. 线程调度的位置:subscribeOn只需要调用一次,指定上游(API请求)在IO线程执行;observeOn可以根据需要多次调用,比如在更新UI前切换到主线程。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.07 20:27:44