基础理论

什么是进程?

进程是程序的一次执行过程,是系统运行程序的基本单位,也是资源调度的基本单位,因此进程是动态的。系统运行一个程序即是一个进程从创建,运行到消亡的过程。在 Java 中,当我们启动 main 函数时其实就是启动了一个 JVM 的进程,而 main 函数所在的线程就是这个进程中的一个线程,也称主线程。
比如:打开的微信是一个进程,QQ也是一个进程。他们可以调用系统的IO设备,相机这些系统资源。

什么是线程?

线程与进程相似,但线程是一个比进程更小的执行单位。一个进程在其执行的过程中可以产生多个线程。与进程不同的是同类的多个线程共享进程的堆和方法区资源,但每个线程有自己的程序计数器、虚拟机栈和本地方法栈,所以系统在产生一个线程,或是在各个线程之间做切换工作时,负担要比进程小得多,也正因为如此,线程也被称为轻量级进程。

java线程和操作系统的线程有啥区别?

JDK 1.2 之前,Java 线程是基于绿色线程(Green Threads)实现的,这是一种用户级线程(用户线程),也就是说 JVM 自己模拟了多线程的运行,而不依赖于操作系统。由于绿色线程和原生线程比起来在使用时有一些限制(比如绿色线程不能直接使用操作系统提供的功能如异步 I/O、只能在一个内核线程上运行无法利用多核),在 JDK 1.2 及以后,Java 线程改为基于原生线程(Native Threads)实现,也就是说 JVM 直接使用操作系统原生的内核级线程(内核线程)来实现 Java 线程,由操作系统内核进行线程的调度和管理。
用户线程和内核线程:

  • 用户线程:由用户空间程序管理和调度的线程,运行在用户空间(专门给应用程序使用)
  • 内核线程:由操作系统内核管理和调度的线程,运行在内核空间(只有内核程序可以访问)

常见的线程模型有三种:

  1. 一对一(一个用户线程对应一个内核线程)
  2. 多对一(多个用户线程映射到一个内核线程)
  3. 多对多(多个用户线程映射到多个内核线程)
    在这里插入图片描述
    在 Windows 和 Linux 等主流操作系统中,Java 线程采用的是一对一的线程模型,也就是一个 Java 线程对应一个系统内核线程,除了Solaris之外,它是个特例.

虚拟线程

什么是虚拟线程?

JDK21后Java正式可以使用虚拟线程,虚拟线程是一种轻量级(用户模式)线程,这种线程是由Java虚拟机调度,而不是操作系统。虚拟线程占用空间小,任务切换开销几乎可以忽略不计,因此可以极大量地创建和使用。
在这里插入图片描述

为什么需要虚拟线程?

JDK中的每个java.lang.Thread实例也就是每个平台线程实例都在底层操作系统线程上运行Java代码,并且平台线程在运行代码的整个生命周期内捕获系统线程。平台线程与底层系统线程是一一对应的。

平台线程实例本质是由系统内核的线程调度程序进行调度,所以平台线程会有以下限制:

  • 资源有限导致系统线程总量有限,进而导致与系统线程一一对应的平台线程有限

  • 平台线程的调度依赖于系统的线程调度程序,当平台线程创建过多,会消耗大量资源用于处理线程上下文切换

  • 每个平台线程都会开辟一块私有的栈空间,大量平台线程会占据大量内存

这些限制导致开发者不能极大量地创建平台线程,为了满足性能需要,需要引入池化技术、添加任务队列构建消费者-生产者模式等方案去让平台线程适配多变的现实场景。显然,开发者们迫切需要一种轻量级线程实现,刚好可以弥补上面提到的平台线程的限制,这种轻量级线程可以满足:

  • 可以大量创建,例如十万级别、百万级别,而不会占据大量内存
  • 由JVM进行调度和状态切换,并且与系统线程"松绑"
  • 用法与原来平台线程差不多,或者说尽量兼容平台线程现存的API

创建虚拟线程的方式

// 1、通过 Thread.ofVirtual() 创建
Runnable fn = () -> {
  // your code here
};

Thread thread = Thread.ofVirtual(fn).start();

// 2、通过 Thread.startVirtualThread() 创建
Thread thread = Thread.startVirtualThread(() -> {
  // your code here
});

// 3、通过 Executors.newVirtualThreadPerTaskExecutor() 创建
var executorService = Executors.newVirtualThreadPerTaskExecutor();
executorService.submit(() -> {
  // your code here
});
class CustomThread implements Runnable {
  @Override
  public void run() {
    System.out.println("CustomThread run");
  }
}

//4、通过 ThreadFactory 创建
CustomThread customThread = new CustomThread();
// 获取线程工厂类
ThreadFactory factory = Thread.ofVirtual().factory();
// 创建虚拟线程
Thread thread = factory.newThread(customThread);
// 启动线程
thread.start();

java线程的状态有哪些?

在这里插入图片描述
NEW: 尚未启动的线程状态,即线程创建,还未调用start方法

RUNNABLE: 就绪状态(调用start,等待调度)+正在运行

BLOCKED: 等待监视器锁时,陷入阻塞状态

WAITING: 等待状态的线程正在等待另一线程执行特定的操作(如notify)

TIMED_WAITING: 具有指定等待时间的等待状态

TERMINATED: 线程完成执行,终止状态

线程的创建/销毁流程

线程的创建流程:

  1. Java: thread.start()
  2. JVM: 分配线程栈、初始化线程控制块
  3. OS: 创建内核线程结构体(task_struct)
  4. OS: 分配内核栈,设置上下文
  5. OS: 将线程加入就绪队列
  6. OS调度器: 选择时机分配CPU时间片

所以说线程被创建了不是立马执行的,线程创建好后进入就绪状态,之后由操作系统调度执行。

线程的销毁流程:

  1. Java: run()方法执行完毕或抛出异常
  2. JVM: 调用线程退出处理钩子
  3. JVM: 清理Java层面的资源
  4. OS: 回收内核资源,线程状态变为Zombie
  5. OS: 父线程(或init进程)回收剩余资源

线程Thread

线程的创建方式

1.继承Thread类

最简单的方式,继承Thread类,重写run方法

public class ExtendsThread extends Thread {
    @Override
    public void run() {
        System.out.println("1......");
    }
// 启动线程
    public static void main(String[] args) {
        new ExtendsThread().start();
    }
}

2.实现Runable接口

实现Runnable接口并重写run方法,这个方法拿不到返回值

public class ImplementsRunnable implements Runnable {
    @Override
    public void run() {
        System.out.println("2......");
    }
    public static void main(String[] args) {
        ImplementsRunnable runnable = new ImplementsRunnable();
        new Thread(runnable).start();
    }
}

3.实现Callable接口

实现Runnable接口并重写call方法,将其包装为FutureTask,这个方式可以拿到线程的返回值

public class ImplementsCallable implements Callable<String> {
    @Override
    public String call() throws Exception {
        System.out.println("3......");
        return "zhuZi";
    }
    public static void main(String[] args) throws Exception {
        ImplementsCallable callable = new ImplementsCallable();
        FutureTask<String> futureTask = new FutureTask<>(callable);
        new Thread(futureTask).start();
        System.out.println(futureTask.get()); // get会阻塞main线程
    }
}

4.使用ExecutorService线程池

可以通过Executors创建线程池,也可以自定义线程池

public class UseExecutorService {
    public static void main(String[] args) {
        ExecutorService poolA = Executors.newFixedThreadPool(2);
        poolA.execute(()->{
            System.out.println("4A......");
        });
        poolA.shutdown(); // 关闭线程池
        // 又或者自定义线程池
        ThreadPoolExecutor poolB = new ThreadPoolExecutor(2, 3, 0,
                TimeUnit.SECONDS, new LinkedBlockingQueue<Runnable>(3),
                Executors.defaultThreadFactory(), new ThreadPoolExecutor.AbortPolicy());
        poolB.submit(()->{
            System.out.println("4B......");
        });
        poolB.shutdown();
    }
}

5.使用CompletableFuture类
CompletableFutureJDK1.8引入的新类,可以用来执行异步任务

public class UseCompletableFuture {
    public static void main(String[] args) throws InterruptedException {
        CompletableFuture<String> cf = CompletableFuture.supplyAsync(() -> {
            System.out.println("5......");
            return "zhuZi";
        });
        // 需要阻塞,否则看不到结果
        Thread.sleep(1000);
    }
}

6.基于ThreadGroup线程组

Java线程可以分组,可以创建多条线程作为一个组

public class UseThreadGroup {
    public static void main(String[] args) {
        ThreadGroup group = new ThreadGroup("groupName");

        new Thread(group, ()->{
            System.out.println("6-T1......");
        }, "T1").start();

        new Thread(group, ()->{
            System.out.println("6-T2......");
        }, "T2").start();

        new Thread(group, ()->{
            System.out.println("6-T3......");
        }, "T3").start();
    }
}

7.使用FutureTask类
这个和之前实现Callable接口的方式差不多,只不过用匿名形式创建Callable

public class UseFutureTask {
    public static void main(String[] args) {
        FutureTask<String> futureTask = new FutureTask<>(() -> {
            System.out.println("7......");
            return "zhuZi";
        });
        new Thread(futureTask).start();
    }
}

8.使用匿名内部类或Lambda
这种方式就是直接内部常见的方式,直接new前面所说的Runnable接口,或者通过Lambda表达式书写

public class UseAnonymousClass {
    public static void main(String[] args) {
        new Thread(new Runnable() {
            @Override
            public void run() {
                System.out.println("8A......");
            }
        }).start();

        new Thread(() -> 
                System.out.println("8B......")
        ).start();
    }
}

9.使用Timer定时器类

public class UseTimer {
    public static void main(String[] args) {
        Timer timer = new Timer();

        timer.schedule(new TimerTask() {
            @Override
            public void run() {
                System.out.println("9......");
            }
        }, 0, 1000); // 里面需要传入两个数字,第一个代表启动后多久开始执行,第二个代表每间隔多久执行一次,单位是ms毫秒
    }
}

10 .使用ForkJoin或Stream并行流
ForkJoinJDK1.7引入的新线程池,基于分治思想实现。而后续JDK1.8parallelStream并行流,默认就基于ForkJoin实现

public class UseForkJoinPool {
    public static void main(String[] args) {
        ForkJoinPool forkJoinPool = new ForkJoinPool();
        forkJoinPool.execute(()->{
            System.out.println("10A......");
        });

        List<String> list = Arrays.asList("10B......");
        list.parallelStream().forEach(System.out::println);
    }
}

线程常用方法

线程生命周期控制方法

