Retrofit 2 + RxJava 1.x 适配器(adapter-rxjava)实战指南:从配置到线程调度与源码原理
2026/9/18 17:19:09 网站建设 项目流程

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 的流式类型(ObservableSingleCompletable),涵盖三种返回类型模式的语义差异、三种线程调度方式、Maven/Gradle 依赖引入方式,并结合仓库源码(RxJavaCallAdapterFactoryRxJavaCallAdapterCallArbiter等)剖析其背压、取消与错误处理实现。读完本文,你将掌握在 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后,服务接口方法即可返回ObservableSingleCompletable

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为响应体类型。SingleCompletable在源码注释中被标注为实验性支持——因为 RxJava 1.x 中这两个类型尚未被官方视为稳定 API,可能存在不兼容变更。

1.2 类型判定的源码逻辑

在 RxJavaCallAdapterFactory.get() 中,工厂先通过getRawType(returnType)判断原始类型是否为ObservableSingleCompletable,不属于三者则返回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() 拿到业务数据 } });

这种模式适合希望“所有结果都在一条流里统一处理”的场景,避免onErroronNext分开分支。

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 实现了SubscriptionProducer两个接口,负责协调「订阅者的请求量」与「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的二次抛出做了防御:捕获OnCompletedFailedExceptionOnErrorFailedExceptionOnErrorNotImplementedException后转交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.xrx.Observablerx.Singlerx.Completablerx.Scheduler均来自 RxJava 1.x 包路径),请勿与 RxJava 2/3 混淆;
  • 开发版本的快照(Snapshot)发布在 Sonatype 的snapshots仓库中(详见 README.md 底部链接),若需要体验最新未发布特性可配置该快照源;
  • 仓库中同时提供了 RxJava 2(retrofit-adapters/rxjava2)与 RxJava 3(retrofit-adapters/rxjava3)的对应适配器模块,新项目应优先考虑新版本 RxJava 对应的适配器。

七、注意事项与迁移提示

  1. Single/Completable为实验性:README 与源码注释均提示,这两个 RxJava 1.x 类型本身不被 RxJava 官方视为稳定,API 可能发生不兼容变更;
  2. HttpException已废弃:适配器模块中的retrofit2.adapter.rxjava.HttpException只是retrofit2.HttpException的弃用子类(见 HttpException.java),建议直接使用 Retrofit 核心包中的retrofit2.HttpException
  3. 请求线程模型:默认同步执行意味着如果不配合subscribeOn/createAsync,请求会阻塞订阅线程——在主线程订阅时务必留意 ANR 风险;
  4. 冷流语义:每次订阅都会通过clone()重新执行请求,这是符合 RxJava 冷 Observable 预期的行为,也意味着多个订阅者会触发多次网络请求;
  5. 迁移到 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询