- 文档
- 教程
- 知识库
【免费下载链接】android-tech-frontier
【停止维护】一个定期翻译国外Android优质的技术、开源库、软件架构设计、测试等文章的开源项目
本文基于开源仓库 android-tech-frontier 中收录的《RxJava Essentials》系列译文编写。该系列对应仓库 rxjava 目录下的 chap1 ~ chap8 共 8 篇章节,本文聚焦于第四章 rxjava/chap4.md,系统讲解如何在 Android 应用中过滤 Observable 发射的数据序列。
导读
在 RxJava 的响应式世界里,Observable 源源不断地将数据"推"给订阅者,但很多时候我们并不需要全部数据——只需要开头几个、结尾几个、去重后的值、或者某个时间窗口内的最新值。本篇技术指南以 Android 应用中"已安装应用列表"这一贯穿全书的实战场景为背景,完整讲解filter、take、takeLast、distinct、first、skip、elementAt、sample、timeout、debounce等过滤类操作符的用法与适用边界。读完本文,你将掌握在 RecyclerView 数据流上按条件筛选、限量取值、去重、按时间采样与超时保护等一整套过滤技能,并理解这些操作符在 RxJava 中的底层语义与 Android 项目中的落地写法。
从上一章继承的实战背景
在上一章(chap3)中,我们创建了第一个响应式 Android 应用:用Observable.create()或Observable.from(apps)发射已安装应用列表,并用 RecyclerView + SwipeRefreshLayout 展示。其中的核心数据模型AppInfo如下(见 rxjava/chap3.md):
@Data @Accessors(prefix = "m") public class AppInfo implements Comparable<Object> { long mLastUpdateTime; String mName; String mIcon; // ... }而填充列表的loadList(List<AppInfo> apps)函数则反复出现于本章多个示例中:
private void loadList(List<AppInfo> apps) { mRecyclerView.setVisibility(View.VISIBLE); Observable.from(apps) .subscribe(new Subscriber<AppInfo>() { @Override public void onCompleted() { mSwipeRefreshLayout.setRefreshing(false); } @Override public void onError(Throwable e) { Toast.makeText(getActivity(), "Something went wrong!", Toast.LENGTH_SHORT).show(); mSwipeRefreshLayout.setRefreshing(false); } @Override public void onNext(AppInfo appInfo) { mAddedApps.add(appInfo); mAdapter.addApplication(mAddedApps.size() - 1, appInfo); } }); }本章(rxjava/chap4.md)正是沿着这条主线,研究可观测序列的本质操作——过滤:如何从发射的数据中选取想要的值、如何获取有限个数的值、如何处理重复与溢出的场景。
需要说明一点:原书中匿名订阅者的写法沿用了new Observable<AppInfo>()的笔误,实际 RxJava 中订阅者应为Subscriber<T>(或Observer<T>),下文示例统一使用Subscriber<T>,语义与原书一致。
过滤序列:filter()
RxJava 通过filter()方法过滤序列中不想要的值。它的核心是谓词(predicate):把谓词传给filter(),只有当谓词返回true时,对应的值才会被发射出去,所有观察者才会接收到它。
回到我们的场景:已安装应用列表已经就绪,但我们只想展示以字母C开头的应用。改造loadList()如下(见 rxjava/chap4.md):
private void loadList(List<AppInfo> apps) { mRecyclerView.setVisibility(View.VISIBLE); Observable.from(apps) .filter((appInfo) -> appInfo.getName().startsWith("C")) .subscribe(new Subscriber<AppInfo>() { @Override public void onCompleted() { mSwipeRefreshLayout.setRefreshing(false); } @Override public void onError(Throwable e) { Toast.makeText(getActivity(), "Something went wrong!", Toast.LENGTH_SHORT).show(); mSwipeRefreshLayout.setRefreshing(false); } @Override public void onNext(AppInfo appInfo) { mAddedApps.add(appInfo); mAdapter.addApplication(mAddedApps.size() - 1, appInfo); } }); }与上一章相比,我们只多了一行链式调用:.filter(appInfo -> appInfo.getName().startsWith("C"))。创建 Observable 之后,所有开头字母不是C的元素都会被过滤掉,只有符合条件的AppInfo才会进入onNext()。
Lambda 背后的 Java 7 写法
为了让filter()的语义更清楚,原书给出了不使用 Lambda(Retrolambda)的 Java 7 等价实现。当时的 Android 工程基于 Java 1.6/7,需要借助 Retrolambda 才能在 Android 上使用 Java 8 的 Lambda 语法(这一点在 rxjava/chap3.md 的工具一节有说明)。Java 7 写法如下:
.filter(new Func1<AppInfo, Boolean>() { @Override public Boolean call(AppInfo appInfo) { return appInfo.getName().startsWith("C"); } })这里的关键点是Func1:它是 RxJava 中"只有一个参数"的函数接口,泛型第一个参数AppInfo是输入类型,第二个参数Boolean是返回值类型。filter()对每一个发射值调用Func1.call(),返回true则放行,返回false则丢弃。
最常用的场景:过滤 null
filter()最常见的用法之一是过滤null对象:
.filter(new Func1<AppInfo, Boolean>() { @Override public Boolean call(AppInfo appInfo) { return appInfo != null; } })这个写法看起来简单,甚至显得有些模板代码化,但它的价值在于:帮我们免去了在onNext()调用里逐个判断null的麻烦,把注意力集中在应用业务逻辑上,让"数据是否有效"这一关注点从消费者侧前移到数据流管道中。
正如原书所强调的:filter()让我们无需关心可观测序列的来源、也无需关心它为什么发射了这么多不同的数据,我们只是想要其中某个子集来创建应用中使用的新序列。这种"从已有序列派生所需序列"的思路,促进了编码中的分离性与抽象性。
获取我们需要的数据:take() 与 takeLast()
当我们不需要整个序列,而只想取开头或结尾的几个元素时,take()和takeLast()就派上用场了。
take(N):取前 N 个
如果我们只想要序列中的前三个元素,发射它们然后让 Observable 完成,可以这样写:
private void loadList(List<AppInfo> apps) { mRecyclerView.setVisibility(View.VISIBLE); Observable.from(apps) .take(3) .subscribe(new Subscriber<AppInfo>() { // onCompleted / onError / onNext 同前 }); }take()接收一个整数 N 作为参数,从原始序列中发射前 N 个元素,然后立刻完成(onCompleted())。它的语义可以概括为:只关心数据流的开头。
takeLast(N):取最后 N 个
如果想要最后 N 个元素,只需使用takeLast():
Observable.from(apps) .takeLast(3) .subscribe(...);注意:takeLast()听起来简单,但它有一个重要前提——由于必须知道"最后 N 个"是什么,它只能作用于一个已经完成(有限)的序列。也就是说,takeLast()会等待源序列发射完毕,缓存全部数据后再挑出末尾 N 个发射。对无限流(如interval()产生的流)使用takeLast()是没有意义的。
有且仅有一次:distinct() 与 distinctUntilChanged()
一个可观测序列可能会在出错时重复发射,或者被设计成重复发射。distinct()和distinctUntilChanged()就是用来处理这类重复问题的。
distinct():全局去重
如果我们想对一个指定的值只处理一次,可以对序列应用distinct()去掉所有重复项。与takeLast()类似,distinct()也作用于完整序列,它需要记录每一个发射过的值才能判断后续值是否重复。因此,处理很大的序列或大批量数据时,务必关注内存使用情况——去重集合会随序列规模增长。
原书用一个练习把本章此前学到的创建操作符串了起来:先take(3)取出一小组可识别的数据项,再用repeat(3)创建一个有重复的更大序列,最后用distinct()去重:
Observable<AppInfo> fullOfDuplicates = Observable.from(apps) .take(3) .repeat(3);fullOfDuplicates把已安装应用的前三个重复了 3 次:共 9 个元素且大量重复。接着:
fullOfDuplicates.distinct() .subscribe(new Subscriber<AppInfo>() { @Override public void onCompleted() { mSwipeRefreshLayout.setRefreshing(false); } @Override public void onError(Throwable e) { Toast.makeText(getActivity(), "Something went wrong!", Toast.LENGTH_SHORT).show(); mSwipeRefreshLayout.setRefreshing(false); } @Override public void onNext(AppInfo appInfo) { mAddedApps.add(appInfo); mAdapter.addApplication(mAddedApps.size() - 1, appInfo); } });结果很直观:9 个元素经distinct()后只保留前 3 个不重复的AppInfo。
distinctUntilChanged():相邻去重
如果一个序列发射的新值与前一个值不同时才通知我们,该怎么处理?原书用温度传感器举例:它每秒发射一次室内温度:
21° 21° 21° 21° 22° ...如果每次拿到新值都去更新温度显示,而绝大多数情况下温度并没有变化,就会造成不必要的 UI 刷新。我们希望忽略连续重复的值,只有在温度确实改变时才收到通知——这正是distinctUntilChanged()的用武之地:它轻易地忽略掉所有相邻的重复,只发射出新的值。
注意它与distinct()的本质区别:
distinct():全局去重,需要记住整个历史发射值,内存开销与序列长度成正比;distinctUntilChanged():只与上一个值比较,几乎无额外内存开销,适用于"变化才通知"的场景(传感器、输入框内容等)。
first 与 last:只看一头一尾
first()和last()方法很容易理解:它们从 Observable 中只发射第一个元素或最后一个元素,发射后序列即完成。
这两个方法都可以传入一个Func1作为谓词参数,用于确定"我们感兴趣的"第一个或最后一个元素——即满足条件的第一个/最后一个。例如,first(predicate)会从序列中找出第一个满足谓词条件的值发射,last(predicate)则找出最后一个满足条件的值。
与它们相对应的变体是firstOrDefault()和lastOrDefault():当可观测序列完成时仍然没有发射任何(符合条件的)值时,这两个函数会发射一个我们预先指定的默认值,避免消费者拿到一个空序列。这在"序列可能为空、但下游必须有值可消费"的场景下非常实用。
skip 与 skipLast:跳过一头一尾
skip()和skipLast()与take()、takeLast()正好对应:它们接收整数 N 作为参数,不让 Observable 发射前 N 个或后 N 个值,序列中其余的值照常发射。
skip(N):跳过前 N 个元素,发射它后面的所有数据;skipLast(N):跳过最后 N 个元素,从源序列中发射剩下的其他元素。
当我们知道一个序列以"没有太多用处的、可控的"元素开头或结尾时(例如日志流的前几条初始化记录、列表末尾的占位数据),可以用这两个操作符把它们剔除掉。
elementAt:只取第 N 个
如果我们只想要可观测序列发射的第 N 个元素,elementAt()正合适:它从序列中发射第 n 个元素,然后立刻完成。例如elementAt(2)会选择序列中的第三个元素(下标从 0 开始),并创建一个只发射该指定元素的新 Observable。
如果序列长度不足——比如想找第五个元素但源序列只有三个元素——就会触发onError()。为了避免这种情况,可以使用elementAtOrDefault():当序列在到达指定位置前就结束时,发射一个默认值。这一"找不到就兜底"的语义与firstOrDefault()、lastOrDefault()一脉相承。
采样:sample() 与 throttleFirst()
回到温度传感器的例子。它每秒发射一次当前室内温度,但温度并不会变化得那么快,我们完全可以用一个更大的发射间隔。在 Observable 后面追加sample(),就能创建一个新的可观测序列:在指定的时间间隔内,由源 Observable 发射最近一次的数值。
Observable<Integer> sensor = [...]; sensor.sample(30, TimeUnit.SECONDS) .subscribe(new Subscriber<Integer>() { @Override public void onCompleted() { } @Override public void onError(Throwable e) { } @Override public void onNext(Integer currentTemperature) { updateDisplay(currentTemperature); } });这个例子中,新 Observable 每隔 30 秒观测一次温度传感器,发射最近一个温度值。sample()支持全部时间单位:毫秒、秒、分、天等等。
需要留意的是,sample()属于时间驱动型操作符,它依赖一个时间调度器。在 rxjava/chap7.md 的 Schedulers 一节明确指出:computation()是buffer()、debounce()、delay()、interval()、sample()、skip()等众多操作符的默认调度器,专门用于计算型工作。因此直接使用sample()时,采样逻辑会运行在 computation 线程池上。
如果我们想让定时器发射的是第一个元素(时间段开始时的那一个)而不是最近的一个,可以使用throttleFirst()——它与sample()(即throttleLast())形成互补,一个取首、一个取尾。
timeout():超时保护
假设我们工作在一个时效性要求很高的环境:温度传感器每秒都发射一个温度值,而我们要求它每隔两秒至少发射一个。此时可以用timeout()函数监听源可观测序列——如果在我们设定的时间间隔内没有收到值,就发射一个错误。
可以把timeout()理解为"Observable 的限时副本":如果在指定的时间间隔内源 Observable 不发射值,timeout()就会触发onError()。
Subscription subscription = getCurrentTemperature() .timeout(2, TimeUnit.SECONDS) .subscribe(new Subscriber<Integer>() { @Override public void onCompleted() { } @Override public void onError(Throwable e) { Log.d("RXJAVA", "You should go check the sensor, dude"); } @Override public void onNext(Integer currentTemperature) { updateDisplay(currentTemperature); } });timeout()一旦超时触发onError(),整个序列就终止了:即便源数据稍后才到达,那个"迟到的元素"也不会再被发射出来。在 rxjava/chap7.md 中还可以看到,timeout()的默认调度器是Schedulers.immediate()(在当前线程立即执行)。
与sample()相同,timeout()也使用TimeUnit对象来指定时间间隔。这个操作符非常适合网络请求、传感器轮询、心跳检测等"必须在规定时间内有数据"的场景。
debounce():防抖
debounce()过滤掉由 Observable 发射的速率过快的数据:如果在一个指定的时间间隔过去后仍旧没有新的数据发射,它才发射最后那一个;反之,如果数据源源不断地快速到达,debounce()就会不断重置内部定时器,直到数据流出现"安静窗口"才放行。
就像sample()和timeout()一样,debounce()使用TimeUnit对象指定时间间隔。它的内部机制是:每收到一个数据就开启一个内部定时器,如果在这个时间间隔内没有新的数据发射,新 Observable 就发射出最后一个数据。
debounce()的典型用途是 UI 输入防抖——例如搜索框输入时,用户停止输入 300 毫秒后才真正发起搜索请求,避免每敲一个键就触发一次网络请求。它与sample()的区别在于:
sample():按固定节奏定时取样(无论数据来得多频繁);debounce():按数据间隙决定是否发射(只有数据"安静"下来才发射)。
总结
本章(rxjava/chap4.md)围绕"过滤可观测序列"这一主题,覆盖了三大类过滤手段:
| 分类 | 操作符 | 用途 |
|---|---|---|
| 条件过滤 | filter() | 按谓词(Func1<T,Boolean>)筛选元素,如过滤 null、按前缀筛选 |
| 限量与定位 | take()/takeLast()、skip()/skipLast()、elementAt()/elementAtOrDefault()、first()/last() | 只取开头/结尾/N 个,跳过开头/结尾,按位置取值 |
| 去重 | distinct()(全局去重,注意内存)、distinctUntilChanged()(相邻去重) | 消除重复发射 |
| 时间驱动 | sample()(定时取最近值)、throttleFirst()(定时取首个值)、timeout()(超时报错)、debounce()(防抖) | 控制发射节奏、限时保护 |
连同上一章学到的创建操作符(create、from、just、repeat、range、interval、timer等,见 rxjava/chap3.md),我们现在已经可以用filter()、skip()、sample()等手段从原始数据流中裁剪出任何想要的 Observable。
在下一章(rxjava/chap5.md)中,我们将学习如何变换一个序列:把函数应用到每个元素上(map家族)、给它们分组和扫描(scan),从而创建出能完成目标的特定 Observable。若想深入了解过滤操作符背后的线程调度原理,可继续阅读 rxjava/chap7.md 中关于Schedulers.computation()、Schedulers.immediate()与各操作符默认调度器的对应关系。
- 文档
- 教程
- 知识库
【免费下载链接】android-tech-frontier
【停止维护】一个定期翻译国外Android优质的技术、开源库、软件架构设计、测试等文章的开源项目
相关推荐
RxJava 过滤算子实战指南:filter、debounce、throttle 与 timeout 全解析
RxJava 过滤算子实战指南:filter、debounce、throttle 与 timeout 全解析 本文围绕 RxJava 官方文档 Filterin
后端异步编程android-tech-frontier《RxJava 开发精要》第 5 章:Observables 变换操作符全解析——从 map 到 cast 的序列塑形实战
android tech frontier《RxJava 开发精要》第 5 章:Observables 变换操作符全解析——从 map 到 cast 的序列塑形
文档教程知识库Nazara Engine内存管理策略:高效资源加载与缓存机制终极指南
Nazara Engine内存管理策略:高效资源加载与缓存机制终极指南 Nazara Engine作为一款跨平台的实时应用程序框架,其 内存管理策略 和 资源缓
游戏开发图形学音视频
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考