方法返回类型功能描述注意事项
start()void启动线程。JVM会调用该线程的 run() 方法,使其进入可运行状态 (RUNNABLE),并等待操作系统调度。一个线程对象只能调用一次 start()。多次调用会抛出 IllegalThreadStateException
run()void线程的执行体。当 start() 被调用后,JVM会调度此方法。如果直接调用 run(),它不会新建线程,而是作为当前线程的一个普通方法执行。
interrupt()void中断线程。这是一个协作机制。它不会强制停止线程,而是设置一个“中断状态”。线程需要在代码中定期检查 isInterrupted() 或响应 InterruptedException (如 sleep, join 抛出) 来决定如何处理中断。
isInterrupted()boolean检查线程是否被中断。调用后不会清除中断状态主要用于线程内部检查自身是否被请求中断。
static interrupted()boolean检查当前线程是否被中断。调用后会清除中断状态主要用于一个线程去检查另一个线程的中断状态并清除状态。注意 static 关键字。
join()void等待线程终止。调用此方法的线程会阻塞,直到 join() 调用的那个线程执行完毕。如果 join() 的线程被中断,会抛出 InterruptedException
join(long millis)void带超时的 join。等待指定时间(毫秒)后,无论目标线程是否完成,调用线程都会继续执行。超时后方法返回。
setDaemon(boolean on)void设置守护线程。必须在线程启动前调用。JVM不会等待守护线程执行完毕。当所有用户线程(非守护线程)结束时,JVM会退出,此时所有守护线程都会被强制终止。

线程状态与元信息获取方法

方法返回类型功能描述
getId()long获取线程的唯一标识符。
getName() / setName(String name)String / void获取或设置线程的名称。
getPriority() / setPriority(int priority)int / void获取或设置线程的优先级。优先级范围是 Thread.MIN_PRIORITY (1) 到 Thread.MAX_PRIORITY (10),默认为 NORM_PRIORITY (5)。高优先级线程被调度的概率更大,但不保证一定先执行。
getState()Thread.State获取线程的当前状态。返回值是一个枚举,包括 NEW, RUNNABLE, BLOCKED, WAITING, TIMED_WAITING, TERMINATED
isAlive()boolean检查线程是否存活。线程在 start() 之后,run() 方法执行完之前,都处于存活状态。
isDaemon()boolean检查是否是守护线程。
static currentThread()Thread静态方法,返回对当前正在执行线程的引用。

线程同步与通信方法

方法返回类型功能描述
wait()void使当前线程等待。调用此方法的线程会释放锁并进入 WAITING 状态,直到其他线程调用同一对象的 notify()notifyAll() 来唤醒它。
wait(long timeout)void带超时的 wait。线程会等待 timeout 毫秒,如果在超时前没有被唤醒,也会自动唤醒,进入 TIMED_WAITING 状态。
notify()void唤醒一个在此对象监视器上等待的单个线程。选择哪个线程是任意的(不保证公平性)。
notifyAll()void唤醒所有在此对象监视器上等待的线程。被唤醒的线程会竞争锁,只有一个能成功获得锁并继续执行。

静态工具方法

方法返回类型功能描述
yield()static void礼让线程。提示当前线程愿意让出对处理器的使用权,从而让同等或更高优先级的线程有机会执行。它是一个建议,不保证执行
sleep(long millis)static void使当前线程休眠指定的毫秒数。进入 TIMED_WAITING 状态。休眠期间不占用 CPU,但依然持有锁。如果休眠期间被中断,会抛出 InterruptedException
onSpinWait()static void(Java 9+) 提示当前线程处于一个“自旋等待”的循环中,如果可能,可以优化底层代码以提高性能。

线程池

前面我们说过,java的平台线程和操作系统内核线程是一对一的关系,所以创建线程的数量有限制,频繁的创建线程也消耗系统资源,Java中为了更好的管理复用创建的线程,就引入了线程池,线程池就是管理一系列线程的资源池,其提供了一种限制和管理线程资源的方式。每个线程池还维护一些基本统计信息,例如已完成任务的数量。
线程池可以带来一下优点:

  • 降低资源消耗。通过重复利用已创建的线程降低线程创建和销毁造成的消耗。
  • 提高响应速度。当任务到达时,任务可以不需要等到线程创建就能立即执行。
  • 提高线程的可管理性。线程是稀缺资源,如果无限制的创建,不仅会消耗系统资源,还会降低系统的稳定性,使用线程池可以进行统一的分配,调优和监控。

Executor框架

Executor 是线程池的基本框架,其包括了线程池的管理,还提供了线程工厂、队列以及拒绝策略等,它主要由三部分组成:

1.任务(Runnable /Callable)

执行任务需要实现的 Runnable 接口 或 Callable接口。Runnable 接口或 Callable 接口 实现类都可以被 ThreadPoolExecutor 或 ScheduledThreadPoolExecutor 执行

2.任务的执行(Executor)

执行机制的核心接口 Executor ,以及继承自 Executor 接口的 ExecutorService 接口。常用的ThreadPoolExecutorScheduledThreadPoolExecutor 这两个关键类实现了 ExecutorService 接口。
在这里插入图片描述

3.异步计算的结果(Future)
当我们把 Runnable接口 或 Callable 接口 的实现类提交给 ThreadPoolExecutorScheduledThreadPoolExecutor 执行。(调用 submit() 方法时会返回一个 FutureTask 对象。
在这里插入图片描述

执行流程:

  1. 主线程首先要创建实现 Runnable 或者 Callable 接口的任务对象。
  2. 把创建完成的实现 Runnable/Callable接口的对象执行任务: ExecutorService.execute(Runnable command))或者ExecutorService.submit(Callable <T> task/Runnable task))。
  3. 如果执行 ExecutorService.submit(…)ExecutorService 将返回一个实现Future接口的对象
  4. 最后,主线程可以执行 FutureTask.get()方法来等待任务执行完成。主线程也可以执行 FutureTask.cancel(boolean mayInterruptIfRunning)来取消此任务的执行

线程池的创建

通过ThreadPoolExecutor创建

这个是最常用也是最基本的线程池创建方式。

    public ThreadPoolExecutor(int corePoolSize,//线程池的核心线程数量
                              int maximumPoolSize,//线程池的最大线程数
                              long keepAliveTime,//当线程数大于核心线程数时,多余的空闲线程存活的最长时间
                              TimeUnit unit,//时间单位
                              BlockingQueue<Runnable> workQueue,//任务队列,用来储存等待执行任务的队列
                              ThreadFactory threadFactory,//线程工厂,用来创建线程,一般默认即可
                              RejectedExecutionHandler handler//拒绝策略,当提交的任务过多而不能及时处理时,我们可以定制策略来处理任务
                               )

核心参数解释:

  • corePoolSize: 核心线程数量,也就是线程池拥有的最少的线程数量,线程池会一直持有这些核心线程,不会销毁。
  • maxmumPoolSize: 最大线程数量,当核心线程数用完后,阻塞队列也塞满了任务,说明线程不够用了,这个时候就可以新开线程来提高线程池的吞吐量,但是不能无限创建,maxmumPoolSize就规定最多可以创建多少
  • keepAliveTime:当线程池持有的线程数大于核心线程数,并且部分线程空闲了没活儿干了,那么这些线程就需要被回收,但是不会立马回收,万一突然又一大推任务来了呢,所以这个参数就规定这些多余的空闲线程存活的实现。
  • unit:keepAliveTime的时间单位
  • workQueue:工作队列,任务来了,但是核心线程数用完了,那么就把任务存到工作队列中,等待被执行。
  • threadFactory:线程工厂,可以自定义实现线程怎么创建。
  • RejectedExecutionHandler:拒绝策略,当核心线程用完了,阻塞队列也塞满了,线程也开到最大限制数,此时线程池达到最大负荷处理不了新添加的任务了,RejectedExecutionHandler就规定这个情况下如何处理新来的任务。

常用阻塞队列

  • ArrayBlockingQueue: 固定有界,实例化时必须指定容量,数组实现,遵循FIFO顺序
  • LinkedBlockingQueue: 可选有界(默认Integer.MAX_VALUE),也可以指定容量,单链表实现,遵循FIFO顺序出队
  • LinkedBlockingDeque:可选有界(默认Integer.MAX_VALUE),也可以指定容量,双向链表实现,可以队头队尾添加元素,灵活性高,遵循FIFO/LIFO顺序
  • PriorityBlockingQueue: 无界队列,底层二叉堆实现(数组二叉堆),自动扩容,按优先级出堆。元素必须实现compare接口
  • DelayQueue: 无界队列,优先级堆实现和PriorityBlockingQueue差不多,自动扩容,按照延迟时间出堆,元素必须实现Delayed接口
  • SynchronousQueue: 没有容量,不存储元素,添加一个元素就必须获取一个元素后才能添加,入队和出队必须成对出现。

常用拒绝策略

  • ThreadPoolExecutor.AbortPolicy:默认策略,抛出 RejectedExecutionException来拒绝新任务的处理,新任务来了不处理,直接抛异常
  • ThreadPoolExecutor.DiscardPolicy:不处理新任务,直接丢弃掉,不报异常
  • ThreadPoolExecutor.DiscardOldestPolicy:此策略将丢弃最早的未处理的任务请求,也就是抛弃工作队列队头任务,把任务新来塞到任务队列中
  • ThreadPoolExecutor.CallerRunsPolicy:调用执行自己的线程运行任务,也就是直接在调用execute方法的线程中运行(run)被拒绝的任务,如果执行程序已关闭,则会丢弃该任务,谁给我的任务谁执行。

线程工厂

public class CustomThreadFactory implements ThreadFactory { // 自定义线程工厂

    private final AtomicInteger threadNum = new AtomicInteger();
    @Override
    public Thread newThread(Runnable r) {
        Thread thread = new Thread(r);
        thread.setName("CustomThreadFactory-"+threadNum.getAndIncrement());
        thread.setDaemon(false);
        thread.setPriority(Thread.NORM_PRIORITY);
        return thread;
    }
}


ExecutorService threadPool = new ThreadPoolExecutor(corePoolSize, maximumPoolSize, keepAliveTime, TimeUnit.MINUTES, workQueue, CustomThreadFactory); // 创建线程时默认就会使用自定义工厂的newThread方法创建线程

线程池中的线程异常后,是主动销毁还是放回线程池复用

  • 使用execute()提交任务:当任务通过execute()提交到线程池并在执行过程中抛出异常时,如果这个异常没有在任务内被捕获,那么该异常会导致当前线程终止,并且异常会被打印到控制台或日志文件中。
    线程池会检测到这种线程终止,并创建一个新线程来替换它,从而保持配置的线程数不变。
  • 使用submit()提交任务:对于通过submit()提交的任务,如果在任务执行中发生异常,这个异常不会直接打印出来。相反,异常会被封装在由submit()返回的Future对象中。当调用Future.get()方法时,
    可以捕获到一个ExecutionException。在这种情况下,线程不会因为异常而终止,它会继续存在于线程池中,准备执行后续的任务。

