性能优化-并发、异步和响应式
并发、异步和响应式
现代程序语言和提供的类库,使得基于多线程的编程得很容易,多线程编程充分利用计算机硬件提供的多个CPU增加系统吞吐量。多线程编程衍生出并发、异步,异步编排和响应式编程风格 。每个风格完成的的需求不一样
| 风格 | 描述 |
|---|---|
| 并发编程 | 充分利用CPU等硬件资源,启用多个线程,并发执行多个请求,不关心调用结果。可以创建Thread和线程池ThreadPoolExecutor |
| 异步 | 请求执行过程中,遇到阻塞调用,如数据库查询或者微服调用,可以把这种阻塞调封装在线程里执行。异步调用结果可通过回调或者Future俩种方式获取 |
| 异步编排 | 多个异步调用按照一定顺序执行,比如请求A和请求B 都执行完毕后,根据俩个的结果再异步执行请求C, 使用CompletableFuture |
| 响应式 | 同架构的事件驱动风格,Producer连续生成一组有限或者无限的数据,经过Consumer的计算分析并汇总,Consumer可以允许并行或者串行处理数据。响应式支持Consumer一种可控速率消费Producer的数据(背压)。支持响应式的有JDK9的java.util.concurrent.Flow和开源工具RxJava,Spring Reactor等 |
无论是并发编程、是异步编程或者是响应式通常不会直接创建Thread执行,这是因为Java的Thread与操作系统的线程一一映射,创建和维护Thread会非常消耗系统资源的,比如
- 一个线程默认占用1M的堆外内存,如果系统有1000个线程,则占用了堆外1G。通过
‑Xss设定占用内存大小,也可以通过pmap操作系统命令查看线程实际使用大小 - Java线程要维护线程的栈帧(参考5.2)等信息,对应的操作系统线程还有维护内核栈,寄存器上下文等信息.线程上下文切换开销巨大。
可见Java Thread是个重量级对象,通常使用的方式是创建一个弹性大小的线程池,如下是ThreadPoolExecutor的构造函数定义
ThreadPoolExecutor(int corePoolSize,
int maximumPoolSize,
long keepAliveTime,
TimeUnit unit,
BlockingQueue<Runnable> workQueue,//任务队列
ThreadFactory threadFactory,
RejectedExecutionHandler rejectHandler)//拒绝策略这里有几个参数关系到线程池的吞吐量和系统性能
| 参数 | 描述 |
|---|---|
| corePoolSize | 核心线程数目,是线程池的的最主要参数,设定了corePoolSize用于并发处理请求 |
| workQueue | 如果所有核心线程线程都在忙,无法处理新的任务(TODO,任务和请求术语统一),则新增的任务将放到队列workQueue |
| maximumPoolSize | 如果workQueue已经满,则ThreadPoolExecutor将创建更多线程执行,现成个数达到maximumPoolSize |
| rejectHandler | 如果workQueue满,且已经达到maximumPoolSize个数,则新增任务交被拒绝,交给rejectHandler处理 |
在使用ThreadPoolExecutor的时候,最主要评估corePoolSize大小以提高并发处理量,如果corePoolSize过大,占用系统资源,且无用线程太多不利于监控工具展示。如果过小,则并发处理能力不足。
还需要设置的maximumPoolSize和workQueue,以应对流量的突变,如果workQueue过小会触发拒绝策略,如果workQueue过大或者无限制容量,则有可能导致内存溢出,或者数分钟后才能拿执行的任务已经过期。
如下使用VisualVM连接JVM后,观察线程池pool-1-thread的运行状况,其中pool-1-thread-1 一直处于等待状态,右侧Running 一列显示时间为0 。其他pool-1-thread线程则正常运行,状态是Running和Park切换中

