Skip to content

性能优化-并发、异步和响应式

About 2352 wordsAbout 8 min

性能新书

2026-09-04

并发、异步和响应式

现代程序语言和提供的类库,使得基于多线程的编程得很容易,多线程编程充分利用计算机硬件提供的多个CPU增加系统吞吐量。多线程编程衍生出并发、异步,异步编排和响应式编程风格 。每个风格完成的的需求不一样

风格描述
并发编程充分利用CPU等硬件资源,启用多个线程,并发执行多个请求,不关心调用结果。可以创建Thread和线程池ThreadPoolExecutor
异步请求执行过程中,遇到阻塞调用,如数据库查询或者微服调用,可以把这种阻塞调封装在线程里执行。异步调用结果可通过回调或者Future俩种方式获取
异步编排多个异步调用按照一定顺序执行,比如请求A和请求B 都执行完毕后,根据俩个的结果再异步执行请求C, 使用CompletableFuture
响应式同架构的事件驱动风格,Producer连续生成一组有限或者无限的数据,经过Consumer的计算分析并汇总,Consumer可以允许并行或者串行处理数据。响应式支持Consumer一种可控速率消费Producer的数据(背压)。支持响应式的有JDK9的java.util.concurrent.Flow和开源工具RxJavaSpring 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的用法

知行合一