如何设定线程池大小

  • CPU 密集型任务(N+1): 这种任务消耗的主要是 CPU 资源,可以将线程数设置为 N(CPU 核心数)+1。比 CPU 核心数多出来的一个线程是为了防止线程偶发的缺页中断,或者其它原因导致的任务暂停而带来的影响。
    一旦任务暂停,CPU 就会处于空闲状态,而在这种情况下多出来的一个线程就可以充分利用 CPU 的空闲时间。
  • I/O 密集型任务(2N): 这种任务应用起来,系统会用大部分的时间来处理 I/O 交互,而线程在处理 I/O 的时间段内不会占用 CPU 来处理,这时就可以将 CPU 交出给其它线程使用。因此在 I/O 密集型任务的应用中,
    我们可以多配置一些线程,具体的计算方法是 2N。

通过Executors工具类创建

1.固定线程数线程池
public static ExecutorService newFixedThreadPool(int nThreads) {
    // 最大线程数 == 核心线程数
    return new ThreadPoolExecutor(nThreads, nThreads,0L, TimeUnit.MILLISECONDS,new LinkedBlockingQueue<Runnable>());
   // !!! LinkedBlockingQueue使用的阻塞队列没有指定容量,默认最大值Intger.MAX_VALUE,当一直提交任务,就容易造成大量任务堆积导致OOM
}

2.单线程线程池
public static ExecutorService newSingleThreadExecutor() {
    // 线程池中只有一个线程可用
    return new FinalizableDelegatedExecutorService (new ThreadPoolExecutor(1, 1,0L, TimeUnit.MILLISECONDS,new LinkedBlockingQueue<Runnable>()));
    // !!! LinkedBlockingQueue使用的阻塞队列没有指定容量,默认最大值Intger.MAX_VALUE,当一直提交任务,就容易造成大量任务堆积导致OOM
}

3. 缓存线程池
public static ExecutorService newCachedThreadPool() {
    // 核心线程数为0,最大线程数Integer.MAX_VALUE
    return new ThreadPoolExecutor(0, Integer.MAX_VALUE,60L, TimeUnit.SECONDS,new SynchronousQueue<Runnable>());
   // !!! 同步队列 SynchronousQueue,没有容量,最大线程数是 Integer.MAX_VALUE,无线的创建线程也会导致OOM
}

4.任务调度线程池
public static ScheduledExecutorService newScheduledThreadPool(int corePoolSize) {
    return new ScheduledThreadPoolExecutor(corePoolSize);
}
public ScheduledThreadPoolExecutor(int corePoolSize) {
    // 最大线程数Integer.MAX_VALUE 
    super(corePoolSize, Integer.MAX_VALUE, 0, NANOSECONDS,new DelayedWorkQueue());
   // !!! 最大线程数Integer.MAX_VALUE并且使用的是DelayedWorkQueue无界阻塞队列容易导致OOM
}

线程池的使用

常用方法:

功能分类方法名返回类型功能描述与使用场景
1. 构造与初始化ThreadPoolExecutor(...)ThreadPoolExecutor核心。创建并初始化一个线程池实例。通过设置七个核心参数(核心/最大线程数、存活时间、工作队列、线程工厂、拒绝策略)来定制线程池的行为,这是使用线程池的第一步。
2. 任务提交execute(Runnable command)void最常用。提交一个 Runnable 任务到线程池。如果线程池和队列都满且达到最大线程数,会根据设定的拒绝策略处理任务。由 Executor 接口定义。
submit(Runnable task)Future<?>提交一个 Runnable 任务,并返回一个 Future 对象。主要用于跟踪任务状态(如是否完成)或等待任务结束,尽管 Runnable 本身没有返回值。
submit(Callable<T> task)Future<T>用于获取返回值。提交一个 Callable 任务,并返回一个 Future<T> 对象。Future.get() 方法会阻塞直到任务完成,并返回 Callable 计算的结果,或抛出执行中的异常。
submit(Runnable task, T result)Future<T>提交一个 Runnable 任务,并预定义一个 result 对象。Future.get() 将直接返回这个 result 对象,提供了一种让 Runnable 任务“返回”特定值的机制。
3. 线程池关闭与终止shutdown()void优雅关闭。停止接受新任务,但会完成队列中已提交的所有任务。已提交且正在执行的任务不受影响。线程池状态变为 SHUTDOWN
shutdownNow()List<Runnable>立即关闭。尝试停止所有正在执行的任务(通过中断),并返回一个包含等待在队列中尚未执行的任务列表。线程池状态变为 STOP
isShutdown()boolean查询线程池是否已经关闭(即已调用 shutdown()shutdownNow())。一旦关闭,不能再提交新任务。
isTerminated()boolean查询线程池是否已经完全终止。即所有任务(包括正在队列中等待的)都已执行完毕,并且所有工作线程都已销毁。
awaitTermination(long timeout, TimeUnit unit)boolean阻塞等待关闭完成。在调用 shutdown() 后,调用此方法可以阻塞当前线程,直到线程池完全终止,或者指定的超时时间到达。返回 true 表示成功终止,false 表示超时。
4. 运时管理与动态调整setCorePoolSize(int newSize)void动态调整核心线程数。可以随时增加或减少核心线程数。减少核心线程数会导致多余的空闲线程在下次检查时被回收。
setMaximumPoolSize(int newSize)void动态调整最大线程数。可以随时增加或减少最大线程数。减少最大线程数且当前线程数超过新值时,多余的空闲线程会立即被回收。
setKeepAliveTime(long time, TimeUnit unit)void动态调整线程空闲存活时间。此参数只对超过核心线程数的那部分线程有效。
allowCoreThreadTimeOut(boolean value)void设置核心线程是否允许超时。默认为 false,即核心线程会永久存活。设置为 true 后,核心线程在超过 keepAliveTime 也会被回收,除非线程池中有任务在执行。
prestartCoreThread()boolean预启动一个核心线程。在任务到来前,手动启动一个核心线程,使其处于就绪状态,减少第一个任务的等待时间。
prestartAllCoreThreads()int预启动所有核心线程。一次性启动所有核心线程。
5. 监控与诊断getActiveCount()int获取当前正在执行任务的线程数。这个值是动态变化的,是衡量线程池负载的关键指标
getPoolSize()int获取当前池中存在的线程总数,包括正在执行任务和空闲的线程。
getCorePoolSize()int获取当前配置的核心线程数
getMaximumPoolSize()int获取当前配置的最大线程数
getCompletedTaskCount()long获取自线程池启动以来,已成功完成任务的总数
getTaskCount()long获取自线程池启动以来,已提交的总任务数。这个值是已完成、正在执行和队列中等待任务的总和。
getQueue()BlockingQueue<Runnable>获取工作队列的引用。通过查看队列的 size(),可以监控任务积压情况。如果队列持续增长,说明处理速度跟不上提交速度。
getLargestPoolSize()int获取线程池曾达到过的最大线程数。这对于性能调优非常有用,可以帮助你判断是否需要调大 maximumPoolSize
getThreadFactory()ThreadFactory获取创建工作线程的线程工厂。
getRejectedExecutionHandler()RejectedExecutionHandler获取当前设置的拒绝策略。

使用示例:

import java.util.concurrent.*;

public class ThreadPoolExecutorExample {

    public static void main(String[] args) {
        // 1. 创建线程池
        // 核心线程2,最大线程4,空闲存活时间60秒,使用有界队列(容量2),拒绝策略为调用者运行
        ThreadPoolExecutor executor = new ThreadPoolExecutor(
                2,
                4,
                60L,
                TimeUnit.SECONDS,
                new ArrayBlockingQueue<>(2),
                Executors.defaultThreadFactory(),
                new ThreadPoolExecutor.CallerRunsPolicy()
        );

        // 监控线程池状态
        Runtime runtime = Runtime.getRuntime();
        ScheduledExecutorService monitorExecutor = Executors.newSingleThreadScheduledExecutor();
        monitorExecutor.scheduleAtFixedRate(() -> {
            System.out.println("=====================================");
            System.out.println("活跃线程数: " + executor.getActiveCount());
            System.out.println("核心线程数: " + executor.getCorePoolSize());
            System.out.println("总线程数: " + executor.getPoolSize());
            System.out.println("已完成任务数: " + executor.getCompletedTaskCount());
            System.out.println("总任务数: " + executor.getTaskCount());
            
            // 监控队列积压
            System.out.println("队列中等待的任务数: " + executor.getQueue().size());
            System.out.println("JVM 可用内存: " + (runtime.maxMemory() - runtime.totalMemory() + runtime.freeMemory()) / 1024 / 1024 + "M");
            System.out.println("=====================================");
        }, 0, 2, TimeUnit.SECONDS);

        // 2. 提交任务
        for (int i = 0; i < 10; i++) {
            final int taskId = i;
            try {
                executor.execute(() -> {
                    try {
                        System.out.println("任务 " + taskId + " 开始执行,由线程 " + Thread.currentThread().getName() + " 处理");
                        // 模拟任务执行耗时
                        TimeUnit.SECONDS.sleep(3);
                        System.out.println("任务 " + taskId + " 执行完毕");
                    } catch (InterruptedException e) {
                        System.out.println("任务 " + taskId + " 被中断");
                        Thread.currentThread().interrupt();
                    }
                });
                System.out.println("成功提交任务 " + taskId);
            } catch (RejectedExecutionException e) {
                System.out.println("任务 " + taskId + " 提交失败,队列已满,触发拒绝策略");
            }
        }

        // 3. 关闭线程池
        try {
            // 先停止接受新任务,但会处理完队列中的任务
            System.out.println("\n准备执行 shutdown()...");
            executor.shutdown();
            
            // 等待线程池真正终止,最多等待1分钟
            if (!executor.awaitTermination(60, TimeUnit.SECONDS)) {
                System.out.println("线程池在1分钟内未终止,执行 shutdownNow()...");
                executor.shutdownNow();
                // 再次等待(可选)
                if (!executor.awaitTermination(60, TimeUnit.SECONDS)) {
                    System.err.println("线程池未能终止");
                }
            }
        } catch (InterruptedException e) {
            executor.shutdownNow();
            Thread.currentThread().interrupt();
        }

        // 关闭监控线程
        monitorExecutor.shutdown();
        System.out.println("主线程结束");
    }
}

线程池的处理流程

  1. 新来一个任务时,如果执行线程数小于核心线程数,就新建一个线程执行当前任务
  2. 如果执行线程数大于核心线程数,加入到阻塞队列中等待执行
  3. 如果阻塞队列已满,但是执行线程数小于最大线程数,就新建一个线程执行当前任务
  4. 如果阻塞队列满了,最大线程数也到达了,那么就按照指定的拒绝策略执行

线程池在spring框架中的使用

