跳到正文
hello world

25. Future:怎么拿到异步任务的结果

发布于阅读量 1

25. Future:怎么拿到异步任务的结果

上一节讲了 executesubmit 的区别。

execute() 更像是把任务丢给线程池执行,不关心结果。

submit() 会返回一个 Future,这个 Future 就是后面拿任务结果的入口。

这一节专门整理 Future

我现在对它的理解是:

Future 不是结果本身,而是一个“将来可以拿到结果的凭证”。

任务提交到线程池以后,可能还没执行,也可能正在执行,也可能已经执行完。

但我先拿到了一个 Future 对象。

后面什么时候需要结果,就通过它去取。


先看一个最简单的 Future

新建类:

com.succos.threadpool.FutureDemo

代码如下:

package com.succos.threadpool;

import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.Future;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;

public class FutureDemo {

    public static void main(String[] args) {

        ThreadPoolExecutor executor = new ThreadPoolExecutor(
                3,
                3,
                60,
                TimeUnit.SECONDS,
                new ArrayBlockingQueue<>(100),
                new ThreadPoolExecutor.CallerRunsPolicy()
        );

        Future<String> future = executor.submit(() -> {
            System.out.println(Thread.currentThread().getName()
                    + " 开始处理 PDF");

            Thread.sleep(3000);

            System.out.println(Thread.currentThread().getName()
                    + " PDF 处理完成");

            return "output/test-watermark.pdf";
        });

        System.out.println("main 提交任务完成,继续往下执行");

        try {
            System.out.println("main 准备获取任务结果");

            String result = future.get();

            System.out.println("main 拿到任务结果:" + result);

        } catch (Exception e) {
            e.printStackTrace();
        }

        executor.shutdown();
    }
}

这里最关键的是:

Future<String> future = executor.submit(() -> {
    Thread.sleep(3000);
    return "output/test-watermark.pdf";
});

submit() 之后,我马上拿到了一个 Future<String>

但这个时候,PDF 任务不一定已经处理完。

任务真正执行是在工作线程里。


future.get() 是等待点

这行代码是重点:

String result = future.get();

它的意思是:

如果任务已经执行完,直接返回结果;
如果任务还没执行完,当前线程就在这里等待。

在这个例子里,执行 future.get() 的是 main 线程。

所以就是:

main 线程等待线程池里的任务执行完成。

等任务执行到:

return "output/test-watermark.pdf";

future.get() 才能拿到这个字符串。

所以 Future 的核心用法就是:

先 submit 提交异步任务;
再 get 获取任务结果。

Future 的泛型代表什么

这里是:

Future<String> future

这个 String 表示任务最终返回的是一个字符串。

因为我提交的任务写的是:

return "output/test-watermark.pdf";

所以返回值类型是 String

如果任务返回的是整数,就可以是:

Future<Integer> future

如果任务返回的是一个自定义结果对象,就可以是:

Future<PdfTaskResult> future

也就是说,Future<T> 里的 T,就是任务结果的类型。

这个和前面学泛型时的感觉是一样的。


Future 不是结果本身

这个地方我觉得要单独强调一下。

Future<String> 不是字符串。

它只是一个代表未来结果的对象。

比如:

Future<String> future = executor.submit(() -> {
    return "PDF 处理成功";
});

这时候 future 不是:

PDF 处理成功

真正的结果要通过:

future.get();

才能拿到。

所以我会这样区分:

Future:结果凭证;
future.get():真正取结果。

get 会阻塞当前线程

future.get() 最大的特点就是会阻塞。

比如任务里面睡 3 秒:

Thread.sleep(3000);

main 执行到:

future.get();

时,如果任务还没结束,main 就会停在那里等。

所以 get() 不能乱放。

如果我在循环里提交一个任务就马上 get(),并发效果会变差。

比如这种写法:

for (File file : files) {
    Future<String> future = executor.submit(() -> {
        processPdf(file);
        return "成功:" + file.getName();
    });

    String result = future.get();
    System.out.println(result);
}

看起来用了线程池,但实际效果接近顺序执行。

因为:

提交 a.pdf;
马上等待 a.pdf 完成;

提交 b.pdf;
马上等待 b.pdf 完成;

提交 c.pdf;
马上等待 c.pdf 完成。

所以批量任务里,我更倾向于:

先把所有任务都 submit;
把 Future 存起来;
最后再统一 get。

这个思路和前面 Thread 里的“先全部 start,再统一 join”是一样的。


批量 Future 的基本写法

新建类:

com.succos.threadpool.FutureBatchDemo

代码如下:

package com.succos.threadpool;

import java.io.File;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.Future;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;

public class FutureBatchDemo {

