跳到正文
hello world

23. 用 ThreadPoolExecutor 批量处理 PDF 水印

发布于阅读量 0

23. 用 ThreadPoolExecutor 批量处理 PDF 水印

前面几节一直在讲线程池的概念。

到这里,可以把线程池真正接到 PDF 水印项目里了。

这一节的目标很简单:

遍历 input 目录下的 PDF;
每个 PDF 封装成一个任务;
提交给 ThreadPoolExecutor;
线程池控制同一时间最多处理 3 个 PDF;
处理后的文件输出到 output 目录;
最后等待所有任务执行完,统计总耗时。

这一版已经比前面的 Thread + Semaphore + CountDownLatch 更接近实际写法了。

因为这次不再手动给每个 PDF 创建线程,而是把每个 PDF 当成一个任务,交给线程池去调度。


这节代码的整体思路

先把流程捋一下。

我希望代码按这个顺序执行:

1. 找到 input 目录;
2. 过滤出 PDF 文件;
3. 创建 output 目录;
4. 创建 PDF 处理线程池;
5. 遍历每个 PDF,提交处理任务;
6. shutdown 线程池;
7. awaitTermination 等待所有任务结束;
8. 打印总文件数和耗时。

这里面最关键的变化是第 5 步。

以前我可能会这样写:

new Thread(() -> {
    processPdf(file);
}).start();

现在换成:

executor.execute(() -> {
    processPdf(file);
});

也就是说,我不直接创建线程了。

我只提交任务。

线程池内部会安排哪个线程来执行。


完整代码

新建类:

com.succos.threadpool.ThreadPoolPdfWatermarkDemo

代码如下:

package com.succos.threadpool;

import com.succos.dto.FileItemContext;
import com.succos.service.PdfWatermarkService;

import java.io.File;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;

public class ThreadPoolPdfWatermarkDemo {

    public static void main(String[] args) {

        File inputDir = new File("input");
        File outputDir = new File("output");

        if (!outputDir.exists()) {
            boolean created = outputDir.mkdirs();

            if (!created) {
                System.out.println("output 目录创建失败");
                return;
            }
        }

        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()
        );

        long start = System.currentTimeMillis();

        for (File file : files) {

            executor.execute(() -> {

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

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

                try {
                    String targetPath = addWatermark(file, outputDir);

                    System.out.println(threadName
                            + " 处理完成:"
                            + file.getName()
                            + ",输出路径:"
                            + targetPath);

                } catch (Exception e) {
                    System.out.println(threadName
                            + " 处理失败:"
                            + file.getName()
                            + ",原因:"
                            + e.getMessage());
                }
            });
        }

        executor.shutdown();

        try {
            boolean finished = executor.awaitTermination(1, TimeUnit.HOURS);

            long end = System.currentTimeMillis();

            if (finished) {
                System.out.println("全部 PDF 处理完成");
            } else {
                System.out.println("等待超时,还有 PDF 没有处理完");
            }

            System.out.println("--------------------------------");
            System.out.println("总文件数:" + files.length);
            System.out.println("总耗时:" + (end - start) + " ms");

        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }

    private static String addWatermark(File file, File outputDir) throws Exception {

        String targetPath = outputDir.getAbsolutePath()
                + File.separator
                + file.getName().replace(".pdf", "-watermark.pdf");

        FileItemContext fileItemContext = new FileItemContext();
        fileItemContext.setSourcePath(file.getAbsolutePath());
        fileItemContext.setTargetPath(targetPath);
        fileItemContext.setWaterMakeText("上下文网");

        PdfWatermarkService service = new PdfWatermarkService();
        service.addWaterMakerOfPDF(fileItemContext);

        return targetPath;
    }
}

这段代码就是线程池处理 PDF 的第一版。

如果你的 PdfWatermarkService 方法名和我这里不一样,把这一行换成你自己的方法即可:

service.addWaterMakerOfPDF(fileItemContext);

线程池配置

这里我用的线程池参数是:

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

这表示:

核心线程数是 3;
最大线程数也是 3;
队列最多放 100 个任务;
如果线程池和队列都满了,就让提交任务的线程自己执行。

因为核心线程数和最大线程数都等于 3,所以这个线程池的行为比较稳定:

同一时间最多 3 个 PDF 正在处理;
其他 PDF 任务进入队列等待;
线程处理完一个 PDF 后,再继续从队列取下一个。

这比每个 PDF 都 new Thread() 好很多。

线程不会随着文件数量无限增加。


每个 PDF 是一个任务

核心提交任务的代码是这一段:

for (File file : files) {

    executor.execute(() -> {

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

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

        try {
            String targetPath = addWatermark(file, outputDir);

            System.out.println(threadName
                    + " 处理完成:"
                    + file.getName()
                    + ",输出路径:"
                    + targetPath);

        } catch (Exception e) {
            System.out.println(threadName
                    + " 处理失败:"
                    + file.getName()
                    + ",原因:"
                    + e.getMessage());
        }
    });
}

这里每次循环都会提交一个 Runnable 任务。

这个任务只负责处理当前这个 PDF。

也就是说:

~~~text
a.pdf 是一个任务;
b.pdf 是一个任务;
c.pdf 是一个任务;