SpringBoot默认使用ThreadPoolTaskExecutor来处理异步任务,如果不进行特殊配置,SpringBoot会使用一个默认的SimpleAsyncTaskExecutor。这种默认实现适用于简单的场景,但对于复杂的应用程序,通常需要自定义线程池配置ThreadPoolTaskExecutor和ThreadPoolExcutor不是用一个线程,springboot为了更好的管理线程池,又对ThreadPoolExcutor进行二次封装,本质上底层还是使用的ThreadPoolExcutor。
看下数据结构就明白了:

public class ThreadPoolTaskExecutor extends ExecutorConfigurationSupport
   implements AsyncListenableTaskExecutor, SchedulingTaskExecutor {
   
// 池大小监控锁对象
private final Object poolSizeMonitor = new Object();

// 核心线程数 (默认: 1)
private int corePoolSize = 1;

// 最大线程数 (默认: Integer.MAX_VALUE)
private int maxPoolSize = Integer.MAX_VALUE;

// 线程空闲时间 (默认: 60秒)
private int keepAliveSeconds = 60;

// 队列容量 (默认: Integer.MAX_VALUE,实际上是无界队列)
private int queueCapacity = Integer.MAX_VALUE;

// 是否允许核心线程超时 (默认: false)
private boolean allowCoreThreadTimeOut = false;

// 任务装饰器,用于在执行前后添加额外逻辑(如上下文传递)
private TaskDecorator taskDecorator;

// 底层的 JDK ThreadPoolExecutor 实例
private ThreadPoolExecutor threadPoolExecutor;

}

参数解释:

属性名对应 ThreadPoolExecutor 参数含义与作用
corePoolSizecorePoolSize核心线程数。即使线程池处于空闲状态,也会一直保持存活的线程数量。当有新任务提交时,如果当前线程数小于 corePoolSize,会立即创建新线程来执行任务,即使其他核心线程处于空闲状态。
maxPoolSizemaximumPoolSize最大线程数。线程池允许创建的最大线程数量。只有在 workQueue (工作队列) 已满,且当前线程数小于 maxPoolSize 时,线程池才会创建新的线程。
queueCapacityworkQueue工作队列容量。一个用于存放等待执行任务的 BlockingQueue。当提交的任务数超过 corePoolSize 时,新任务会被放入这个队列中等待。如果队列也满了,才考虑创建新线程(直到达到 maxPoolSize)。
keepAliveSecondskeepAliveTime线程空闲存活时间。当线程池中的线程数量超过 corePoolSize 时,多余的线程如果在 keepAliveSeconds 时间内没有执行任何任务,就会被销毁。这样可以在任务量减少时,回收资源,降低资源消耗。
threadNamePrefix(Spring 特有)线程名前缀。Spring 提供的便利功能,所有由这个线程池创建的线程,其名称都会以此前缀开头,如 task-executor-1。这对于在日志、JMX 或性能分析中追踪线程来源非常有帮助。
allowCoreThreadTimeOutallowCoreThreadTimeOut是否允许核心线程超时。如果设置为 true,则核心线程在 keepAliveSeconds 时间内空闲时也会被销毁。默认为 false,核心线程会一直存活。
rejectedExecutionHandlerhandler拒绝策略。当线程和队列都已满,并且尝试提交新任务时,会触发拒绝策略。Spring 提供了内置的几种策略:
AbortPolicy (默认): 抛出 RejectedExecutionException
CallerRunsPolicy: 由提交任务的线程(调用者)自己执行该任务。
DiscardOldestPolicy: 丢弃队列中等待最久的任务,然后尝试再次提交当前任务。
DiscardPolicy: 直接丢弃当前任务,不抛出异常。
你也可以自定义策略。

自定义线程池:

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import java.util.concurrent.Executor;
@Configuration
public class AsyncConfig {
    @Bean(name = "taskExecutor")
    public Executor taskExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(5);
        executor.setMaxPoolSize(10);
        executor.setQueueCapacity(25);
        executor.setThreadNamePrefix("MyExecutor-");
        executor.initialize();
        return executor;
    }
}

使用@Async注解

SpringBoot提供了@Async注解,使得方法可以异步执行。使用@Async 时,Spring会自动使用配置的TaskExecutor。

首先,需要在应用程序主类或者配置类上启用异步支持:

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.scheduling.annotation.EnableAsync;

@SpringBootApplication 
@EnableAsync // 开启异步支持
public class Application {
    public static void main(String[] args) {
        SpringApplication.run(Application.class, args);
    }
}

在代码中使用@Async注解并且指定使用某一个线程池

import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Service;
@Service
public class MyService {
    @Async("taskExecutor") // 注解使用
    public void asyncMethod() {
        System.out.println("Execute method asynchronously - " + Thread.currentThread().getName());
    }
}


public class MyService {
    private final Executor taskExecutor;
    @Autowired  // 注入使用
    public MyService(Executor taskExecutor) {
        this.taskExecutor = taskExecutor;
    }
    public void executeTask() {
        taskExecutor.execute(() -> {
            System.out.println("Execute task in thread - " + Thread.currentThread().getName());
        });
    }
}

线程池在以下场景中非常有用

  • 处理异步任务:如文件上传、邮件发送等。

  • 并行处理:如批量数据处理、并行计算等。

  • 提高系统吞吐量:如高并发请求处理等。

线程池原理

核心数据结构

private final AtomicInteger ctl = new AtomicInteger(ctlOf(RUNNING, 0)); // 状态和线程数:用来标记线程池状态(高3位),线程个数(低29位,所以线程个数最多是2^29-1)
private static final int COUNT_BITS = Integer.SIZE - 3;  // 线程个数掩码位数:高3位保留用于表示线程池的状态。剩下的29位则用来表示线程个数。Integer.SIZE = 32
private static final int CAPACITY   = (1 << COUNT_BITS) - 1;  // 线程最大个数

// 线程池状态(高3位)
private static final int RUNNING    = -1 << COUNT_BITS;  // 11100000...000 线程池处于正常状态,可以接受新的任务,同时会按照预设的策略来处理已有任务的执行。
private static final int SHUTDOWN   =  0 << COUNT_BITS;  // 00000000...000 线程池处于关闭状态,不再接受新的任务,但是会继续执行已有任务直到执行完成。执行线程池对象的shutdown()时进入该状态。
private static final int STOP       =  1 << COUNT_BITS;  // 00100000...000 线程池处于关闭状态,不再接受新的任务,同时会中断正在执行的任务,清空线程队列。执行shutdownNow()时进入该状态。
private static final int TIDYING    =  2 << COUNT_BITS;  // 01000000...000 所有任务已经执行完毕,线程池进入该状态会开始进行一些结尾工作,比如及时清理线程池的一些资源。
private static final int TERMINATED =  3 << COUNT_BITS;  // 01100000...000 线程池已经完全停止,所有的状态都已经结束了,线程池处于最终的状态。

// 包装和解析方法
private static int runStateOf(int c)     { return c & ~CAPACITY; }     // 获取状态
private static int workerCountOf(int c)  { return c & CAPACITY; }      // 获取线程数
private static int ctlOf(int rs, int wc) { return rs | wc; }           // 组合状态和线程数
private final ReentrantLock mainLock = new ReentrantLock(); // 可重入锁,用于加锁往workers里加线程
private final Condition termination = mainLock.newCondition(); //  用于线程间通信,通过await()和signal(),signalAll()让线程等待或唤醒
private final HashSet<Worker> workers = new HashSet<Worker>(); // 一个HashSet,存有所有工作线程
private volatile boolean allowCoreThreadTimeOut; // 是否允许核心线程超时:默认是false,即核心线程永不超时,如果是true,核心线程将根据参数的超时时间存活
private final BlockingQueue<Runnable> workQueue; // 阻塞队列

Worker数据结构:

private final class Worker extends AbstractQueuedSynchronizer implements Runnable {
    
    final Thread thread;      // 实际执行任务的线程
    Runnable firstTask;       // 初始任务(可能为null)
    volatile long completedTasks; // 完成的任务计数
    
    Worker(Runnable firstTask) {
        setState(-1); // 初始状态,禁止中断直到runWorker
        this.firstTask = firstTask;
        this.thread = getThreadFactory().newThread(this); // 创建线程,以自身为Runnable
    }
    
    public void run() {
        runWorker(this);  // 委托给外部类的runWorker方法
    }
}

工作流程

在这里插入图片描述

线程池工作原理:

  1. 判断核心线程数: 新加入任务,判断corePoolSize是否到最大值;如果没到最大值就创建核心线程执行新任务,如果到最大值就判断是否有空闲的核心线程;
  2. 判断加入队列: 如果有空闲的核心线程,则空闲核心线程执行新任务,如果没空闲的核心线程,则尝试加入FIFO阻塞队列;
  3. 判断队列容量: 若加入成功,则等待空闲核心线程将队头任务取出并执行,若加入失败(例如队列满了),则判断是否到最大值;
  4. 判断丢弃策略: 如果没到最大值就创建非核心线程执行新任务,如果到了最大值就执行丢弃策略,默认丢弃新任务;
  5. 非核心线程自动回收: 线程数大于corePoolSize时,空闲线程将在keepAliveTime后回收,直到线程数等于核心线程数。这些核心线程也不会被回收。

源码解析

执行任务execute方法:

public void execute(Runnable command) {
    if (command == null)
        throw new NullPointerException();
    int c = ctl.get();
    if (workerCountOf(c) < corePoolSize) { // 执行线程数量小于核心线程数,新增核心线程执行任务
        if (addWorker(command, true))
            return;
        c = ctl.get();
    }
    if (isRunning(c) && workQueue.offer(command)) { // 核心线程数满了,尝试将任务添加到工作队列中等待执行
        int recheck = ctl.get();
        if (! isRunning(recheck) && remove(command)) // 再次校验线程池状态,如果不是running状态了就要将刚刚添加的任务从工作队列中移除
            reject(command); // 移除成功后,执行拒绝策略
        else if (workerCountOf(recheck) == 0)  // 到这步说明添加任务到工作队列成功,并且线程池状态支正常可以执行任务,如果线程池数量为空,就新增一个空闲线程执行刚刚新增的任务
            addWorker(null, false);
    }
    else if (!addWorker(command, false)) // 添加任务到阻塞队列失败说明工作队列满了,就新增非核心线程执行任务
        reject(command); // 如果新增非核心线程执行任务失败,则执行拒绝策略
}

添加工作线程addWorker方法:

private boolean addWorker(Runnable firstTask, boolean core) {
    retry: // 这个标签与 for (;;) 循环结合使用,调用continue retry时重新开始整个循环,调用break retry时结束整个循环
    for (int c = ctl.get();;) {
        // 1.1 如果当前线程池状态是SHUTDOWN或者STOP,提交的任务又不空,或者工作队列为空,则添加worker失败
        if (runStateAtLeast(c, SHUTDOWN) && (runStateAtLeast(c, STOP) || firstTask != null || workQueue.isEmpty()))
            return false;

        for (;;) {
           // 1.2 如果执行线程数到达上限,添加woker也失败
            if (workerCountOf(c) >= ((core ? corePoolSize : maximumPoolSize) & COUNT_MASK))
                return false;
           // 1.3 ctl线程计数+1,cas成功就跳出循环执行后面的添加worker逻辑
            if (compareAndIncrementWorkerCount(c))
                break retry;
            c = ctl.get();
          // 1.4 计数+1失败,再次校验线程池状态,如果是SHUTDOWN状态重新尝试增加ctl
            if (runStateAtLeast(c, SHUTDOWN))
                continue retry;
        }
    }
    
    // 线程池线程数量ctl+ 1成功,开始新增worker
    boolean workerStarted = false;
    boolean workerAdded = false;
    Worker w = null;
    try {
        w = new Worker(firstTask); // 2. 新增一个worker实例,创建线程
        final Thread t = w.thread;
        if (t != null) {
            final ReentrantLock mainLock = this.mainLock;
            mainLock.lock();  // 全局锁,加锁添加worker防止并发冲突
            try {
                int c = ctl.get();
                 // 再次校验线程池状态
                if (isRunning(c) || (runStateLessThan(c, STOP) && firstTask == null)) {
                    if (t.getState() != Thread.State.NEW)
                        throw new IllegalThreadStateException();
                    workers.add(w); //3. workers是一个HashSet,用于保存所有工作线程,这里将新增的worker加入到集合中
                    workerAdded = true;
                    int s = workers.size();
                    if (s > largestPoolSize) // largestPoolSize记录着线程池中出现过的最大线程数量
                        largestPoolSize = s;
                }
            } finally {
                mainLock.unlock();
            }
            if (workerAdded) { 4. worker新增成功,启动线程执行任务,这里其实最终调用了woker的runWorker方法
                container.start(t);
                workerStarted = true;
            }
        }
    } finally {
        if (! workerStarted) // 5. 新增失败回退刚刚的操作,将worker从hashset中移除,并且ctl - 1 
            addWorkerFailed(w);
    }
    return workerStarted;
}

runWorker ()执行线程:

上面addWorker()开启线程时,调用的Worker对象的start()方法。Worker对象封装了原线程。Worker类是Runnable的子类,它的run()方法只有一行,调用了runWorker()方法

final void runWorker(Worker w) {
    Thread wt = Thread.currentThread();
    Runnable task = w.firstTask;
    w.firstTask = null;
    w.unlock(); // 1. 解锁操作,防止后面加锁出现死锁
    boolean completedAbruptly = true;
    try {
        // 2.开始执行task,如果当前task为空,则从工作队列中获取新的task
        while (task != null || (task = getTask()) != null) {
            w.lock(); // 加锁保证代表当前worker正在执行工作
            // 3.如果线程池此时异常了,或者当前线程被设置了请求中断,就设置线程为可中端状态
            if ((runStateAtLeast(ctl.get(), STOP) || (Thread.interrupted() && runStateAtLeast(ctl.get(), STOP))) && !wt.isInterrupted())
                wt.interrupt();
            try {
                beforeExecute(wt, task); // 4. 在任务执行前做一些操作:这个方法内为空,需要我们继承ThreadPoolExecutor重写这个方法
                try {
                    task.run(); // 5.启动线程执行任务
                    afterExecute(task, null); // 6. 在任务执行前做一些操作:这个方法内为空,需要我们继承ThreadPoolExecutor重写这个方法
                } catch (Throwable ex) {
                    afterExecute(task, ex);
                    throw ex;
                }
            } finally {
                task = null; // 任务置为null便于GC
                w.completedTasks++; // 执行任务数+1
                w.unlock(); // 解锁
            }
        }
        completedAbruptly = false;
    } finally {
       //7. 退出循环说明没有task执行了,或者发生异常,执行线程退出逻辑
        processWorkerExit(w, completedAbruptly);
    }
}

submit提交任务:

public Future<?> submit(Runnable task) {
    if (task == null) throw new NullPointerException();
    RunnableFuture<Void> ftask = newTaskFor(task, null); // 封装为future
    execute(ftask); // 最终还是调用execute方法
    return ftask;
}

shutdown():关闭线程

/**
     * 关闭线程池
     */
    public void shutdown() {
        // 1.加锁
        final ReentrantLock mainLock = this.mainLock;
        mainLock.lock();
        try {
            // 2.校验关闭权限:校验调用者是否有权限关闭线程池及其中的所有线程
            checkShutdownAccess();
            // 3.修改线程池状态:将线程池状态修改为SHUTDOWN
            advanceRunState(SHUTDOWN);
            // 4.中断空闲线程:遍历线程池workers中的工作线程,并尝试中断那些处于空闲状态的线程。
            interruptIdleWorkers();
            // 5.钩子函数,实际内容为空,用户继承重写可以编写关闭线程池后要执行的逻辑
            onShutdown(); // hook for ScheduledThreadPoolExecutor
        } finally {
            mainLock.unlock();
        }
        // 6.查看当前线程池是否可以变为TERMINATED状态
        // 从 SHUTDOWN 状态修改为 TIDYING,在修改为 TERMINATED
        tryTerminate();
    }

ForkJoinPool线程池

基础介绍

什么是分治任务模型

简单来说就是对于一个大任务,一下子解决很困难,可以把这个大任务拆分为小任务,逐个击破,每个小任务解决了,大任务自然就解决了。

分治任务模型可分为两个阶段:

  1. 第一个阶段是任务分解,就是迭代地将任务分解为子任务,直到子任务可以直接计算出结果;

  2. 第二个阶段是 结果合并,即逐层合并子任务的执行结果,直到获得最终结果。

在这里插入图片描述

假设一个场景:

需要对100万的数据进行排序,一下子整体排序肯定很困难,采用分治思想,先将数据分为100份,每一个份数据很少自然就比较好排序,排好后再分为50份,因为之前100份的每一份都是有序的,所以分为的50份中的每一份就很好排序,以此类推直到整体排好序,这其实就是归并排序算发的思想,先局部有序,再整体有序。

放在实际代码中怎么实现呢?这么大的数据如果用一个线程去做就会很慢,自然想到ThreadPoolExcutor线程池,把数据分为100份扔到线程池的任务队列中,线程池开个20个线程对每一份数据进行局部排序,100份各自排序好后,在将100份分成50份,塞到任务队列中交给线程池排序,以此类推直到整体全部排好序,想法是可以的,但是使用现有的线程池实现这种场景比较麻烦。

所以JDK1.7后引入了一种新的 Fork/Join 并行计算框架,主要就用于支持这这种分治任务模型,不是用来替代ThreadPoolExcutor,而是ForkJoinPool底层实现就是为了更好的适用于分治任务模型,像java8引入的stream并行流和ComplateFuture默认使用的就是ForkJoinPool。

基本使用

如何创建

1.使用构造函数

public ForkJoinPool(int parallelism,
                    ForkJoinWorkerThreadFactory factory,
                    UncaughtExceptionHandler handler,
                    boolean asyncMode)
ForkJoinPool customPool = new ForkJoinPool(3)
ForkJoinPool customPool = new ForkJoinPool(3,cutomfactory,exceptionhander,false);

参数解释:

  • parallelism (并行度): 工作线程的数量。通常是设置为 Runtime.getRuntime().availableProcessors() (CPU 核心数) 以获得最佳性能。

  • ForkJoinWorkerThreadFactory factory (工作线程工厂): 用于创建工作线程的工厂。默认工厂就足够,但自定义工厂可以让你设置线程名、优先级、是否为守护线程等。

  • UncaughtExceptionHandler handler (异常处理器): 当工作线程抛出未捕获的异常时调用的处理器。默认为 null(打印堆栈并终止线程)。在生产环境中,自定义一个统一的异常处理器(如记录日志)是好习惯。

  • asyncMode (异步模式):

    • false (默认): 使用 LIFO (后进先出) 的工作窃取策略,非常适合递归分治任务。
    • true: 使用 FIFO (先进先出) 的模式,适合处理独立的、无依赖的任务流。

默认无参构造函数:

public ForkJoinPool() {
    this(Math.min(MAX_CAP, Runtime.getRuntime().availableProcessors()),
         defaultForkJoinWorkerThreadFactory, null, false);
}

2.使用Executors工具类来创建ForkJoinPool

// 创建一个 ForkJoinPool 线程池
ExecutorService forkJoinPool = Executors.newWorkStealingPool();// 默认CPU 核心数
ExecutorService forkJoinPool = Executors.newWorkStealingPool(3); // 指定线程数

3.**ForkJoinPool.commonPool()**创建

ForkJoinPool forkJoinPool = ForkJoinPool.commonPool();
// - 简单的并行流 (parallelStream) - 不需要精细控制的 CompletableFuture - 小型或示例项目

实例常用方法

ForkJoinPool实例方法返回类型功能描述与使用场景
核心任务执行
invoke(ForkJoinTask task)<T> T阻塞执行并等待任务完成。提交的线程会一直等待,直到该任务及其所有子任务执行完毕,然后返回最终结果。适用于需要同步等待结果的场景。
submit(ForkJoinTask task)ForkJoinTask<T>异步提交任务。立即返回一个 Future 对象,调用线程不会被阻塞。后续可通过返回的 ForkJoinTask 对象(Future)来检查状态或获取结果。这是最常用的提交方式之一。
execute(ForkJoinTask task)void异步执行任务。不关心任务结果,也不提供任何用于跟踪任务的对象。调用此方法后任务即被放入队列,调用线程继续执行。
execute(Runnable task)void异步执行一个 Runnable 任务。与上面的 execute 不同,它直接接受 Runnable。提交后无法跟踪任务状态或获取返回值。
submit(Runnable task, T result)ForkJoinTask<T>异步提交一个 Runnable 任务,并提供一个预定义的结果对象。当任务完成时,调用 future.get() 将返回你预先传入的 result 对象。
submit(Callable<T> task)ForkJoinTask<T>异步提交一个 Callable 任务,可以返回一个结果。返回的 ForkJoinTask 是一个 Future,可通过 get() 获取 Callable 的计算结果。
invokeAll(ForkJoinTask... tasks)<T> T (first task’s)提交多个任务阻塞等待所有任务完成。此方法会等待 tasks 数组中的每一个任务都执行完毕后才返回。返回值是第一个任务的结果(如果是 RecursiveTask)。
生命周期管理
shutdown()void开始一个有序的关闭过程。停止接受新任务,但会继续执行已经提交到队列中和正在执行中的所有任务。线程池状态变为 SHUTDOWN
shutdownNow()List<Runnable>尝试立即停止所有活动任务。它尝试通过中断工作线程来停止正在执行的任务,并返回一个包含所有未开始执行的任务的列表。线程池状态变为 STOP
awaitTermination(long timeout, TimeUnit unit)boolean等待线程池终止。在调用 shutdown() 后,使用此方法来阻塞当前线程,直到线程池完全终止,或者指定的超时时间到达。返回 true 表示成功终止,false 表示超时。
isShutdown()boolean查询线程池是否已经开始关闭(即已经调用了 shutdown()shutdownNow())。
isTerminated()boolean查询线程池是否已经完全终止。即所有任务都已结束(包括正在执行和队列中等待的),并且所有工作线程都已销毁。
isQuiescent()boolean查询线程池是否处于空闲状态(Quiescent)。这意味着当前没有任务正在执行,也没有任务准备好可以执行(即所有工作线程的队列都是空的)。与 isTerminated() 不同,线程池此时仍然处于活动状态,可以接收新任务。