注意: 线程有多个状态,定义在java.lang.Thread.State ,这将在下一节Lock中说明。上图绿色表示JVM判断线程正常运行相当于处于NEW或者RUNNABLE状态。
线程池每个参数设置的合理值需要通过可观测系统长期观察确定。良好的设置应该保证正常流量下使用corePoolSize 即可处理所有并发,而不需要队列缓冲大量任务。 缓冲大量任务即可能造成内存溢出,有可能因为演示处理任务导致处理结果往往过期(比如请求端已经超时退出,而任务还缓冲在线程池队列中)
有了线程池,可以用其submit方法实现异步调用并在需要时候获取调用结果,定义如下
//ThreadPoolExecutor的submit方法
public <T> Future<T> submit(Callable<T> task)
//Callable定义
public interface Callable<V> {
V call() throws Exception;
}如下是使用线程池实现异步调用,并获取返回结果例子
Future<Integer> future = pool.submit(() -> {
System.out.println("异步任务,模拟一个RPC调用");
return rpcCall();
});
System.out.println("主线程继续处理其他任务");
//获取调用结果
Integer result = future.get();也可以使用CompletableFuture来实现异步调用和异步编排,CompletableFuture提供了四个静态方法来创建异步任务:
public static CompletableFuture<Void> runAsync(Runnable runnable)
public static CompletableFuture<Void> runAsync(Runnable runnable
, Executor executor)
public static <U> CompletableFuture<U> supplyAsync(Supplier<U> supplier)
public static <U> CompletableFuture<U> supplyAsync(Supplier<U> supplier
, Executor executor)runAsync方法接收一个异步任务Runnable,不需要关心执行结果,supplyAsync则需要提供一个Supplier实现,Supplier返回一个执行结果, 如下是Supplier的定义
public interface Supplier<T> {
T get();
}CompletableFuture的所有xxxAsync方法都提供了可选的Executor参数,如果调用带有Executor参数,则CompletableFuture使用Executor来执行异步任务,Executor是一个接口,线程池ThreadPoolExecutor就是Executor的一个实现,因此,我们在使用CompletableFuture,可以传入一个线程池。这也是最常用得使用方式,如果没有传入,则使用CompletableFuture提供的一个默认线程池,不推荐使用默认线程池,因为不容易被系统管理(本书所有知识点的高性能和高可用,都不推荐默认参数)。
//无返回值
public static void runAsync() throws Exception {
CompletableFuture<Void> future = CompletableFuture.runAsync(() -> {
System.out.println("运行ing");
});
}
//有返回值
public static void supplyAsync() throws Exception {
CompletableFuture<Long> future = CompletableFuture.supplyAsync(() -> {
return 1l;
});
long ret = future.get();
System.out.println("ret = "+ret);
}如果期望俩个异步调用并行执行完毕才执行下一步操作,可以调用allOf方法,比如(TODO,来自豆包的例子,最好改一下)
CompletableFuture<String> f1 = supplyAsync(() -> rpcA());
CompletableFuture<String> f2 = supplyAsync(() -> rpcB());
CompletableFuture<String> f3 = supplyAsync(() -> rpcC());
CompletableFuture.allOf(f1,f2,f3).join();如果要处理持续的数据流或者请求,并期望实现高性能和高可用。 可以采用响应式编程,在JDK9中,concurrent.Flow支持响应式,或者使用开源的RxJava,Spring Reactor类库等。Akka则实现了分布式的响应式编程。响应式编程的主要目标是通过异步方式,给与客户端请求快速的响应并支持流量控制(背压)。可以用在响应式的场景有
- API网关实现,比如 Spring Gateway 。网关设计为不会被下游处理能力慢而阻塞自身
- 物联网应用,比如设备告警持续上报,除了需要处理告警业务外,也具备流量控制,考虑到如果设备告警未恢复,会一直告警。因此设置只处理最新的告警上报
- 复杂的异步编排,比如搜索场景,当用户有进一步输入的时候,取消上一次的搜索
- R2DBC ,响应式关系数据库连接,对标JDBC,不同的是JDBC是同步阻塞操作,R2DBC通常与Spring‑WebFlux构成响应式的web应用
下面代码使用RxJava 框架处理一个数据流,并对不同任务类型使用不同线程池。对于IO操作相关的,指派使用一个线程池. 对于计算密集型的,使用主线程执行
import io.reactivex.rxjava3.core.Flowable;
import io.reactivex.rxjava3.schedulers.Schedulers;
public class RxJavaSample {
public static void main(String[] args) {
Flowable.range(1, 10)
.parallel(2)
.runOn(Schedulers.computation())
.map(v -> v * v) //计算密集型
.sequential()
.blockingSubscribe(System.out::println);
}
}代码依次输出1到10的平方说明如下
| API | 说明如下 |
|---|---|
| range | 产生1到10个Integer数据流,还可以通过fromCallable自定义数据流。调用返回一个Flowable<Integer> |
| parallel | 使用多少个并行处理上行数据 |
| runOn | 使用那种线程池执行,这里使用Schedulers.computation(),专门用于计算密集型任务.如果是IO密集型,调用Schedulers.io(), 或者自己定义线程池:Schedulers.from(executor) |
| map | 一个数据处理任务 |
| sequential | 按照顺序整理上一步骤处理结果 |
| blockingSubscribe | 结果回调,顺序输出1,4,9,16 ....100 后返回。即阻塞调用。如果是subscribe,则不阻塞调用 |
RxJava提供多种方式用于处理上游生产数据过快,超过下游消费能力的方式,最简单的调用Flowable.onBackpressureBuffer系列方法制定策略
Flowable.range(1, 10).onBackpressureBuffer(10) //缓冲10个,如果溢出,则中断执行并抛出异常
Flowable.range(1, 10).onBackpressureLatest() // 缓冲区永远最多只保留一条最新数据由于篇幅限制,更多的CompletableFuture的异步编排方式参考其API, 参考RxJava官网文档了解Flowable的用法
