过滤操作符
take(2):取前 2 个。1,2,3,4→ 取1,2skip(2):跳过前 2 个。1,2,3,4→ 取3,4first():取第一个takeLast(2):取最后 2 个
变换操作符
flatMap:把发射数据的流映射成新的流,然后合并成一个流(可并发、顺序不保证)。
Observable.just("A", "B")
.flatMap(s -> Observable.just(s + "1", s + "2"))
.subscribe(System.out::println);
// 发射 A1 A2 B1 B2,也可能是 A1 B1 A2 B2,顺序不保证IO线程池和Computation线程池
RxJava Scheduler = Java 线程池 + 调度策略封装 + Observable 链集成 + 统一取消 + 定时支持
Computation线程池(Schedulers.computation())
- 线程数量:固定 = CPU核心数(8核CPU就是8个线程)
- 设计目的:CPU密集型任务—纯计算、不等待外部资源
- 适合场景:
- 定时器(
interval、timer、delay) - 数据排序、过滤、数学运算
- 图片像素处理
- 加密/解密
- 定时器(
- 特点:线程固定,不会无限扩展。因为CPU密集型任务的瓶颈是CPU核心数,线程再多也快不了,反而会增加上下文切换开销
IO线程池(Schedulers.io())
- 线程数量:动态伸缩,按需创建,最多可达几十上百个
- 设计目的:IO密集型任务—大量时间等待,CPU基本空闲
- 适用场景:
- 网络请求(等待服务器响应)
- 数据库读写(等待磁盘IO)
- 文件读写
- 接口调用
- 特点:线程数可以很多,因为IO大部分时间都在阻塞等待,一个线程等待时CPU可以切给其他的线程用
Computation 线程池(固定 8 线程,假设 8 核 CPU):
Thread-1: [计算计算计算计算] ← CPU 一直在跑
Thread-2: [计算计算计算计算]
Thread-3: [定时器滴答滴答滴答]
...8个
IO 线程池(动态伸缩):
Thread-1: [发请求]---等待---等待---[收响应] ← 大部分时间空闲
Thread-2: [发请求]---等待---[收响应]
Thread-3: [读文件]---等待---[读完]
...可以很多个
一个例子解释 map、flatMap、compose
根据用户名查询用户 →获取该用户的订单 →然后计算订单总额度 →显示出来
getUser("张三")
.compose(Rxu.subOnIO())
.flatMap(user -> getOrders(user))
.map(
orders -> {
double total = 0;
for (Order order: orders){
total += order.price;
}
return total;
}
)
.observeOn(AndroidSchedulers.mainThread())
.subscribe(
total -> textView.setText("总金额"+total),
error -> ReportError(error)
);
没有compose(),每次线程转换都单独写
getUser("张三")
.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread())
.subscirbe(...);
getOrders(user)
.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread())
.subscribe(...);getUser("张三")
.subscribeOn(Schedulers.io())
.flatMap(user -> getOrders(user))
.map(
orders -> {
double total = 0;
for (Order order: orders){
total += order.price;
}
return total;
}
)
.observeOn(AndroidSchedulers.mainThread())
.subscribe(
total -> textView.setText("总金额"+total),
error -> ReportError(error)
);用itm()
getUser("张三")
.flatMap(user -> getOrders(user))
.map(
orders -> {
double total = 0;
for (Order order: orders){
total += order.price;
}
return total;
}
)
.compose(Rxu.itm())
.subscribe(
total -> textView.setText("总金额"+total),
error -> ReportError(error)
);