ForkJoinTask

ThreadPoolExecutor线程池提交的任务都是Runable或者Callable的实现类,同样的ForkJoinPool提交的任务也是需要是ForkJoinTask的继承类,当然它也可以接受Runable或者Callable的实现类。

ForkJoinTask是一个抽象类定义了继承类需要实现的方法模板,根据需不需要返回值jdk又提供了ForkJoinTask两个继承子类:

  1. RecursiveAction:继承这个类实现compute方法,该方法没有返回值
  2. **RecursiveTask: **继承这个类实现compute方法,该方法有返回值

所以我们如果需要提交执行任务,那么提交的任务只需要根据需不需要返回值继承上面两个类中的一个即可,类比Runable和Callable。

常用方法

方法名返回类型功能描述与使用场景
核心方法
fork()ForkJoinTask<T>分叉 / 异步提交。将当前任务作为一个新的子任务提交到 ForkJoinPool 中。此方法不会阻塞,会立即返回,调用线程可以继续执行其他工作。这是将任务分解成多个小任务的关键方法。
join()<T> T连接 / 等待结果。调用此方法的线程会阻塞,直到此任务的计算完成。如果任务是 RecursiveTask(有返回值),则返回其计算结果;如果任务是 RecursiveAction(无返回值),则返回 null重要: join() 方法会处理任务执行期间抛出的任何 unchecked exception,并将其包装在 ExecutionException 中重新抛出。
辅助生命周期方法
helpQuiesce()ForkJoinTask<T>辅助终止。当调用此方法时,会阻塞调用线程,直到所有由当前任务派生的子任务(非其兄弟任务)都已执行完毕。这有助于在提交一个任务后,确保它的所有“后代”任务都被处理完。
quietlyJoin()void静默的 join。行为与 join() 类似,会使调用线程阻塞直到任务完成。但是,如果任务执行过程中抛出异常,此方法会静默地忽略它。通常仅在你不关心任务的执行异常时使用。
获取结果与状态
get()T标准的 Future 方法,获取任务的最终结果。如果任务尚未完成,则阻塞调用线程。如果任务在执行过程中抛出了异常,get() 会将其包装在 ExecutionException 中重新抛出。
get(long timeout, TimeUnit unit)T带超时的 get 方法。在指定的时间内等待任务完成并获取结果。如果在超时时间内任务仍未完成,则抛出 TimeoutException
isDone()boolean标准 Future 方法。检查任务是否已经完成(无论成功、失败还是被取消)。
isCompletedNormally()boolean检查任务是否已成功完成且没有抛出异常。如果任务完成时抛出了异常,此方法返回 false
isCompletedAbnormally()boolean检查任务是否因为抛出异常或被取消而异常终止
isCancelled()boolean检查任务是否被成功取消
getException()Throwable如果任务由于抛出异常而终止,则返回抛出的异常。如果任务已正常完成,则返回 null。如果任务尚未完成,则抛出 IllegalStateException

使用示例

假设:我们要计算 1 到 1 亿的和,为了加快计算的速度,我们自然想到算法中的分治原理,将 1 亿个数字分成 1 万个任务,每个任务计算 1 万个数值的综合,利用 CPU 的并发计算性能缩短计算时间。

// 并行数组求和
class ArraySumTask extends RecursiveTask<Long> {
    private static final int THRESHOLD = 1000;
    private final int[] array;
    private final int start, end;
    
    ArraySumTask(int[] array, int start, int end) {
        this.array = array; this.start = start; this.end = end;
    }
    
    @Override
    protected Long compute() {
        if (end - start <= THRESHOLD) {
            // 直接计算小任务
            long sum = 0;
            for (int i = start; i < end; i++) sum += array[i];
            return sum;
        } else {
            // 拆分任务
            int mid = (start + end) / 2;
            ArraySumTask left = new ArraySumTask(array, start, mid);
            ArraySumTask right = new ArraySumTask(array, mid, end);
            
            left.fork(); // 异步执行左半部分
            long rightResult = right.compute(); // 同步计算右半部分
            long leftResult = left.join(); // 获取左半部分结果
       
            return leftResult + rightResult;
        }
    }
}


// 执行任务
public class Main {
    public static void main(String[] args) {
        ForkJoinPool forkJoinPool = new ForkJoinPool();
        ForkJoinTask<Long> submit = forkJoinPool.submit(new ArraySumTask(1, 100000000)); // 提交任务给线程池
        Long result = submit.join(); // 等待任务执行完成
        System.out.println(result);
        forkJoinPool.shutdown(); // 关闭线程池
    }
}

工作原理

当我们通过 ForkJoinPool 的 invoke 或 submit 方法提交任务时,ForkJoinPool 会根据一定的路由规则将任务分配到一个任务队列中。如果任务执行过程中创建了子任务,那么子任务会被提交到对应工作线程的任务队列中。

ThreadPoolExecutor中,多个线程共用一个阻塞任务队列,执行完任务后再从阻塞队列中拉取任务。而ForkJoinPool中每一个线程都有一个自己的任务队列,当线程发现自己的队列里没有任务了,才会到别的线程的队列里获取任务执行。

ForkJoinWorkerThread类似与ThreadPoolExecutor中Worker,每一个ForkJoinPool 中的每一个ForkJoinWorkerThread就是一个线程,对应着有自己的工作队列workQueue

public class ForkJoinWorkerThread extends Thread {
    
    // 所在的线程池
    final ForkJoinPool pool;

    // 当前线程下的任务队列
    final ForkJoinPool.WorkQueue workQueue;

    // 初始化时的构造方法
    protected ForkJoinWorkerThread(ForkJoinPool pool) {
        super("aForkJoinWorkerThread");
        this.pool = pool;
        this.workQueue = pool.registerWorker(this);
    }
}

static final class WorkQueue {
    int stackPred;          // 队列顶部(拥有者端)
    int config;                
    int base;              // 队列底部 (窃取端)    
    ForkJoinTask<?>[] array;   // 任务数组
    final ForkJoinWorkerThread owner // 拥有者线程
}

ForkJoinPool 中有一个数组形式的成员变量 workQueue[](ExecutorThreadPool中是用hashset存储线程),其对应一个队列数组,每个队列对应一个消费线程。丢入线程池的任务,根据特定规则进行转发。
在这里插入图片描述
当工作线程的任务队列为空时,它是否无事可做呢?

肯定的不是,如果是这样实现效率不就低了么。ForkJoinPool 引入了一种称为"任务窃取"的机制。当工作线程空闲时,它可以从其他工作线程的任务队列中"窃取"任务,这样就保证线程不会空闲占用系统资源。
在这里插入图片描述
为了避免多线程窃取任务导致的数据竞争问题,ForkJoinPool 中的任务队列采用双端队列的形式。工作线程从任务队列的一个端获取任务,而"窃取任务"从另一端进行消费。

CompletableFuture异步编程

基本介绍

之前想要一个线程执行任务并返回结果,就需要将任务封装为一个FutureTask交给线程执行,然后使用get阻塞获取到结果。但是这样有一个弊端,主线程get()得到结果需要一直阻塞等待,即使使用isDone()方法轮询去查看线程执行状态,但是这样也非常浪费cpu资源。

在这里插入图片描述
当Future的线程进行了一个非常耗时的操作,那我们的主线程也就阻塞了。 当我们在简单业务上,可以使用Future的另一个重载方法get(long,TimeUnit)来设置超时时间,避免我们的主线程被无穷尽地阻塞。

并且遇到将两个异步计算合并为一个,这两个异步计算之间相互独立,同时第二个又依赖于第一个的结果场景,又或者等待Future集合中的所有任务都完成,面对这些场景Future不适用了。

所以说java8引入了CompletableFuture,,结合了Future的优点,提供了非常强大的Future的扩展功能,可以帮助我们简化异步编程的复杂性,提供了函数式编程的能力,可以通过回调的方式处理计算结果。

基本使用

创建实例

1.通过构造函数创建

CompletableFuture<String> future = new CompletableFuture<>(); // 无参构造函数

2.通过静态方法创建

public static <U> CompletableFuture<U> supplyAsync(Supplier<U> supplier);
public static <U> CompletableFuture<U> supplyAsync(Supplier<U> supplier, Executor executor);

public static CompletableFuture<Void> runAsync(Runnable runnable);
public static CompletableFuture<Void> runAsync(Runnable runnable, Executor executor);

参数讲解:

  • supplier: labmada表达式

  • Runnable: 实现Runnable类的任务

  • Executor: 可以指定线程池,如果不指定默认使用系统的公共线程池ForkJoinPool,而且这些线程都是守护线程。

有两种格式,一种是supply开头的方法,一种是run开头的方法

  • supply开头:这种方法,可以返回异步线程执行之后的结果
  • run开头:这种不会返回结果,就只是执行线程任务

获取结果

方法说明
get()阻塞直到获得结果。如果任务失败,抛出 ExecutionException
get(long timeout, TimeUnit unit)带超时的 get。超时则抛出 TimeoutException
join()阻塞直到获得结果。如果任务失败,抛出未检查的 CompletionExceptionCancellationException。语法上更简洁。
getNow(T valueIfAbsent)非阻塞尝试获取结果。如果任务已完成,则返回结果;否则返回传入的 valueIfAbsent
isDone()检查任务是否完成(无论成功或失败)。
isCancelled()检查任务是否被取消。
complete(T value)手动完成CF。如果已经完成,则无效。返回一个 boolean 表示是否设置成功。
completeExceptionally(Throwable ex)手动以异常完成CF

计算完成后续操作-complete

1. 无论前阶段是正常完成还是发生异常,都会执行此 BiConsumer。不会修改阶段的结果。结果会原封不动地传递下去。
public CompletableFuture<T>   whenComplete(BiConsumer<? super T,? super Throwable> action)  // 同步处理
public CompletableFuture<T>   whenCompleteAsync(BiConsumer<? super T,? super Throwable> action) // 异步处理返回结果,使用默认线程池
public CompletableFuture<T>   whenCompleteAsync(BiConsumer<? super T,? super Throwable> action, Executor executor) // 异步处理结果指定线程池

简单来说就是把上一步执行的结果或者异常,作为参数传给complate,让complate继续处理上一步的结果

示例:

