RxJava基础操作符和高级操作符
RxJava1. 理解核心概念与基础操作1.1 核心四要素 (The Observable Contract)1.2 冷 Observable 与 热 Observable1.3 基础操作符 (Operators)1.3.1 创建操作符 (Creating):1.3.2 变换操作符 (Transforming):1.3.3 过滤操作符 (Filtering):1.3.4 组合操作符 (Combining):1.4. 线程调度 (Schedulers)2. 深入原理与高级操作2.1 操作符链与 Lift 原理2.2 背压 (Backpressure) - RxJava 2/3 的核心概念2.3. 高级操作符2.4. Subject (主题)2.5. 整体变换 (compose 和 transform)1. 理解核心概念与基础操作这是学习 RxJava 的基石必须牢固掌握。1.1 核心四要素 (The Observable Contract)RxJava 的核心是 观察者模式 (Observer Pattern)它定义了四个关键角色•Observable(被观察者/数据源): 产生数据流的源头。可以发射 多个数据项以成功或失败结束。•Observer(观察者/订阅者): 接收并处理 Observable 发射的数据。它定义了四个方法• onSubscribe(Disposable d): 订阅建立时调用可用于取消订阅。• onNext(T value):接收到一个数据项时调用。• onError(Throwable e): 发生错误时调用流终止。• onComplete():流成功完成时调用流终止。•Subscription(订阅): 连接 Observable 和 Observer 的纽带。通过调用 subscribe() 方法建立。•Disposable(可dispose的): 代表一个订阅关系调用其 dispose() 方法可以取消订阅停止接收 数据并释放资源。Hello World 示例:import io.reactivex.rxjava3.core.Observable; import io.reactivex.rxjava3.disposables.Disposable; public class RxJavaHelloWorld { public static void main(String[] args) { // 1. 创建 Observable (数据源) Observable String observable Observable.just(Hello, RxJava!); // 2. 创建 Observer (观察者) MyObserver observer new MyObserver(); // 3. 建立订阅 (Subscription) Disposable disposable observable.subscribe(observer); // 4. (可选) 取消订阅 // disposable.dispose(); } static class MyObserver implements io.reactivex.rxjava3.core.Observer String { Override public void onSubscribe(io.reactivex.rxjava3.disposables.Disposable d) { System.out.println(订阅建立); } Override public void onNext(String s) { System.out.println(接收到: s); } Override public void onError(Throwable e) { System.err.println(发生错误: e.getMessage()); } Override public void onComplete() { System.out.println(流已结束); } } } // 输出: // 订阅建立 // 接收到: Hello // 接收到: RxJava! // 流已结束1.2 冷 Observable 与 热 Observable• 冷 Observable (Cold Observable): 只有在有订阅者订阅时才开始发射数据。每个订阅者都会收到完整的数据流。just, fromArray, create 等创建的通常是冷的。• 热 Observable (Hot Observable): 无论是否有订阅者都会发射数据。订阅者只能接收到订阅之后发射的数据。PublishSubject, ReplaySubject 等创建的是热的。1.3 基础操作符 (Operators)操作符是 RxJava 的灵魂用于对数据流进行转换、过滤、组合等。1.3.1 创建操作符 (Creating):• just(T…): 发射指定的项。• fromArray(T[]) / fromIterable(Iterable? extends T): 从数组或集合发射数据。• create(ObservableOnSubscribe): 通过自定义逻辑创建 Observable。• interval(long, TimeUnit): 定时发射递增的数字。• range(int start, int count): 发射一个范围内的整数序列。1.3.2 变换操作符 (Transforming):• map(FunctionT, R): 将每个发射的项通过一个函数转换成另一个值。Observable.just(1, 2, 3 ) .map(Integer::parseInt) // 将 String 转为 Integer .subscribe(System.out::println); // 输出: 1, 2, 3• flatMap(FunctionT, ObservableSource? extends R): 将每个发射的项转换成一个新的 Observable然后将这些 Observable 发射的数据“拍平”并合并到一个单一的 Observable 中。常用于异步操作。Observable.just(url1, url2) .flatMap(url - fetchDataFromNetwork(url)) // 假设 fetchDataFromNetwork 返回 ObservableData .subscribe(data - System.out.println(Received: data));• concatMap(FunctionT, ObservableSource? extends R): 类似 flatMap但保证内部 Observable 的发射顺序与原始项的顺序一致。1.3.3 过滤操作符 (Filtering):• filter(Predicate): 只发射满足指定条件的项。Observable.just(1, 2, 3, 4, 5) .filter(x - x % 2 0) // 只保留偶数 .subscribe(System.out::println); // 输出: 2, 4• take(long count): 只取前 N 项。• skip(long count): 跳过前 N 项。1.3.4 组合操作符 (Combining):• merge(ObservableSource? extends T…): 将多个 Observable 发射的数据合并不保证顺序。• concat(ObservableSource? extends T…): 按顺序连接多个 Observable前一个完成后才订阅下一个。• zip(ObservableSource, ObservableSource, BiFunctionT1, T2, R): 将多个 Observable 的数据按顺序一一配对组合。ObservableString names Observable.just(Alice, Bob); ObservableInteger ages Observable.just(25, 30); Observable.zip(names, ages, (name, age) - name is age years old.) .subscribe(System.out::println); // 输出: // Alice is 25 years old. // Bob is 30 years old.• combineLatest(ObservableSource, ObservableSource, BiFunctionT1, T2, R): 当任意一个 Observable 发射新数据时与其它 Observable 的最新数据组合。1.4. 线程调度 (Schedulers)RxJava 强大的异步能力离不开调度器。• subscribeOn(Scheduler scheduler): 指定 上游数据产生、操作符执行在哪个线程执行。只对第一个 subscribeOn有效。• observeOn(Scheduler scheduler): 指定下游onNext, onError, onComplete在哪个线程执行。可以多次调用,切换线程。常用调度器:• Schedulers.io(): 用于 I/O 密集型操作如网络请求、文件读写。• Schedulers.computation(): 用于 CPU 密集型计算。• Schedulers.newThread(): 为每个任务创建一个新线程。• AndroidSchedulers.mainThread() (RxAndroid): 在 Android 主线程执行用于更新 UI。示例:Observable.just(url) .subscribeOn(Schedulers.io()) // 网络请求在 IO 线程 .map(url - fetchData(url)) // 转换操作也在 IO 线程 .observeOn(AndroidSchedulers.mainThread()) // 切换到主线程更新 UI .subscribe(data - updateUI(data));2. 深入原理与高级操作掌握了基础后需要理解其内部机制和更复杂的操作。2.1 操作符链与 Lift 原理理解操作符是如何串联起来的至关重要。每个操作符如 map, filter内部通常会调用 lift(Operator) 方法创建一个新的 Observable这个新的 Observable 会包装上游的 Observable并在订阅时将下游的 Observer 包装成一个新的 Observer从而在数据传递过程中插入转换或过滤逻辑。这是一个典型的装饰者模式。2.2 背压 (Backpressure) - RxJava 2/3 的核心概念在异步场景下如果生产者Observable发射数据的速度远快于消费者Observer处理数据的速度会导致内存溢出。RxJava 2 引入了 Flowable 来专门处理背压问题。• Flowable vs Observable:• Observable: 无背压支持适用于 GUI 事件、短序列等。• Flowable: 有背压支持适用于网络请求、数据库操作、文件 I/O 等可能产生大量数据的场景。• 背压策略 (BackpressureStrategy):• BUFFER: 缓存所有数据可能导致 OOM。• DROP: 丢弃无法处理的数据。• LATEST: 只保留最新的数据丢弃旧的。• ERROR: 抛出 MissingBackpressureException。• MISSING: 不做任何处理需要手动通过 onBackpressureXXX 操作符处理。2.3. 高级操作符• switchMap: 类似 flatMap但当源 Observable 发射一个新项时会取消订阅并丢弃前一个内部 Observable 产生的数据只处理最新的内部 Observable。• debounce: 如果在一个指定的时间段内没有新的数据发射则发射最近的一个数据。常用于搜索框防抖。• distinct / distinctUntilChanged: 过滤掉重复的数据。• retry / retryWhen: 在发生错误时重试。• timeout: 如果在指定时间内没有收到数据则抛出超时异常。• doOnNext / doOnError / doOnComplete / doFinally: 用于副作用side-effect如日志记录、资源清理不影响数据流本身。2.4. Subject (主题)Subject 既是 Observable 又是 Observer可以手动向数据流中“推送”数据是创建热 Observable 的主要方式。• PublishSubject: 向所有订阅者广播数据订阅后才能收到数据。• ReplaySubject: 缓存所有发射过的数据新的订阅者会收到所有历史数据。• BehaviorSubject: 缓存最后一个数据新的订阅者会立即收到这个最新数据然后接收后续数据。• AsyncSubject: 只在流完成时向所有订阅者发射最后一个数据。BehaviorSubject 示例 (模拟状态管理):BehaviorSubjectString userNameSubject BehaviorSubject.createDefault(Guest); // 订阅者 A userNameSubject.subscribe(name - System.out.println(User A sees: name)); // 输出: User A sees: Guest // 订阅者 B (稍后订阅) userNameSubject.subscribe(name - System.out.println(User B sees: name)); // 输出: User B sees: Guest // 更新用户名 userNameSubject.onNext(Alice); // 输出: User A sees: Alice // 输出: User B sees: Alice2.5. 整体变换 (compose 和 transform)当一组操作符需要在多个地方复用时可以使用 compose 或 transform。• compose(ObservableTransformerT, R): 用于 Observable。• to(ObservableConverterT, R): 更通用的转换。示例- 封装网络请求的通用线程切换:public class SchedulersTransformerT implements ObservableTransformerT, T { Override public ObservableSourceT apply(ObservableT upstream) { return upstream.subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()); } } // 使用 apiService.getData() .compose(new SchedulersTransformer()) // 应用通用变换 .subscribe(...);

相关新闻

RPCS3模拟器深度优化指南:让你的PS3游戏在PC上流畅运行

RPCS3模拟器深度优化指南:让你的PS3游戏在PC上流畅运行

RPCS3模拟器深度优化指南:让你的PS3游戏在PC上流畅运行 【免费下载链接】rpcs3 PlayStation 3 emulator and debugger 项目地址: https://gitcode.com/GitHub_Trending/rp/rpcs3 你是否曾经梦想在PC上重温《最后生还者》、《神秘海域》或《战神3》等PS3经典大…

2026/9/18 6:45:22 阅读更多 →
SPI并行模式原理与实战:突破嵌入式通信带宽瓶颈

SPI并行模式原理与实战:突破嵌入式通信带宽瓶颈

1. SPI并行模式:从串行瓶颈到高速通道的演进搞嵌入式通信的工程师,对SPI(Serial Peripheral Interface)肯定不陌生。作为板上芯片间通信的“老将”,它的四线制(SCLK, MOSI, MISO&…

2026/9/17 9:20:59 阅读更多 →
PyPtt单元测试教程:确保你的PTT脚本稳定可靠

PyPtt单元测试教程:确保你的PTT脚本稳定可靠

PyPtt单元测试教程:确保你的PTT脚本稳定可靠 【免费下载链接】PyPtt The best PTT library 项目地址: https://gitcode.com/gh_mirrors/py/PyPtt PyPtt是一款强大的PTT(批踢踢)库,为开发者提供了与PTT论坛交互的丰富功能。…

2026/9/18 13:08:18 阅读更多 →

最新新闻

OpenHands 实战:TaoToken 跑通 SWE-bench Verified 全流程

OpenHands 实战:TaoToken 跑通 SWE-bench Verified 全流程

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/19 16:55:36 阅读更多 →
用 OpenDesign 复刻 Airtable 设计系统:从视觉规范到语义化 Design Token 的完整落地指南

用 OpenDesign 复刻 Airtable 设计系统:从视觉规范到语义化 Design Token 的完整落地指南

用 OpenDesign 复刻 Airtable 设计系统:从视觉规范到语义化 Design Token 的完整落地指南 【免费下载链接】open-design 🎨 Best DeepSeek Harness Design Plugin. The open-source Claude Design alternative. 🖥️ Local-first desktop app…

2026/9/19 16:55:36 阅读更多 →
Cline 报 401?TaoToken 这样核对模型 ID 和 Base URL

Cline 报 401?TaoToken 这样核对模型 ID 和 Base URL

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/19 16:55:36 阅读更多 →
Caffe EuclideanLoss 层完全指南:平方和(L2)回归损失层的原理、配置与实战

Caffe EuclideanLoss 层完全指南:平方和(L2)回归损失层的原理、配置与实战

Caffe EuclideanLoss 层完全指南:平方和(L2)回归损失层的原理、配置与实战 【免费下载链接】caffe Caffe: a fast open framework for deep learning. 项目地址: https://gitcode.com/gh_mirrors/ca/caffe 本文以 Caffe 官方教程 docs…

2026/9/19 16:55:36 阅读更多 →
QuickRecorder:macOS 录屏快速上手指南

QuickRecorder:macOS 录屏快速上手指南

QuickRecorder:macOS 录屏快速上手指南 【免费下载链接】QuickRecorder A lightweight screen recorder based on ScreenCapture Kit for macOS / 基于 ScreenCapture Kit 的轻量化多功能 macOS 录屏工具 项目地址: https://gitcode.com/GitHub_Trending/qu/Quick…

2026/9/19 16:55:36 阅读更多 →
Roc 语言 Try.map_both 实战:从 REPL 快照测试看懂 Ok/Err 双分支映射语义

Roc 语言 Try.map_both 实战:从 REPL 快照测试看懂 Ok/Err 双分支映射语义

Roc 语言 Try.map_both 实战:从 REPL 快照测试看懂 Ok/Err 双分支映射语义 【免费下载链接】roc A fast, friendly, functional language. 项目地址: https://gitcode.com/GitHub_Trending/ro/roc 本篇技术指南以 Roc 仓库中的 REPL 快照测试 test/snapshots…

2026/9/19 16:54:36 阅读更多 →

日新闻

BP神经网络时序预测:滑窗长度与多窗口平均策略

BP神经网络时序预测:滑窗长度与多窗口平均策略

简介:面向机器学习、深度学习与数据建模学习者的一份完整研究文献,聚焦BP神经网络在农业产量预测中的应用。文档以1980—2018年全国棉花产量为样本,系统讲解数据归一化处理、激活函数原理、多层神经网络结构搭建及训练流程,展示敏…

2026/9/19 0:00:30 阅读更多 →
Transformer训练实时监控实战:基于MindSpore的损失曲线可视化方案

Transformer训练实时监控实战:基于MindSpore的损失曲线可视化方案

上个月调一个Deformable DETR模型,在单卡上要跑将近两天。第二天早上我下意识打开终端翻日志,发现loss从凌晨两点就开始往上爬,一路从0.8涨到1.35,整整六个小时没人发现。那六个小时的训练不仅白跑,还霸占着卡——等于…

2026/9/19 0:00:30 阅读更多 →
OpenCloud 中的 Go 类型安全转换库 spf13/cast:从零值回退到泛型 API 的完整实战指南

OpenCloud 中的 Go 类型安全转换库 spf13/cast:从零值回退到泛型 API 的完整实战指南

OpenCloud 中的 Go 类型安全转换库 spf13/cast:从零值回退到泛型 API 的完整实战指南 【免费下载链接】opencloud 🌤️ OpenCloud is the open source platform for file management, sharing and collaboration. Simple and sovereign. 项目地址: htt…

2026/9/19 0:00:30 阅读更多 →

周新闻

AI SDK Harness 依赖更新指南:掌握 harness 包 SDK 依赖的升级、桥接同步与一致性校验

AI SDK Harness 依赖更新指南:掌握 harness 包 SDK 依赖的升级、桥接同步与一致性校验

AI SDK Harness 依赖更新指南:掌握 harness 包 SDK 依赖的升级、桥接同步与一致性校验 【免费下载链接】ai The AI Toolkit for TypeScript. From the creators of Next.js, the AI SDK is a free open-source library for building AI-powered applications and ag…

2026/9/19 3:59:36 阅读更多 →
Refine v5 Ant Design NumberField 组件实战:基于 Intl 的本地化数字格式化

Refine v5 Ant Design NumberField 组件实战:基于 Intl 的本地化数字格式化

Refine v5 Ant Design NumberField 组件实战:基于 Intl 的本地化数字格式化 【免费下载链接】refine A React Framework for building internal tools, admin panels, dashboards & B2B apps with unmatched flexibility. 项目地址: https://gitcode.com/GitH…

2026/9/19 3:53:08 阅读更多 →
Flutter应用改名全指南:从Android到iOS的配置与工具实践

Flutter应用改名全指南:从Android到iOS的配置与工具实践

刚接一个外包项目时,甲方要求把工程里临时用的应用名改成正式产品名。我本来觉得“改名”这种小事,打开配置文件改一行不就完了?结果真动手才发现,Flutter项目里“应用名称”根本不是一处配置,而是一整套散落在 Androi…

2026/9/19 4:02:43 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/16 22:31:27 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/15 21:39:18 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/16 22:32:59 阅读更多 →