    public static void main(String[] args) {

        File inputDir = new File("input");

        File[] files = inputDir.listFiles(file ->
                file.isFile() && file.getName().toLowerCase().endsWith(".pdf")
        );

        if (files == null || files.length == 0) {
            System.out.println("input 目录下没有 PDF 文件");
            return;
        }

        ThreadPoolExecutor executor = new ThreadPoolExecutor(
                3,
                3,
                60,
                TimeUnit.SECONDS,
                new ArrayBlockingQueue<>(100),
                new ThreadPoolExecutor.CallerRunsPolicy()
        );

        List<Future<String>> futureList = new ArrayList<>();

        long start = System.currentTimeMillis();

        for (File file : files) {

            Future<String> future = executor.submit(() -> {
                System.out.println(Thread.currentThread().getName()
                        + " 开始处理:"
                        + file.getName());

                Thread.sleep(3000);

                System.out.println(Thread.currentThread().getName()
                        + " 处理完成:"
                        + file.getName());

                return "成功:" + file.getName();
            });

            futureList.add(future);
        }

        System.out.println("所有任务已经提交,开始统一获取结果");

        for (Future<String> future : futureList) {
            try {
                String result = future.get();
                System.out.println("任务结果:" + result);

            } catch (Exception e) {
                System.out.println("任务执行失败:" + e.getMessage());
            }
        }

        executor.shutdown();

        long end = System.currentTimeMillis();

        System.out.println("--------------------------------");
        System.out.println("全部结果收集完成");
        System.out.println("总文件数:" + files.length);
        System.out.println("总耗时:" + (end - start) + " ms");
    }
}

这段代码的结构比单个 Future 更接近实际批量任务。

流程是:

先遍历所有 PDF;
每个 PDF submit 一个任务;
把每个 Future 存进 List;
所有任务提交完以后;
再遍历 Future 获取结果。

这样线程池才能真正并发处理多个 PDF。


Future.get 的异常

future.get() 需要处理异常。

常见的是这两个:

InterruptedException
ExecutionException

InterruptedException 表示当前等待结果的线程被中断了。

ExecutionException 表示任务执行过程中抛了异常。

比如任务里面写了:

int x = 1 / 0;

任务线程会抛出异常。

但这个异常不会直接在 submit() 那一行抛出来,而是会被包装到 Future 里。

等我调用:

future.get();

时,再以 ExecutionException 的形式抛出来。

所以如果用了 submit(),但从来不调用 get(),任务里的异常可能不会那么明显。


单个任务失败不会影响其他任务提交

批量处理 PDF 时,某一个 PDF 失败,不应该影响其他 PDF。

比如:

a.pdf 成功;
b.pdf 损坏,失败;
c.pdf 继续成功。

如果每个任务都用自己的 Future,我就可以单独处理每个结果。

比如:

for (Future<String> future : futureList) {
    try {
        String result = future.get();
        System.out.println("任务结果:" + result);
    } catch (Exception e) {
        System.out.println("某个任务失败:" + e.getMessage());
    }
}

这样某个 future.get() 失败,只会进入当前这次循环的 catch

不会直接影响其他 Future 继续获取结果。

当然,真实项目里最好返回更明确的结果对象,而不是只打印异常。


用结果对象表示 PDF 处理结果

字符串可以学习用,但实际项目里我更喜欢返回对象。

比如定义一个内部类:

static class PdfTaskResult {

    private boolean success;

    private String fileName;

    private String targetPath;

    private String message;

    public static PdfTaskResult success(String fileName, String targetPath) {
        PdfTaskResult result = new PdfTaskResult();
        result.success = true;
        result.fileName = fileName;
        result.targetPath = targetPath;
        result.message = "处理成功";
        return result;
    }

    public static PdfTaskResult fail(String fileName, String message) {
        PdfTaskResult result = new PdfTaskResult();
        result.success = false;
        result.fileName = fileName;
        result.message = message;
        return result;
    }
}

这样一个 PDF 的处理结果就比较清楚:

是否成功;
文件名;
输出路径;
失败原因。

比单纯返回字符串更适合后续统计。


Future 返回结果对象示例

新建类:

com.succos.threadpool.FuturePdfResultDemo

代码如下:

package com.succos.threadpool;

import java.io.File;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.Future;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;

public class FuturePdfResultDemo {

