CodeGym /课程 /JAVA 25 SELF /ExecutorService、Callable、Future:启动任务

ExecutorService、Callable、Future:启动任务

JAVA 25 SELF
第 54 级 , 课程 1
可用

1. ExecutorService:像专业人士一样管理线程

为什么不应仅仅通过 new Thread 创建线程

在多线程的初期,一切看起来都很简单:

Thread t = new Thread(() -> {
    // 执行一些操作
});
t.start();

这种方式可以工作,但当任务很多时,很快就会成为负担。每次调用 new Thread() 都会创建一个新线程,而几十或上百个线程会给系统带来负担。而且管理也不方便:需要关注它们何时结束、出现错误时该怎么办、如何停止与复用它们。

这时登场的是 ExecutorService——智能的线程调度器。你只需把任务交给它,它会自行决定由哪个线程以及何时执行。结果是运行更快、更稳定,也更省心。

ExecutorService 的工作机制

ExecutorService 的原理简单而高效。

  • 其内部有一个线程池——预先创建的一组工作线程(可以是固定的或动态的)。
  • 任务进入队列,由空闲线程获取执行。
  • 服务管理生命周期:你可以等待完成、优雅地停止线程池并释放资源。

创建 ExecutorService

最常见的方式——使用 Executors 类的工厂方法:

import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

ExecutorService executor = Executors.newFixedThreadPool(4); // 4 个线程
  • newFixedThreadPool(N)——由 N 个线程组成的线程池(适合多数场景)。
  • newCachedThreadPool()——可缓存的动态线程池,按需创建线程(注意:任务暴增时可能会耗尽内存)。
  • newSingleThreadExecutor()——单线程(顺序执行)。

示例:通过 Runnable 提交到 ExecutorService

executor.submit(() -> {
    System.out.println("来自线程池的问候!");
});

当你使用完 ExecutorService 后,必须正确地关闭它:

executor.shutdown(); // 禁止添加新任务,等待当前任务完成

重要:如果不调用 shutdown(),程序可能不会退出——线程池中的线程会继续等待新任务。

2. Runnable 与 Callable:任务各不相同

在 Java 5 之前,如果你想在线程中执行某些操作,会实现接口 Runnable。它表示一个不返回结果且不抛出受检异常的任务。

Runnable task = () -> {
    System.out.println("只是干活,不返回任何东西!");
};
executor.submit(task);

Callable:带结果(且可抛出异常)的任务

有时我们希望任务不仅“做点什么”,还要返回结果——例如,计算和、计算结果、或从服务器获取数据。为此使用接口 Callable<T>

import java.util.concurrent.Callable;

Callable<Integer> sumTask = () -> {
    int sum = 0;
    for (int i = 1; i <= 100; i++) sum += i;
    return sum;
};
  • 方法 call() 返回类型为 T 的结果。
  • 方法 call() 可以抛出受检异常。

类比:Runnable——“去洗碗”(结果不重要),Callable——“去拿杯茶并告诉我它的温度”(结果很重要)。

启动 Callable:若要获取结果,请使用 executor.submit(...)。它会返回一个 Future<T> 对象。

3. Future:对结果的承诺

Future 是一种“对未来返回结果的承诺”。当你将任务提交给 ExecutorService 时,会得到一个 Future,稍后可通过它获取结果、查询任务是否完成,或取消任务。

Future 的主要方法

  • T get()——获取结果(会一直等待,直到任务完成)。
  • boolean isDone()——任务是否已完成。
  • boolean cancel(boolean mayInterruptIfRunning)——尝试取消任务。
  • boolean isCancelled()——任务是否已被取消。

示例:启动 Callable 并获取结果

import java.util.concurrent.*;

public class ParallelSumApp {
    public static void main(String[] args) throws Exception {
        ExecutorService executor = Executors.newFixedThreadPool(2);

        Callable<Integer> sumTask = () -> {
            int sum = 0;
            for (int i = 1; i <= 100; i++) sum += i;
            return sum;
        };

        Future<Integer> future = executor.submit(sumTask);

        System.out.println("任务已启动,可以去做点别的事...");

        // 获取结果(如果任务尚未完成,该方法会阻塞当前线程)
        Integer result = future.get();
        System.out.println("计算结果: " + result);

        executor.shutdown();
    }
}
  • 任务被提交到线程池。
  • 在任务执行期间,主线程可以去做其他事情。
  • 当需要结果时,调用 future.get()——如果任务仍在执行,线程会等待。
  • 一旦任务完成,就会返回结果。

