RxJava 2.x 之图解创建、订阅、发射流程

2024-02-19 06:10

本文主要是介绍RxJava 2.x 之图解创建、订阅、发射流程,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

  • 从一个例子开始
  • 创建过程
  • 订阅过程
  • 发射过程
  • 小结
从一个例子开始
Observable.create(new ObservableOnSubscribe<Integer>() {@Overridepublic void subscribe(ObservableEmitter<Integer> emitter) throws Exception {for (int i = 0; i < 3; i++) {emitter.onNext(i);}emitter.onComplete();Log.d(TAG, "subscribe " + Thread.currentThread().getName());}}).subscribeOn(Schedulers.newThread()).map(new Function<Integer, String>() {@Overridepublic String apply(Integer value) throws Exception {Log.d(TAG, "apply " + Thread.currentThread().getName());return "apply " + value;}}).observeOn(AndroidSchedulers.mainThread()).subscribeWith(new ResourceObserver<String>() {@Overridepublic void onNext(String value) {Log.d(TAG, "onNext " + value);}@Overridepublic void onError(Throwable e) {Log.d(TAG, "onError");}@Overridepublic void onComplete() {Log.d(TAG, "onComplete " + Thread.currentThread().getName());}});

来看看输出:

10-26 16:55:17.418 32696-561/com.onzhou.study D/MainActivity: apply RxNewThreadScheduler-1
10-26 16:55:17.418 32696-561/com.onzhou.study D/MainActivity: apply RxNewThreadScheduler-1
10-26 16:55:17.418 32696-561/com.onzhou.study D/MainActivity: create RxNewThreadScheduler-1
10-26 16:55:17.427 32696-32696/com.onzhou.study D/MainActivity: onNext apply 0
10-26 16:55:17.427 32696-32696/com.onzhou.study D/MainActivity: onNext apply 1
10-26 16:55:17.427 32696-32696/com.onzhou.study D/MainActivity: onNext apply 2
10-26 16:55:17.427 32696-32696/com.onzhou.study D/MainActivity: onComplete main

可以看到创建发送转换过程都在子线程中,而最后的回调是在主线程中

整个过程笔者整理成一张图,一步一步来跟进分析

创建过程
  • 第一步:通过create操作符创建了一个ObservableCreate类型的Observable,由于是基于匿名内部类创建的,因此持有的是实现了ObservableOnSubscribe接口的HomeActivity实例

  • 第二步:通过subscribeOn操作符创建了一个ObservableSubscribeOn类型的Observable,且其内部的source持有上个步骤的ObservableCreate实例

  • 第三步:通过map操作符创建了一个ObservableMap类型的Observable,且其内部持有上个步骤传入的ObservableSubscribeOn实例

  • 第四步:通过observeOn操作符创建了一个ObservableObserveOn类型的Observable,且其内部持有上个步骤的ObservableMap实例

  • 第五步:通过subscribeWith方法完成订阅,由于是基于匿名内部类创建的,因此传入的实际上是实现了ResourceObserverHomeActivity实例

订阅过程

上述的几个步骤其实已经完成的基本的创建过程了,最后我们拿到的实际是ObservableObserveOn的实例,下面开始订阅流程。

  • 第一步:subscribeWith方法,传入的observer是实现了ResourceObserver接口HomeActivity实例,通过subscribeActual发起订阅,内部实际调用的是source.subscribe方法,由于source持有的是上面传入的ObservableMap实例,因此这一步骤实际调用的是,ObservableMap实例中的subscribe方法,传入的参数就是ObserveOnObserver实例(构造参数主要是实现了ResourceObserver的实例即:HomeActivity)

  • 第二步:进入ObservableMap实例subscribe方法中,通过subscribeActual发起订阅,实际调用的是source.subscribe方法,传入的是MapObserver实例(构造参数为之前传递的ObserveOnObserver实例),由于source持有的是ObservableSubscribeOn的实例,因此最终调用的其实是ObservableSubscribeOn实例中的subscribe方法

  • 第三步:进入ObservableSubscribeOn实例subscribe方法中,通过subscribeActual发起订阅,完成MapObserver实例对SubscribeOnObserver的订阅,接着由NewThreadScheduler线程调度器完成对应的任务(该任务的执行是在线程中执行的),SubscribeTask实现了Runnable接口,最终会回调run方法,执行source.subscribe方法,这里的source持有的就是最开始的ObservableCreate实例

@Overridepublic void subscribeActual(final Observer<? super T> s) {final SubscribeOnObserver<T> parent = new SubscribeOnObserver<T>(s);//这里的s就是上个步骤的MapObserver实例s.onSubscribe(parent);//这里的scheduler就是我们最开始指定的Schedulers.newThread 即NewThreadScheduler线程调度器parent.setDisposable(scheduler.scheduleDirect(new SubscribeTask(parent)));
}

  • 第四步:进入ObservableCreate实例subscribe方法中,通过subscribeActual发起订阅,这里的source持有的是HomeActivity实例,直接调用subscribe方法,传入参数是构建的最顶层的发射器CreateEmitter实例

  • 第五步:上述的几个过程实际已经完成了订阅的过程,最后经过层层传递,持有的最顶层的是CreateEmitter实例,即我们最终的被观察者
发射过程