    public static void main(String[] args) {

        File inputDir = new File("input");

        File[] files = inputDir.listFiles(file ->
                file.isFile() && file.getName().toLowerCase().endsWith(".pdf")
        );

        if (files == null || files.length == 0) {
            System.out.println("input 目录下没有 PDF 文件");
            return;
        }

        ThreadPoolExecutor executor = new ThreadPoolExecutor(
                3,
                3,
                60,
                TimeUnit.SECONDS,
                new ArrayBlockingQueue<>(100),
                new ThreadPoolExecutor.CallerRunsPolicy()
        );

        List<Future<PdfTaskResult>> futureList = new ArrayList<>();

        long start = System.currentTimeMillis();

        for (File file : files) {

            Future<PdfTaskResult> future = executor.submit(() -> {

                String threadName = Thread.currentThread().getName();

                System.out.println(threadName + " 开始处理:" + file.getName());

                try {
                    Thread.sleep(3000);

                    String targetPath = "output/"
                            + file.getName().replace(".pdf", "-watermark.pdf");

                    System.out.println(threadName + " 处理完成:" + file.getName());

                    return PdfTaskResult.success(file.getName(), targetPath);

                } catch (Exception e) {
                    return PdfTaskResult.fail(file.getName(), e.getMessage());
                }
            });

            futureList.add(future);
        }

        int successCount = 0;
        int failCount = 0;

        for (Future<PdfTaskResult> future : futureList) {
            try {
                PdfTaskResult result = future.get();

                if (result.success) {
                    successCount++;

                    System.out.println("成功:"
                            + result.fileName
                            + " -> "
                            + result.targetPath);
                } else {
                    failCount++;

                    System.out.println("失败:"
                            + result.fileName
                            + ",原因:"
                            + result.message);
                }

            } catch (Exception e) {
                failCount++;
                System.out.println("获取任务结果失败:" + e.getMessage());
            }
        }

        executor.shutdown();

        long end = System.currentTimeMillis();

        System.out.println("--------------------------------");
        System.out.println("全部 PDF 处理结束");
        System.out.println("总文件数:" + files.length);
        System.out.println("成功数量:" + successCount);
        System.out.println("失败数量:" + failCount);
        System.out.println("总耗时:" + (end - start) + " ms");
    }

    static class PdfTaskResult {

        private boolean success;

        private String fileName;

        private String targetPath;

        private String message;

        public static PdfTaskResult success(String fileName, String targetPath) {
            PdfTaskResult result = new PdfTaskResult();
            result.success = true;
            result.fileName = fileName;
            result.targetPath = targetPath;
            result.message = "处理成功";
            return result;
        }

        public static PdfTaskResult fail(String fileName, String message) {
            PdfTaskResult result = new PdfTaskResult();
            result.success = false;
            result.fileName = fileName;
            result.message = message;
            return result;
        }
    }
}

这里我先用 Thread.sleep(3000) 模拟处理。

如果要换成真实 PDF 水印,只需要把任务里的模拟逻辑换成 addWatermark(file),然后返回真实的 targetPath


为什么任务内部捕获异常并返回失败对象

在上面的代码里,我在任务内部写了:

try {
    Thread.sleep(3000);

    String targetPath = "output/"
            + file.getName().replace(".pdf", "-watermark.pdf");

    return PdfTaskResult.success(file.getName(), targetPath);

} catch (Exception e) {
    return PdfTaskResult.fail(file.getName(), e.getMessage());
}

这样做的好处是:无论任务成功还是失败,Future 最后都能返回一个 PdfTaskResult

如果任务失败,我不是让异常直接抛出去,而是转成一个失败结果。

这样最后统一统计时会比较方便。

当然,还有另一种写法:任务内部不捕获异常,让异常在 future.get() 时抛出来。

两种都可以。

我个人更喜欢在批量处理里返回统一结果对象。

因为这样每个任务都有明确状态。


Future 的局限

Future 能拿结果,但它也有局限。

比如:

get() 会阻塞;
多个 Future 组合不够方便;
任务完成后继续执行下一步不够自然;
异常处理写起来比较散;
超时、取消还要另外处理。

比如我想表达:

PDF 处理完成后,自动生成下载地址;
然后保存数据库;
然后通知用户。

Future 写起来就会比较啰嗦。

这也是后面要学习 CompletableFuture 的原因。

CompletableFuture 更适合异步任务编排。

但在学习线程池阶段,Future 还是很重要。

因为它是理解“异步任务结果”的基础。


Future 适合什么场景

我现在会这样看:

如果只是提交一个任务,并且后面要拿它的结果,Future 很合适。

比如:

异步处理一个 PDF,返回输出路径;
批量处理多个 PDF,最后统计成功失败;
提交几个查询任务,最后拿结果汇总。

但如果任务之间有复杂依赖,比如:

任务 A 完成后自动执行任务 B;
任务 A 和任务 B 都完成后合并结果;
任意一个任务完成就继续;
任务失败后走兜底逻辑。

Future 就不够顺手了。

这种场景更适合 CompletableFuture


这一节小结

这一节我主要记住几点:

1. Future 是异步任务的结果凭证,不是结果本身;
2. future.get() 才是真正获取结果;
3. get() 会阻塞当前线程;
4. Future<T> 里的 T 就是任务返回值类型;
5. 批量任务要先全部 submit,再统一 get;
6. Future 可以返回字符串,也可以返回自定义结果对象;
7. PDF 批量处理里,用结果对象更方便统计成功和失败。

用一句话总结:

Future 解决的是“任务交给线程池执行以后,我后面怎么拿到它的结果”这个问题。

下一节继续看 Future 的超时、取消和状态判断。

因为 get() 默认会一直等,如果某个 PDF 卡住了,程序不能无限等下去。