Retrofit 2 + RxJava 1.x 适配器(adapter-rxjava)实战指南:从配置到线程调度与源码原理
【免费下载链接】retrofitA type-safe HTTP client for Android and the JVM项目地址: https://gitcode.com/gh_mirrors/re/retrofit
导读
本文基于当前仓库中 Retrofit 官方模块 retrofit-adapters/rxjava 的说明文档展开,系统讲解如何让 Retrofit 服务接口直接返回 RxJava 1.x 的流式类型(Observable、Single、Completable),涵盖三种返回类型模式的语义差异、三种线程调度方式、Maven/Gradle 依赖引入方式,并结合仓库源码(RxJavaCallAdapterFactory、RxJavaCallAdapter、CallArbiter等)剖析其背压、取消与错误处理实现。读完本文,你将掌握在 Retrofit 项目中接入adapter-rxjava的完整方案,并理解其底层工作机理,为迁移 RxJava 2/3 或排查线程与背压问题打下基础。
一、模块定位:为 Retrofit 适配 RxJava 1.x 流式类型
adapter-rxjava是 Retrofit 官方提供的一个CallAdapter(调用适配器)扩展模块,其作用是把 Retrofit 内部基于Call的同步/异步请求模型,桥接为 RxJava 1.x 的响应式流类型。该模块的gradle.properties中将其描述为:
"A Retrofit CallAdapter for RxJava's stream types."
在 RxJavaCallAdapterFactory.java 的类注释中明确说明:将本工厂加入Retrofit后,服务接口方法即可返回Observable、Single或Completable。
1.1 支持的返回类型
根据 README.md 与源码实现,该适配器支持以下返回类型:
| 返回类型 | 说明 |
|---|---|
Observable<T> | 直接发射反序列化后的响应体 |
Observable<Response<T>> | 发射包装了全部 HTTP 响应信息的Response对象 |
Observable<Result<T>> | 发射Result包装对象,成功与失败都以onNext形式发射 |
Single<T> | 单次发射的响应体版本(实验性) |
Single<Response<T>> | 单次发射的 Response 包装版本(实验性) |
Single<Result<T>> | 单次发射的 Result 包装版本(实验性) |
Completable | 忽略响应体,只关心请求是否成功完成(实验性) |
其中T为响应体类型。Single与Completable在源码注释中被标注为实验性支持——因为 RxJava 1.x 中这两个类型尚未被官方视为稳定 API,可能存在不兼容变更。
1.2 类型判定的源码逻辑
在 RxJavaCallAdapterFactory.get() 中,工厂先通过getRawType(returnType)判断原始类型是否为Observable、Single或Completable,不属于三者则返回null(让后续的 CallAdapter 工厂继续处理):
Class<?> rawType = getRawType(returnType); boolean isSingle = rawType == Single.class; boolean isCompletable = rawType == Completable.class; if (rawType != Observable.class && !isSingle && !isCompletable) { return null; }对于泛型参数,代码要求必须是ParameterizedType(即必须写成Observable<Foo>而非裸的Observable),否则抛出IllegalStateException。随后解析Observable/Single的泛型上界:若是Response<Foo>或Result<Foo>则进一步解包,得到真正的 HTTP 响应体类型responseType。
二、快速上手:注册工厂并定义服务接口
2.1 注册RxJavaCallAdapterFactory
在构建Retrofit实例时,通过addCallAdapterFactory()注册适配器工厂(README.md 原示例):
Retrofit retrofit = new Retrofit.Builder() .baseUrl("https://example.com/") .addCallAdapterFactory(RxJavaCallAdapterFactory.create()) .build();2.2 定义返回响应式类型的服务方法
注册之后,服务接口的返回值即可使用上文列出的任意响应式类型,例如:
interface MyService { @GET("/user") Observable<User> getUser(); }调用时即可享受 RxJava 的链式操作与线程调度能力:
myService.getUser() .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe(user -> showUser(user), throwable -> handleError(throwable));注意:
addCallAdapterFactory的注册顺序会影响匹配优先级——Retrofit会按注册顺序依次询问每个工厂,因此若同时使用多个适配器,应将RxJavaCallAdapterFactory放在合适的位置。
三、三种响应类型模式的语义与源码依据
README 中列出的Observable<T>、Observable<Response<T>>、Observable<Result<T>>三种形式,其行为差异在 RxJavaCallAdapterFactory.java 的类注释中有权威说明,三种模式对应三种不同的OnSubscribe包装类:
| 返回类型 | 成功(2XX) | 失败(非 2XX) | 网络错误 | 底层包装 |
|---|---|---|---|---|
Observable<T> | onNext(反序列化后的 body) | onError(HttpException) | onError(IOException) | BodyOnSubscribe.java |
Observable<Response<T>> | onNext(Response 对象) | onNext(Response 对象) | onError(IOException) | 直接传递,无包装 |
Observable<Result<T>> | onNext(Result.response(resp)) | onNext(Result.response(resp)) | onNext(Result.error(e)) | ResultOnSubscribe.java |
3.1 Body 模式:Observable<T>
这是最常用的模式。从 BodyOnSubscribe.java 可以看到:BodySubscriber.onNext(Response)先判断response.isSuccessful():
- 成功:将
response.body()交给下游onNext; - 失败(非 2XX):构造
retrofit2.adapter.rxjava.HttpException(其源码仅继承了retrofit2.HttpException并标记@Deprecated,见 HttpException.java),通过onError终止流。
因此使用Observable<T>时,业务上“请求已发出但 HTTP 返回 4xx/5xx”的情况会被当作错误流处理。
3.2 Response 模式:Observable<Response<T>>
此模式下适配器不进行任何包装,直接把retrofit2.Response<T>逐条发射给订阅者。由于Response本身携带了code()、headers()、isSuccessful()等信息,你可以在onNext中自行判断 HTTP 状态码,非 2XX 不会触发onError,只有网络层错误(IOException)才会进入onError。
3.3 Result 模式:Observable<Result<T>>
Result是适配器模块提供的公开类型,定义于 Result.java。从 ResultOnSubscribe.java 可见,无论请求成功还是失败,都通过onNext发射,随后onCompleted:
@Override public void onNext(Response<R> response) { subscriber.onNext(Result.response(response)); } @Override public void onError(Throwable throwable) { try { subscriber.onNext(Result.error(throwable)); } catch (Throwable t) { subscriber.onError(t); return; } subscriber.onCompleted(); }Result提供三个核心访问方法(Result.java):
response():仅在isError()为 false 时非空,返回Response<T>;error():仅在isError()为 true 时非空,返回底层异常;若异常是IOException表示网络传输问题,其他异常类型属于意外失败(配置错误、编程错误等);isError():判断请求是否以错误结束。
observable.subscribe(result -> { if (result.isError()) { // 网络错误或转换错误 } else { // result.response().body() 拿到业务数据 } });这种模式适合希望“所有结果都在一条流里统一处理”的场景,避免onError与onNext分开分支。
3.4 模式选择小结
- 只要 2XX:用
Observable<T>,非 2XX 直接走onError; - 需要检查状态码/响应头:用
Observable<Response<T>>; - 希望网络错误与 HTTP 错误统一走
onNext分支:用Observable<Result<T>>。
四、线程调度:默认同步与三种控制方式
README 明确指出:默认情况下所有响应式类型都在当前订阅线程上同步执行请求。仓库提供三种方式控制请求发生的线程:
4.1 方式一:对返回的响应式类型调用subscribeOn
RxJava 的subscribeOn决定订阅发生(即请求发起)的线程:
observable.subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe(...);这是最灵活的方式,Scheduler完全由调用方按需选择。
4.2 方式二:createAsync()使用 OkHttp 内部线程池
RxJavaCallAdapterFactory.createAsync() 创建出的工厂会将isAsync置为true,使请求通过 OkHttp 的Call.enqueue()异步执行,运行在 OkHttp 内部的线程池中。仓库测试 AsyncTest.java 与 CancelDisposeTest.java 均使用该工厂验证异步与取消行为:
Retrofit retrofit = new Retrofit.Builder() .baseUrl(...) .addCallAdapterFactory(RxJavaCallAdapterFactory.createAsync()) .build();4.3 方式三:createWithScheduler(Scheduler)指定默认订阅调度器
RxJavaCallAdapterFactory.createWithScheduler(Scheduler) 为所有由该工厂创建的流设置默认的subscribeOn调度器。注意两点:
- 传入
null会抛出NullPointerException("scheduler == null"),这一点由测试 RxJavaCallAdapterFactoryTest.nullSchedulerThrows 明确验证; - 该方式仍是同步请求模型,只是请求发生在指定调度器的线程上。
RxJavaCallAdapterFactory factory = RxJavaCallAdapterFactory.createWithScheduler(Schedulers.io());测试 ObservableWithSchedulerTest.java、SingleWithSchedulerTest.java 与 CompletableWithSchedulerTest.java 均采用该方式验证各类响应式类型的线程行为。
4.4 三种方式对比
| 方式 | 请求线程 | 底层执行路径 | 适用场景 |
|---|---|---|---|
默认create()+subscribeOn | 由调用方每次指定 | Call.execute()同步(CallExecuteOnSubscribe.java) | 灵活控制,推荐用于 RxJava 熟练场景 |
createAsync() | OkHttp 内部线程池 | Call.enqueue()异步(CallEnqueueOnSubscribe.java) | 不想关心调度器、直接异步 |
createWithScheduler(s) | 指定的默认调度器 | Call.execute()+subscribeOn(s) | 统一为整个应用固定订阅线程 |
五、底层原理:适配器如何把Call变成Observable
5.1adapt()的组装流程
RxJavaCallAdapter.adapt(Call) 是整个桥接的核心:
OnSubscribe<Response<R>> callFunc = isAsync ? new CallEnqueueOnSubscribe<>(call) : new CallExecuteOnSubscribe<>(call); OnSubscribe<?> func; if (isResult) { func = new ResultOnSubscribe<>(callFunc); } else if (isBody) { func = new BodyOnSubscribe<>(callFunc); } else { func = callFunc; // Response 模式 } Observable<?> observable = Observable.create(func); if (scheduler != null) { observable = observable.subscribeOn(scheduler); } if (isSingle) { return observable.toSingle(); } if (isCompletable) { return observable.toCompletable(); } return observable;流程可归纳为:Call → OnSubscribe<Response<T>> → (Body/Result 包装)→ Observable → 可选 subscribeOn → toSingle()/toCompletable()。
5.2 每个订阅者都会clone()一份 Call
Call是一次性(one-shot)类型,同一个实例不能重复执行。因此两个OnSubscribe实现(CallExecuteOnSubscribe.call()、CallEnqueueOnSubscribe.call())在call()方法中都会先执行originalCall.clone(),保证每个订阅者拥有独立的请求执行——这正是冷 Observable 语义:每次订阅都会真正发起一次 HTTP 请求。
5.3CallArbiter:背压、取消与状态机
CallArbiter.java 实现了Subscription与Producer两个接口,负责协调「订阅者的请求量」与「HTTP 响应的到达」之间的竞态,其内部用AtomicInteger维护四态状态机:
| 状态 | 含义 |
|---|---|
STATE_WAITING(0) | 等待订阅者请求或响应到达 |
STATE_REQUESTED(1) | 订阅者已请求数据 |
STATE_HAS_RESPONSE(2) | 响应已到达,等待被取走 |
STATE_TERMINATED(3) | 已终止 |
两个关键方法:
request(long amount):订阅者发起背压请求。若amount == 0直接忽略;若状态为WAITING则置为REQUESTED;若状态为HAS_RESPONSE则置为TERMINATED并投递响应;emitResponse(...):网络层产生响应。若订阅者已请求(REQUESTED)则立即投递;若还在WAITING则暂存响应并转入HAS_RESPONSE,等待订阅者请求。
unsubscribe()会置位unsubscribed并调用call.cancel()取消底层 OkHttp 请求——这就是 CancelDisposeTest.java 所验证的「取消订阅即取消请求」行为。
5.4 错误处理与异常安全
CallArbiter.deliverResponse()与BodyOnSubscribe均对onError/onCompleted的二次抛出做了防御:捕获OnCompletedFailedException、OnErrorFailedException、OnErrorNotImplementedException后转交RxJavaPlugins.getInstance().getErrorHandler().handleError()处理,避免异常吞掉或逃逸,这是仓库中多组ThrowingTest/ThrowingSafeSubscriberTest测试所覆盖的健壮性设计。
六、依赖引入:Maven 与 Gradle
README 的 Download 章节提供了两种主流构建工具的引入方式(latest.version请替换为实际版本号,可在仓库 gradle/libs.versions.toml 与各模块gradle.properties中查看当前仓库使用的版本管理方式):
Maven:
<dependency> <groupId>com.squareup.retrofit2</groupId> <artifactId>adapter-rxjava</artifactId> <version>latest.version</version> </dependency>Gradle:
implementation 'com.squareup.retrofit2:adapter-rxjava:latest.version'此外需要注意:
- 本模块依赖RxJava 1.x(
rx.Observable、rx.Single、rx.Completable、rx.Scheduler均来自 RxJava 1.x 包路径),请勿与 RxJava 2/3 混淆; - 开发版本的快照(Snapshot)发布在 Sonatype 的
snapshots仓库中(详见 README.md 底部链接),若需要体验最新未发布特性可配置该快照源; - 仓库中同时提供了 RxJava 2(retrofit-adapters/rxjava2)与 RxJava 3(retrofit-adapters/rxjava3)的对应适配器模块,新项目应优先考虑新版本 RxJava 对应的适配器。
七、注意事项与迁移提示
Single/Completable为实验性:README 与源码注释均提示,这两个 RxJava 1.x 类型本身不被 RxJava 官方视为稳定,API 可能发生不兼容变更;HttpException已废弃:适配器模块中的retrofit2.adapter.rxjava.HttpException只是retrofit2.HttpException的弃用子类(见 HttpException.java),建议直接使用 Retrofit 核心包中的retrofit2.HttpException;- 请求线程模型:默认同步执行意味着如果不配合
subscribeOn/createAsync,请求会阻塞订阅线程——在主线程订阅时务必留意 ANR 风险; - 冷流语义:每次订阅都会通过
clone()重新执行请求,这是符合 RxJava 冷 Observable 预期的行为,也意味着多个订阅者会触发多次网络请求; - 迁移到 RxJava 2/3:若计划升级,可参考本仓库的 rxjava2 与 rxjava3 模块,它们的
RxJava2CallAdapterFactory/RxJava3CallAdapterFactoryAPI 与本文所述工厂保持同构,迁移成本主要体现在 RxJava 本身的 API 差异(如Flowable的引入、Function接口包名变化等)。
八、测试与验证
仓库 retrofit-adapters/rxjava/src/test 提供了覆盖全面、可直接参考的测试用例,可作为理解与验证本文所述行为的权威依据:
- RxJavaCallAdapterFactoryTest.java:工厂对返回类型的解析、
null调度器抛 NPE、非 RxJava 类型返回null等; - ObservableTest.java、SingleTest.java、CompletableTest.java:各类型的成功/失败语义;
- AsyncTest.java、CancelDisposeTest.java:
createAsync()异步执行与取消订阅行为; - ResultTest.java:
Result包装类型的响应/错误判定。
通过这些测试与源码的相互印证,可以确认本文所述的所有行为均有仓库内实现与用例支撑。
总结
adapter-rxjava用约十个类的轻量实现,把 Retrofit 的Call请求模型无缝接入 RxJava 1.x 生态:RxJavaCallAdapterFactory负责识别返回类型并决定 Body/Response/Result 三种语义,CallArbiter以状态机优雅解决背压、竞态与取消问题,而create()、createAsync()、createWithScheduler(Scheduler)三个工厂方法覆盖了从「调用方自定义调度」到「OkHttp 线程池」再到「全局默认调度器」的全部线程诉求。掌握本文内容后,你既能在老项目中快速接入 RxJava 1.x 响应式请求,也能基于对底层机制的了解平滑迁移到 RxJava 2/3 适配器。
【免费下载链接】retrofitA type-safe HTTP client for Android and the JVM项目地址: https://gitcode.com/gh_mirrors/re/retrofit
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考