CompletableFuture<Integer> future = CompletableFuture.supplyAsync(() -> {
    return 1 + 1;
});
CompletableFuture<Integer> future2 = future.whenCompleteAsync((result, error) -> {
    System.out.println("步骤一结果: " + result); // 处理上一步的结果
    if(error != null) {
        error.printStackTrace(); // 打印步骤1异常
    }
});

计算完成后续操作 - handle

// 无论前阶段是正常完成还是发生异常,都会执行此 BiFunction。第一个参数是结果,第二个参数是异常(null表示无异常)。
public <U> CompletableFuture<U>     handle(BiFunction<? super T,Throwable,? extends U> fn)
public <U> CompletableFuture<U>     handleAsync(BiFunction<? super T,Throwable,? extends U> fn)
public <U> CompletableFuture<U>     handleAsync(BiFunction<? super T,Throwable,? extends U> fn, Executor executor)

handle方法和complete方法区别在于,complete方法第二次处理数据后不能返回第二次处理的结果,handle处理第二阶段的数据后可以返回第二次处理的结果,但是注意啊,如果传递的是数组,对象这些,在complete做出修改了,那么complete返回的数据就是修改后的数据。如果知道java参数传递的原理就会清楚这点

示例:

CompletableFuture<Integer> future = CompletableFuture.supplyAsync(() -> {
    return 1 + 1;
});
CompletableFuture<Integer> complete = future.whenCompleteAsync((result, error) -> {
    System.out.println("步骤一结果: " + result);
    if(error != null) {
        error.printStackTrace(); // 打印步骤1异常
    }
    result = result + 1;
    // return result + 1; 报错,因为whenCompleteAsync返回值只能是void
});

CompletableFuture<Integer> handle = future.handle((result, error) -> {
    System.out.println("步骤一结果: " + result);
    if(error != null) {
        error.printStackTrace(); // 打印步骤1异常
    }
    return result + 1;
});

// 阻塞等待
Integer i = complete.get(); // 结果是 2
Integer i2 = handle.get(); // 结果是 3

计算完成后续操作-apply

// apply方法和handle方法一样,都是结束计算之后的后续操作,不同是,handle方法会给出异常,可以让用户自己在第二步内部处理
// 而apply方法只有一个返回结果,如果上一步异常了,会被直接抛出,这一步就不会被执行
public <U> CompletableFuture<U>     thenApply(Function<? super T,? extends U> fn)
public <U> CompletableFuture<U>     thenApplyAsync(Function<? super T,? extends U> fn)
public <U> CompletableFuture<U>     thenApplyAsync(Function<? super T,? extends U> fn, Executor executor)

简单来说就是,apply是上一步如果出现异常了,异常需要上一步中处理,如果不处理程序就抛出异常结束执行,不会到apply这一步handle和complete是上一步出现异常了没有处理而是收集起来程序还是正常执行下去,并且异常作为参数传递到给handle或complete中处理。

示例:

// 1.无异常
CompletableFuture<Integer> future = CompletableFuture.supplyAsync(() -> {
    return 1 + 1;
});

CompletableFuture<Integer> apply = future.thenApply(result -> { // 只接受结果 
    return result + 3;
});
Integer i1 = apply.get(); // 结果 5
System.out.println(i1);


2.上层出现异常
CompletableFuture<Integer> future = CompletableFuture.supplyAsync(() -> {
    int a = 1/0; // 这里抛出了异常
    return 1 + 1;
});

CompletableFuture<Integer> apply = future.thenApply(result -> { // 第二步不会执行
    return result + 3;
});
Integer i1 = apply.get(); // 没有结果程序抛出异常
System.out.println(i1);

apply和handle一样都可以返回新值,传递给后续步骤处理。

计算完成后续操作-accept

// 上一步处理的结果,这里消费结果,不会返回新值
public CompletableFuture<Void>  thenAccept(Consumer<? super T> action)
public CompletableFuture<Void>  thenAcceptAsync(Consumer<? super T> action)
public CompletableFuture<Void>  thenAcceptAsync(Consumer<? super T> action, Executor executor)

示例:

CompletableFuture<Integer> future = CompletableFuture.supplyAsync(() -> {
    return 1 + 1;
});

CompletableFuture<Integer> complete = future.whenComplete((result, error) -> {
    System.out.println(result);
});
Integer integer = complete.get(); // 2

CompletableFuture<Void> accept = future.thenAccept(result -> {
    System.out.println(result);
});
Void unused = accept.get(); // null

和前面的异同:

  1. accept和apply一样不能收到上一步处理的异常,complete和handle可以收到上一步的异常,所以上一步出现了异常没有这里程序直接就报错了,不会到accept和apply
  2. accept不能返回新处理的值,它返回的是CompletableFuture,get或者join就会发现是null,complete虽然也不能返回新值,但是它可以返回是CompletableFuture,这里的T其就是上一步的结果

捕获异常结果-exceptionally

// 捕捉到所有中间过程的异常,方法会给我们一个异常作为参数,我们可以处理这个异常,同时返回一个默认值
public CompletableFuture<T> exceptionally(Function<Throwable, ? extends T> fn)

示例:

CompletableFuture<Integer> future = CompletableFuture.supplyAsync(() -> {
    // 返回null
    return null;
});

CompletableFuture<String> exceptionally = future.thenApply(result -> {
    // 制造一个空指针异常NPE
    int i = result;
    return i;
}).thenApply(result -> {
    // 这里不会执行,因为上面出现了异常
    String words = "现在是" + result + "点钟";
    return words;
}).exceptionally(error -> {
    // 我们选择在这里打印出异常
    error.printStackTrace();
    // 并且当异常发生的时候,我们返回一个默认的文字
    return "出错啊~";
});

exceptionally.thenAccept(System.out::println);

也就是说我们可以把这个方法放在所有处理步骤之后,如果有一步出错了,直接就跳到这个exceptionally方法中处理了,如果没有异常这个方法就不会处理,有点异常cath的意思。

任务编排

组合多个completableFuture任务编排

apply组合

apply可以返回处理后的新值并包装为一个新的CompletableFuture<T>所以可以链式编程一般,将任务分为好几步完成,每一步都需要等待上一步处理完成才能进行。所以是串行执行的

示例:

CompletableFuture<Integer> future = CompletableFuture.supplyAsync(() -> {
    // 第一步处理
    return 1 + 1;
}).thenApply(result -> {
    // 第二步处理
    return result * 2;
}).thenApply(result -> {
    // 第三步处理
    return result + 3;
});

future.thenAccept(result -> {
    System.out.println(result);
});

compose组合

它也是和apply类似处理上一步的结果,返回对应的值,不同的是apply可以直接返回新值他会自动包装为CompletableFuture,而compose需要封装为CompletableFuture才能返回。

示例:

public static void main(String[] args) throws InterruptedException, ExecutionException {

    CompletableFuture<Integer> future = CompletableFuture.supplyAsync(() -> {
        // 第一步处理
        return 1 + 1;
    });

    CompletableFuture<Integer> compose = future.thenCompose(x -> {
        return CompletableFuture.supplyAsync(() -> { // 手动封装为CompletableFuture类
            // 第二步处理
            return x + 1;
        });
    }).thenCompose(x -> {
        return CompletableFuture.supplyAsync(() -> {
            // 第三步处理
            return x + 1;
        });
    });

    System.out.println(compose.get());
}

为啥要封装为CompletableFuture<T>返回?使用apply不是更简单一点吗,只能说使用场景不同

看一个例子:

class UserService {
    public static CompletableFuture<User> getUserById(int id) {
        return CompletableFuture.supplyAsync(() -> {
            return new User(id, "User " + id);
        });
    }
}
class OrderService {
    public static CompletableFuture<Order[]> getOrdersByUser(User user) {
        return CompletableFuture.supplyAsync(() -> {
            return new Order[]{new Order("Order-001"), new Order("Order-002")};
        });
    }
}
public class ThenApplyExample {
    public static void main(String[] args) {
        ExecutorService executor = Executors.newFixedThreadPool(2);
        CompletableFuture<User> userFuture = UserService.getUserById(1); // 获取用户
        // 使用 thenApply
        CompletableFuture<CompletableFuture<Order[]>> nestedFuture = userFuture.thenApply(user -> OrderService.getOrdersByUser(user)); // 获取用户订单数据,这里就返回的是嵌套对象
        CompletableFuture<Order>> ordersFuture = nestedFuture.join(); // 这里需要处理嵌套
        Order[] orderList = ordersFutureAsync.join(); // 第二次join才能得到最终的订单数
    }
}

使用compose获取

