过滤操作符

  • take(2):取前 2 个。1,2,3,4 → 取 1,2
  • skip(2):跳过前 2 个。1,2,3,4 → 取 3,4
  • first():取第一个
  • 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密集型任务—纯计算、不等待外部资源
  • 适合场景:
    • 定时器(intervaltimerdelay
    • 数据排序、过滤、数学运算
    • 图片像素处理
    • 加密/解密
  • 特点:线程固定,不会无限扩展。因为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)
	);