4. 实战:多个任务并等待完成

经常需要一次性启动多个任务,并等待它们全部完成。例如,你要处理一个数据数组,将其拆分为多段,并分别在独立任务中计算每段之和。

示例:分块求数组元素之和

import java.util.*;
import java.util.concurrent.*;

public class ParallelArraySum {
    public static void main(String[] args) throws Exception {
        int[] array = new int[1000];
        Arrays.setAll(array, i -> i + 1); // 填充 1 到 1000 的数字

        ExecutorService executor = Executors.newFixedThreadPool(4);

        int chunkSize = array.length / 4;
        List<Future<Integer>> futures = new ArrayList<>();

        for (int i = 0; i < 4; i++) {
            int from = i * chunkSize;
            int to = (i == 3) ? array.length : (i + 1) * chunkSize;

            Callable<Integer> sumTask = () -> {
                int sum = 0;
                for (int j = from; j < to; j++) sum += array[j];
                System.out.println("从 " + from + " 到 " + (to - 1) + " 的和 = " + sum);
                return sum;
            };

            futures.add(executor.submit(sumTask));
        }

        int totalSum = 0;
        for (Future<Integer> f : futures) {
            totalSum += f.get(); // 按顺序等待每个任务
        }

        System.out.println("总和: " + totalSum);

        executor.shutdown();
    }
}

这里将数组分为 4 部分。对每一部分创建一个任务(Callable)来计算其和。所有任务提交给 ExecutorService,返回 Future。最后汇总所有任务的结果并相加。

在真实场景中,使用 invokeAll 更方便一次性等待全部任务完成。

5. 使用 Future 时的错误处理

当你调用 future.get() 时,如果任务以异常结束,它会被包装成 ExecutionException 抛出。这一点很重要:如果任务出了问题,你只会在调用 get() 时知道。

示例:异常处理

Callable<Integer> errorTask = () -> {
    throw new IllegalArgumentException("出现了一些问题!");
};

Future<Integer> badFuture = executor.submit(errorTask);

try {
    badFuture.get();
} catch (ExecutionException e) {
    System.out.println("任务以错误结束: " + e.getCause());
}
  • 在任务内部抛出了异常。
  • 调用 get() 时,它会被“包装”为 ExecutionException
  • 真正的原因可以通过 getCause() 获取。

6. 实用细节

如何取消任务

Future<?> f = executor.submit(() -> {
    while (true) {
        // 无限循环的工作
        if (Thread.currentThread().isInterrupted()) {
            System.out.println("有人请求我结束!");
            break;
        }
    }
});

Thread.sleep(100); // 稍等一下
f.cancel(true); // 尝试取消任务
  • cancel(true) 会尝试中断尚未完成的任务。
  • 在任务内部,建议检查 Thread.currentThread().isInterrupted() 并优雅退出。

shutdown 与 shutdownNow

shutdown()——柔和停止:禁止添加新任务,并让当前任务正常结束。最常用。

shutdownNow()——强制停止:尝试中断活动线程,并返回尚未开始的任务列表。请谨慎使用。

invokeAll 与 invokeAny

invokeAll(Collection<Callable<T>> tasks) 会启动所有传入的任务,并等待它们全部完成。返回 Future 列表。

invokeAny(Collection<Callable<T>> tasks) 只等待第一个成功完成的任务,返回其结果,并取消其余任务。适用于只关心第一个成功响应的场景。

7. 使用 ExecutorService、Callable 和 Future 的常见错误

错误 1:未关闭 ExecutorService。 如果忘记调用 shutdown(),程序在 main 结束后可能仍会“挂起”,因为线程池中的线程在等待新任务。

错误 2:提交任务后立刻等待结果。 如果在 submit() 之后立即调用 get(),就无法获得异步性的好处——线程仍会等待。应当并行去做有用的工作,在确实需要结果时再获取。

错误 3:忽略任务中的异常。 如果在调用 get() 时不处理 ExecutionException,可能会错过任务中发生的重要错误。

错误 4:在没有同步的情况下使用共享的可变变量。 如果多个任务处理同一份数据——需要同步或使用线程安全的集合。

错误 5:创建过多的线程。 不要将线程池大小设置得远大于 CPU 核心数——这甚至可能让执行更慢。

错误 6:忘记取消不再需要的任务。 如果任务已经不需要了,请通过 cancel() 取消它,以避免浪费资源。

评论
TO VIEW ALL COMMENTS OR TO MAKE A COMMENT,
GO TO FULL VERSION