上述的过程已经完成了订阅过程,在最后订阅完成之后,最终会通过source.subscribe方法,其实就是调用HomeActivity实例的subscribe方法,完成元素发射

@Override
public void subscribe(ObservableEmitter<Integer> emitter) throws Exception {for (int i = 0; i < 3; i++) {emitter.onNext(i);}emitter.onComplete();Log.d(TAG, "subscribe " + Thread.currentThread().getName());
}

我们在最顶层的被观察者里通过ObservableEmitter实例onNext方法完成元素的发射,最终又会通过一层一层的Observer转发到最原始的实现了ResourceObserver接口观察者中来

注意:

  • 这里的被观察者里的所有发射过程实际上都是在NewThreadScheduler线程调度器分配的线程里完成的
  • 发射的元素被传递到下层的ObservableObserveOn类中的ObserveOnObserver实例onNext方法,实际执行的是HandlerScheduler.HandlerWorkerschedule方法,最终就是通过我们持有的主线程的handler切换到主线程中

小结

整个创建过程订阅过程发射过程看起来山路十八弯,但是如果你一步一步跟进查看,会发现整个流程实际上是很清晰的,整个过程起点终点很明确,
而中间产生的一系列ObservableObserver你都可以看作是代理类,用来转发订阅以及最终的元素发射

这篇关于RxJava 2.x 之图解创建、订阅、发射流程的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/723731

相关文章

在Ubuntu上部署SpringBoot应用的操作步骤

《在Ubuntu上部署SpringBoot应用的操作步骤》随着云计算和容器化技术的普及,Linux服务器已成为部署Web应用程序的主流平台之一,Java作为一种跨平台的编程语言,具有广泛的应用场景,本... 目录一、部署准备二、安装 Java 环境1. 安装 JDK2. 验证 Java 安装三、安装 mys

Springboot的ThreadPoolTaskScheduler线程池轻松搞定15分钟不操作自动取消订单

《Springboot的ThreadPoolTaskScheduler线程池轻松搞定15分钟不操作自动取消订单》:本文主要介绍Springboot的ThreadPoolTaskScheduler线... 目录ThreadPoolTaskScheduler线程池实现15分钟不操作自动取消订单概要1,创建订单后

JAVA中整型数组、字符串数组、整型数和字符串 的创建与转换的方法

《JAVA中整型数组、字符串数组、整型数和字符串的创建与转换的方法》本文介绍了Java中字符串、字符数组和整型数组的创建方法,以及它们之间的转换方法,还详细讲解了字符串中的一些常用方法,如index... 目录一、字符串、字符数组和整型数组的创建1、字符串的创建方法1.1 通过引用字符数组来创建字符串1.2

SpringCloud集成AlloyDB的示例代码

《SpringCloud集成AlloyDB的示例代码》AlloyDB是GoogleCloud提供的一种高度可扩展、强性能的关系型数据库服务,它兼容PostgreSQL,并提供了更快的查询性能... 目录1.AlloyDBjavascript是什么?AlloyDB 的工作原理2.搭建测试环境3.代码工程1.

Java调用Python代码的几种方法小结

《Java调用Python代码的几种方法小结》Python语言有丰富的系统管理、数据处理、统计类软件包,因此从java应用中调用Python代码的需求很常见、实用,本文介绍几种方法从java调用Pyt... 目录引言Java core使用ProcessBuilder使用Java脚本引擎总结引言python

SpringBoot操作spark处理hdfs文件的操作方法

《SpringBoot操作spark处理hdfs文件的操作方法》本文介绍了如何使用SpringBoot操作Spark处理HDFS文件,包括导入依赖、配置Spark信息、编写Controller和Ser... 目录SpringBoot操作spark处理hdfs文件1、导入依赖2、配置spark信息3、cont

springboot整合 xxl-job及使用步骤

《springboot整合xxl-job及使用步骤》XXL-JOB是一个分布式任务调度平台,用于解决分布式系统中的任务调度和管理问题,文章详细介绍了XXL-JOB的架构,包括调度中心、执行器和Web... 目录一、xxl-job是什么二、使用步骤1. 下载并运行管理端代码2. 访问管理页面,确认是否启动成功

Java中的密码加密方式

《Java中的密码加密方式》文章介绍了Java中使用MD5算法对密码进行加密的方法,以及如何通过加盐和多重加密来提高密码的安全性,MD5是一种不可逆的哈希算法,适合用于存储密码,因为其输出的摘要长度固... 目录Java的密码加密方式密码加密一般的应用方式是总结Java的密码加密方式密码加密【这里采用的

Java中ArrayList的8种浅拷贝方式示例代码

《Java中ArrayList的8种浅拷贝方式示例代码》:本文主要介绍Java中ArrayList的8种浅拷贝方式的相关资料,讲解了Java中ArrayList的浅拷贝概念,并详细分享了八种实现浅... 目录引言什么是浅拷贝?ArrayList 浅拷贝的重要性方法一:使用构造函数方法二:使用 addAll(

解决mybatis-plus-boot-starter与mybatis-spring-boot-starter的错误问题

《解决mybatis-plus-boot-starter与mybatis-spring-boot-starter的错误问题》本文主要讲述了在使用MyBatis和MyBatis-Plus时遇到的绑定异常... 目录myBATis-plus-boot-starpythonter与mybatis-spring-b