Retrofit 2 + RxJava 1.x 适配器(adapter-rxjava)实战指南:从配置到线程调度与源码原理

发布时间:2026/9/18 17:19:28
Retrofit 2 + RxJava 1.x 适配器(adapter-rxjava)实战指南:从配置到线程调度与源码原理
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 RxJavas stream types.在 RxJavaCallAdapterFactory.java 的类注释中明确说明将本工厂加入Retrofit后服务接口方法即可返回Observable、Single或Completable。1.1 支持的返回类型根据 README.md 与源码实现该适配器支持以下返回类型返回类型说明ObservableT直接发射反序列化后的响应体ObservableResponseT发射包装了全部 HTTP 响应信息的Response对象ObservableResultT发射Result包装对象成功与失败都以onNext形式发射SingleT单次发射的响应体版本实验性SingleResponseT单次发射的 Response 包装版本实验性SingleResultT单次发射的 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即必须写成ObservableFoo而非裸的Observable否则抛出IllegalStateException。随后解析Observable/Single的泛型上界若是ResponseFoo或ResultFoo则进一步解包得到真正的 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) ObservableUser getUser(); }调用时即可享受 RxJava 的链式操作与线程调度能力myService.getUser() .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe(user - showUser(user), throwable - handleError(throwable));注意addCallAdapterFactory的注册顺序会影响匹配优先级——Retrofit会按注册顺序依次询问每个工厂因此若同时使用多个适配器应将RxJavaCallAdapterFactory放在合适的位置。三、三种响应类型模式的语义与源码依据README 中列出的ObservableT、ObservableResponseT、ObservableResultT三种形式其行为差异在 RxJavaCallAdapterFactory.java 的类注释中有权威说明三种模式对应三种不同的OnSubscribe包装类返回类型成功2XX失败非 2XX网络错误底层包装ObservableTonNext(反序列化后的 body)onError(HttpException)onError(IOException)BodyOnSubscribe.javaObservableResponseTonNext(Response 对象)onNext(Response 对象)onError(IOException)直接传递无包装ObservableResultTonNext(Result.response(resp))onNext(Result.response(resp))onNext(Result.error(e))ResultOnSubscribe.java3.1 Body 模式ObservableT这是最常用的模式。从 BodyOnSubscribe.java 可以看到BodySubscriber.onNext(Response)先判断response.isSuccessful()成功将response.body()交给下游onNext失败非 2XX构造retrofit2.adapter.rxjava.HttpException其源码仅继承了retrofit2.HttpException并标记Deprecated见 HttpException.java通过onError终止流。因此使用ObservableT时业务上“请求已发出但 HTTP 返回 4xx/5xx”的情况会被当作错误流处理。3.2 Response 模式ObservableResponseT此模式下适配器不进行任何包装直接把retrofit2.ResponseT逐条发射给订阅者。由于Response本身携带了code()、headers()、isSuccessful()等信息你可以在onNext中自行判断 HTTP 状态码非 2XX 不会触发onError只有网络层错误IOException才会进入onError。3.3 Result 模式ObservableResultTResult是适配器模块提供的公开类型定义于 Result.java。从 ResultOnSubscribe.java 可见无论请求成功还是失败都通过onNext发射随后onCompletedOverride public void onNext(ResponseR 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.javaresponse()仅在isError()为 false 时非空返回ResponseTerror()仅在isError()为 true 时非空返回底层异常若异常是IOException表示网络传输问题其他异常类型属于意外失败配置错误、编程错误等isError()判断请求是否以错误结束。observable.subscribe(result - { if (result.isError()) { // 网络错误或转换错误 } else { // result.response().body() 拿到业务数据 } });这种模式适合希望“所有结果都在一条流里统一处理”的场景避免onError与onNext分开分支。3.4 模式选择小结只要 2XX用ObservableT非 2XX 直接走onError需要检查状态码/响应头用ObservableResponseT希望网络错误与 HTTP 错误统一走onNext分支用ObservableResultT。四、线程调度默认同步与三种控制方式README 明确指出默认情况下所有响应式类型都在当前订阅线程上同步执行请求。仓库提供三种方式控制请求发生的线程4.1 方式一对返回的响应式类型调用subscribeOnRxJava 的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变成Observable5.1adapt()的组装流程RxJavaCallAdapter.adapt(Call) 是整个桥接的核心OnSubscribeResponseR 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 → OnSubscribeResponseT → Body/Result 包装→ Observable → 可选 subscribeOn → toSingle()/toCompletable()。5.2 每个订阅者都会clone()一份 CallCall是一次性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 与 GradleREADME 的 Download 章节提供了两种主流构建工具的引入方式latest.version请替换为实际版本号可在仓库 gradle/libs.versions.toml 与各模块gradle.properties中查看当前仓库使用的版本管理方式Mavendependency groupIdcom.squareup.retrofit2/groupId artifactIdadapter-rxjava/artifactId versionlatest.version/version /dependencyGradleimplementation com.squareup.retrofit2:adapter-rxjava:latest.version此外需要注意本模块依赖RxJava 1.xrx.Observable、rx.Single、rx.Completable、rx.Scheduler均来自 RxJava 1.x 包路径请勿与 RxJava 2/3 混淆开发版本的快照Snapshot发布在 Sonatype 的snapshots仓库中详见 README.md 底部链接若需要体验最新未发布特性可配置该快照源仓库中同时提供了 RxJava 2retrofit-adapters/rxjava2与 RxJava 3retrofit-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.javacreateAsync()异步执行与取消订阅行为ResultTest.javaResult包装类型的响应/错误判定。通过这些测试与源码的相互印证可以确认本文所述的所有行为均有仓库内实现与用例支撑。总结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),仅供参考