如何用RxJava2/RxAndroid按序批量处理多报告的链式接口请求?
嘿,这个场景我之前做项目的时候刚好碰到过,用RxJava的几个操作符组合起来就能完美满足你的需求,咱们一步步来拆解实现思路和代码:
核心思路拆解
你的需求核心有三个关键点:
- 按顺序处理每个Report:必须等前一个Report的三个API都处理完(不管成功失败),再开始下一个
- 单个Report的API链式调用:API1→API2→API3,每个API的结果作为下一个的输入
- 错误隔离:单个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的输入,保证请求顺序执行。
额外注意事项
- 如果你的API是用Retrofit实现的,那它返回的Observable本身就是冷Observable,每次订阅都会发起新的请求,完全适配这个逻辑。
- 如果需要在某个API失败时做更复杂的处理(比如重试几次再放弃),可以在单个API的Observable里添加
retry()操作符,再结合onErrorResumeNext。 - 线程调度的位置:
subscribeOn只需要调用一次,指定上游(API请求)在IO线程执行;observeOn可以根据需要多次调用,比如在更新UI前切换到主线程。
内容的提问来源于stack exchange,提问作者Vinamra
相关产品推荐
相关产品推荐