public class ThenComposeExample {
    public static void main(String[] args) {
        ExecutorService executor = Executors.newFixedThreadPool(2);
        CompletableFuture<User> userFuture = UserService.getUserById(1); // 获取用户
        // 使用 thenCompose
        CompletableFuture<Order[]> ordersFuture = userFuture.thenCompose(user -> OrderService.getOrdersByUser(user)); // 获取订单这里返回的对象就是订单future
        Order[] finalOrders = ordersFuture.join(); // 一次join就得到订单数据了
        executor.shutdown();
    }

通过上面这个例子就可以看到,apply对返回的数据会包装为一个CompletableFuture,如果返回的T也是CompletableFuture就会嵌套包装为CompletableFuture<CompletableFuture>,而compose也是返回CompletableFuture但是这个是我们手动封装CompletableFuture的对象,它不会进行二次封装。

combine组合

apply和compose都是串行,并没有实现多个任务独立执行的效果。combine就是将两个独立任务合的结果整合到一个任务中执行。

使用示例:

public class BasicExample {
    public static void main(String[] args) throws Exception {
        // 两个独立的异步计算
        CompletableFuture<Integer> future1 = CompletableFuture.supplyAsync(() -> {
            System.out.println("计算任务1: 10 * 2");
            return 10 * 2;
        });
        
        CompletableFuture<Integer> future2 = CompletableFuture.supplyAsync(() -> {
            System.out.println("计算任务2: 5 + 8");
            return 5 + 8;
        });
        
        // 合并两个结果
        CompletableFuture<Integer> combinedFuture = future1.thenCombine(future2, 
            (result1, result2) -> {
                System.out.println("合并结果: " + result1 + " + " + result2);
                return result1 + result2;
            });
        
        System.out.println("最终结果: " + combinedFuture.get()); // 输出: 33
    }
}

等待所有任务完成 - allOf

allOf方法,当所有给定的任务完成后,返回一个全新的已完成CompletableFuture

public static CompletableFuture<Void> allOf(CompletableFuture<?>... cfs)

使用示例:

CompletableFuture<Integer> future1 = CompletableFuture.supplyAsync(() -> {
            try {
                //使用sleep()模拟耗时操作
                TimeUnit.SECONDS.sleep(2);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            return 1;
        });

        CompletableFuture<Integer> future2 = CompletableFuture.supplyAsync(() -> {
            return 2;
        });
        CompletableFuture.allOf(future1, future1);
        // 输出3
        System.out.println(future1.join()+future2.join());

获取率第一个完成任务的结果——anyOf

仅等待Future集合种最快结束的任务完成(有可能因为他们试图通过不同的方式计算同一个值),并返回它的结果。
小贴士 :如果最快完成的任务出现了异常,也会先返回异常,如果害怕出错可以加个exceptionally() 去处理一下可能发生的异常并设定默认返回值

public static CompletableFuture<Object> anyOf(CompletableFuture<?>... cfs)
CompletableFuture<Integer> future = CompletableFuture.supplyAsync(() -> {
            throw new NullPointerException();
        });

        CompletableFuture<Integer> future2 = CompletableFuture.supplyAsync(() -> {
            try {
                // 睡眠3s模拟延时
                TimeUnit.SECONDS.sleep(3);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            return 1;
        });
        CompletableFuture<Object> anyOf = CompletableFuture
                .anyOf(future, future2)
                .exceptionally(error -> {
                    error.printStackTrace();
                    return 2;
                });
        System.out.println(anyOf.join());

多组合场景

模拟一个多租户 SaaS 平台,处理用户订单、计算价格、发送通知等异步操作。

public class MultiTenantSaaSExample {
    
    // 租户信息
    static class Tenant {
        String id;
        String name;
        String plan; // BASIC, PRO, ENTERPRISE
        double discountRate;
    }
    
    // 用户订单
    static class Order {
        String orderId;
        String tenantId;
        String userId;
        List<OrderItem> items;
        double totalAmount;
        double finalAmount;
        String status;
    }
    
    static class OrderItem {
        String productId;
        String productName;
        double price;
        int quantity;
    }
    
    // 模拟数据库和服务类
    static class TenantService {
        // 获取租户信息
        static CompletableFuture<Tenant> getTenantAsync(String tenantId) {
            return CompletableFuture.supplyAsync(() -> {
                // 模拟不同租户方案
                switch (tenantId) {
                    case "tenant_001": return new Tenant(tenantId, "公司A", "PRO", 0.1);
                    case "tenant_002": return new Tenant(tenantId, "公司B", "BASIC", 0.0);
                    case "tenant_003": return new Tenant(tenantId, "公司C", "ENTERPRISE", 0.15);
                    default: return new Tenant(tenantId, "默认公司", "BASIC", 0.0);
                }
            });
        }
    }
    
    static class UserService {
        // 获取用户邮箱
        static CompletableFuture<String> getUserEmailAsync(String userId) {
            return CompletableFuture.supplyAsync(() -> {
                return userId + "@company.com";
            });
        }
        
        // 检查用户权限
        static CompletableFuture<Boolean> checkUserPermissionAsync(String userId, String permission) {
            return CompletableFuture.supplyAsync(() -> {
                try { Thread.sleep(100); } catch (InterruptedException e) {}
                return !userId.contains("blocked"); // 模拟权限检查
            });
        }
    }
    
    static class PricingService {
        // 应用折扣
        static CompletableFuture<Double> applyDiscountAsync(double amount, Tenant tenant) {
            return CompletableFuture.supplyAsync(() -> {
                return amount * (1 - tenant.discountRate);
            });
        }
        
        // 计算税费
        static CompletableFuture<Double> calculateTaxAsync(double amount) {
            return CompletableFuture.supplyAsync(() -> {
                return amount * 0.1; // 10% 税费
            });
        }
    }
    
    static class NotificationService {
        // 发送邮件通知
        static CompletableFuture<Void> sendEmailAsync(String email, String subject, String content) {
            return CompletableFuture.runAsync(() -> {
                System.out.println(Thread.currentThread().getName() + " - 发送邮件到: " + email);
                System.out.println("主题: " + subject);
                System.out.println("内容: " + content);
            });
        }
        
        // 发送短信通知
        static CompletableFuture<Void> sendSMSAsync(String phone, String message) {
            return CompletableFuture.runAsync(() -> {
                System.out.println(Thread.currentThread().getName() + " - 发送短信到: " + phone);
                System.out.println("短信内容: " + message);
            });
        }
    }
    
    static class AuditService {
        // 记录审计日志
        static CompletableFuture<Void> logAuditAsync(String action, String tenantId, String details) {
            return CompletableFuture.runAsync(() -> {
                System.out.println(Thread.currentThread().getName() + " - 记录审计日志: " + action);
                System.out.println("租户: " + tenantId + ", 详情: " + details);
                try { Thread.sleep(50); } catch (InterruptedException e) {}
            });
        }
    }
    
    // 主处理服务
    static class OrderProcessingService {
        
        /**
         * 处理订单的完整流程
         */
        public static CompletableFuture<Void> processOrderAsync(Order order) {
            System.out.println("开始处理订单: " + order.orderId + ", 租户: " + order.tenantId);
            
            // 1. thenCompose - 链式依赖操作:先获取租户信息,然后基于租户信息计算价格
            CompletableFuture<Order> pricedOrderFuture = TenantService.getTenantAsync(order.tenantId)
                .thenCompose(tenant -> {
                    System.out.println("获取到租户: " + tenant.name + ", 方案: " + tenant.plan);
                    return applyPricingAsync(order, tenant);
                });
            
            // 2. thenApply - 转换结果:在价格计算完成后,更新订单状态
            CompletableFuture<Order> processedOrderFuture = pricedOrderFuture
                .thenApply(processedOrder -> {
                    processedOrder.status = "PRICED";
                    System.out.println("订单定价完成: " + processedOrder.orderId + 
                                     ", 最终金额: " + processedOrder.finalAmount);
                    return processedOrder;
                });
            
            // 3. thenCombine - 合并两个独立操作:检查用户权限和获取用户邮箱
            CompletableFuture<Order> validatedOrderFuture = processedOrderFuture
                .thenCombine(
                    UserService.checkUserPermissionAsync(order.userId, "CREATE_ORDER"),
                    (orderResult, hasPermission) -> {
                        if (!hasPermission) {
                            throw new SecurityException("用户 " + order.userId + " 没有创建订单的权限");
                        }
                        orderResult.status = "VALIDATED";
                        System.out.println("订单验证通过: " + orderResult.orderId);
                        return orderResult;
                    }
                );
            
            // 4. thenCompose + thenApply - 复杂的链式操作
            CompletableFuture<Void> notificationFuture = validatedOrderFuture
                .thenCompose(validatedOrder -> 
                    UserService.getUserEmailAsync(validatedOrder.userId)
                        .thenApply(email -> {
                            System.out.println("获取到用户邮箱: " + email);
                            return new Object[]{validatedOrder, email};
                        })
                )
                .thenCompose(data -> {
                    Order orderData = (Order) data[0];
                    String email = (String) data[1];
                    
                    // 准备通知内容
                    String subject = "订单确认 - " + orderData.orderId;
                    String content = String.format(
                        "尊敬的客户,您的订单已确认。总金额: %.2f, 最终金额: %.2f",
                        orderData.totalAmount, orderData.finalAmount
                    );
                    
                    // 发送邮件通知
                    return NotificationService.sendEmailAsync(email, subject, content);
                });
            
            // 5. allOf - 等待多个并行操作完成
            CompletableFuture<Void> auditFuture = validatedOrderFuture
                .thenCompose(validatedOrder -> {
                    // 并行执行多个审计操作
                    CompletableFuture<Void> audit1 = AuditService.logAuditAsync(
                        "ORDER_CREATED", validatedOrder.tenantId, 
                        "订单创建: " + validatedOrder.orderId
                    );
                    
                    CompletableFuture<Void> audit2 = AuditService.logAuditAsync(
                        "PRICING_APPLIED", validatedOrder.tenantId,
                        "价格计算: " + validatedOrder.finalAmount
                    );
                    
                    // 使用 allOf 等待所有审计操作完成
                    return CompletableFuture.allOf(audit1, audit2);
                });
            
            // 6. thenAccept - 消费最终结果
            CompletableFuture<Void> finalStep = validatedOrderFuture
                .thenAccept(finalOrder -> {
                    finalOrder.status = "COMPLETED";
                    System.out.println("=== 订单处理完成 ===");
                    System.out.println("订单ID: " + finalOrder.orderId);
                    System.out.println("租户: " + finalOrder.tenantId);
                    System.out.println("总金额: " + finalOrder.totalAmount);
                    System.out.println("最终金额: " + finalOrder.finalAmount);
                    System.out.println("状态: " + finalOrder.status);
                });
            
            // 7. thenHandle - 异常处理
            CompletableFuture<Void> handledFuture = finalStep
                .thenHandle((result, exception) -> {
                    if (exception != null) {
                        System.err.println("订单处理失败: " + exception.getMessage());
                        // 记录失败日志
                        AuditService.logAuditAsync(
                            "ORDER_FAILED", order.tenantId, 
                            "失败原因: " + exception.getMessage()
                        );
                        throw new CompletionException(exception);
                    }
                    System.out.println("订单处理成功完成");
                    return result;
                });
            
            // 8. 使用 allOf 等待所有主要操作完成
            return CompletableFuture.allOf(
                notificationFuture,
                auditFuture,
                handledFuture
            );
        }
        
        /**
         * 应用价格计算:折扣 + 税费
         */
        private static CompletableFuture<Order> applyPricingAsync(Order order, Tenant tenant) {
            // 并行计算折扣和税费
            CompletableFuture<Double> discountedFuture = PricingService.applyDiscountAsync(
                order.totalAmount, tenant
            );
            
            CompletableFuture<Double> taxFuture = PricingService.calculateTaxAsync(order.totalAmount);
            
            // 使用 thenCombine 合并折扣和税费结果
            return discountedFuture.thenCombine(taxFuture, (discountedAmount, tax) -> {
                order.finalAmount = discountedAmount + tax;
                System.out.println(String.format(
                    "价格计算 - 原价: %.2f, 折扣后: %.2f, 税费: %.2f, 最终: %.2f",
                    order.totalAmount, discountedAmount, tax, order.finalAmount
                ));
                return order;
            });
        }
        
        /**
         * 批量处理多个租户的订单
         */
        public static CompletableFuture<Void> processBatchOrdersAsync(List<Order> orders) {
            System.out.println("\n=== 开始批量处理 " + orders.size() + " 个订单 ===");
            
            // 为每个订单创建处理任务
            List<CompletableFuture<Void>> orderFutures = orders.stream()
                .map(OrderProcessingService::processOrderAsync)
                .collect(Collectors.toList());
            
            // 使用 allOf 等待所有订单处理完成
            return CompletableFuture.allOf(
                orderFutures.toArray(new CompletableFuture[0])
            ).thenAccept(v -> {
                System.out.println("=== 批量订单处理全部完成 ===");
            });
        }
    }
}

参考文章

Logo

有“AI”的1024 = 2048,欢迎大家加入2048 AI社区

更多推荐