跳到正文
hello world

15. Semaphore 实战:限制 PDF 并发处理数量

发布于阅读量 2

上一节已经把 Semaphore 的基本用法说清楚了。

它的核心就是:

acquire() 获取许可证;
release() 释放许可证;
许可证有几个,同一时间最多就允许几个线程进入。

这一节我把它放到 PDF 批量处理里,写一个相对完整一点的版本。

目标很明确:

input 目录下有多个 PDF;
每个 PDF 都可以作为一个任务;
但是同一时间最多只允许 3 个 PDF 正在处理;
主线程要等所有 PDF 都处理完,再统计总耗时。

上一节的代码只是限制了并发数量,但主线程没有等待所有任务结束。这一节补上这个问题。


这节要解决什么问题

如果我只是这样写:

for (File file : files) {
    new Thread(() -> {
        processPdf(file);
    }).start();
}

问题有两个。

一个是同时处理的 PDF 数量不可控。

如果 input 目录下有 50 个 PDF,就可能同时启动 50 个处理任务。

另一个是主线程不知道所有任务什么时候结束。

所以这节代码要同时解决两个点:

Semaphore:限制同一时间最多 3 个任务处理 PDF;
join:让 main 等所有子线程执行完。

注意,这里先用 join() 等待线程结束。

后面会学 CountDownLatch,它会比这种写法更适合等待一批任务。

但现在这个阶段,用 Thread + Semaphore + join 更容易看清楚执行过程。


完整代码

新建类:

com.succos.thread.SemaphorePdfLimitDemo

代码如下:

package com.succos.thread;

import java.io.File;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Semaphore;

public class SemaphorePdfLimitDemo {

    private static final Semaphore semaphore = new Semaphore(3);

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

        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;
        }

        long start = System.currentTimeMillis();

        List<Thread> threadList = new ArrayList<>();

        for (File file : files) {

            Thread thread = new Thread(() -> {

                boolean acquired = false;

                try {
                    System.out.println(file.getName()
                            + " 等待许可证,线程:"
                            + Thread.currentThread().getName());

                    semaphore.acquire();
                    acquired = true;

                    System.out.println(file.getName()
                            + " 拿到许可证,开始处理,线程:"
                            + Thread.currentThread().getName());

                    // 这里先用 sleep 模拟 PDF 处理
                    Thread.sleep(3000);

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

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

                } finally {
                    if (acquired) {
                        semaphore.release();

                        System.out.println(file.getName()
                                + " 释放许可证,线程:"
                                + Thread.currentThread().getName());
                    }
                }

            }, "pdf-thread-" + file.getName());

            threadList.add(thread);
        }

        for (Thread thread : threadList) {
            thread.start();
        }

        for (Thread thread : threadList) {
            thread.join();
        }

        long end = System.currentTimeMillis();

        System.out.println("--------------------------------");
        System.out.println("全部 PDF 处理完成");
        System.out.println("总文件数:" + files.length);
        System.out.println("总耗时:" + (end - start) + " ms");
    }
}

这段代码里,我暂时还是用:

Thread.sleep(3000);

模拟 PDF 处理。

先把并发控制看清楚,再替换成真正的水印方法。


执行流程

这里先创建了一个信号量:

private static final Semaphore semaphore = new Semaphore(3);

意思是同一时间最多允许 3 个线程进入处理逻辑。

然后每个 PDF 文件都会创建一个线程:

Thread thread = new Thread(() -> {
    // 处理当前 PDF
}, "pdf-thread-" + file.getName());

但是线程启动以后,不是马上就能处理 PDF。

它要先执行:

semaphore.acquire();

拿到许可证以后,才会继续往下执行。

如果许可证已经被前面的线程拿完了,就会在这里等待。


为什么要先收集线程,再统一 start

代码里我没有创建一个线程就马上 join()

而是先把线程放到集合里:

List<Thread> threadList = new ArrayList<>();

for (File file : files) {
    Thread thread = new Thread(() -> {
        // 任务逻辑
    }, "pdf-thread-" + file.getName());

    threadList.add(thread);
}

然后统一启动:

for (Thread thread : threadList) {
    thread.start();
}

最后统一等待:

for (Thread thread : threadList) {
    thread.join();
}

这个顺序很重要。

如果写成这样:

thread.start();
thread.join();

就会变成启动一个,等一个,再启动下一个。

那基本就没有并发效果了。

所以我现在会形成一个固定习惯:

先创建任务;
再统一启动;
最后统一等待。

这个习惯后面用 CompletableFuture.allOf() 时也一样。


Semaphore 在哪里起作用

真正限制并发的是这段:

semaphore.acquire();
acquired = true;

System.out.println(file.getName()
        + " 拿到许可证,开始处理,线程:"
        + Thread.currentThread().getName());

Thread.sleep(3000);

只有拿到许可证的线程,才能打印“开始处理”。

假设有 10 个 PDF,虽然我创建了 10 个线程,但同一时间最多只有 3 个线程能通过 acquire()

其他线程会卡在这里:

semaphore.acquire();

等前面的线程释放许可证。

这和普通锁不一样。

普通锁一般同一时间只允许 1 个线程进入。

而这里是:

同一时间允许 3 个线程进入。

这就是 Semaphore(3) 的意义。


release 必须放在 finally 里

这段代码也很关键:

finally {
    if (acquired) {
        semaphore.release();
    }
}

不管 PDF 处理成功还是失败,只要前面拿到了许可证,最后都要释放。

否则许可证会被占住。

比如 3 个线程都拿到了许可证,其中一个线程处理 PDF 时报错了,但没有释放许可证。

那可用许可证就从 3 个变成了 2 个。

如果这种错误多发生几次,后面的线程可能就一直卡在 acquire(),程序看起来就像“假死”了一样。

所以 release() 一定要放到 finally 里。

而且最好配合:

boolean acquired = false;

只有真正拿到了许可证,才释放。


为什么还要 join

Semaphore 只负责控制有几个线程能同时进入处理逻辑。

它不负责通知主线程“所有任务都结束了”。

所以我还需要:

for (Thread thread : threadList) {
    thread.join();
}

这段代码的作用是让 main 线程等待所有子线程结束。

如果没有这段,下面的统计代码就可能提前执行:

long end = System.currentTimeMillis();

System.out.println("全部 PDF 处理完成");
System.out.println("总耗时:" + (end - start) + " ms");

那统计出来的耗时就不准确。

所以这里其实用了两个工具配合:

Semaphore:控制同时处理的数量;
join:等待所有线程结束。

它们解决的问题不一样。


运行时能观察到什么

如果 input 目录下有 8 个 PDF,Semaphore(3) 的效果大概是:

第 1 批:3 个 PDF 开始处理;
其他 PDF 等待许可证;

第 1 批里有任务完成,释放许可证;
等待中的某个 PDF 拿到许可证,开始处理;

一直重复;
直到所有 PDF 都处理完成。

控制台日志可能类似:

a.pdf 拿到许可证,开始处理,线程:pdf-thread-a.pdf
b.pdf 拿到许可证,开始处理,线程:pdf-thread-b.pdf
c.pdf 拿到许可证,开始处理,线程:pdf-thread-c.pdf

d.pdf 等待许可证,线程:pdf-thread-d.pdf
e.pdf 等待许可证,线程:pdf-thread-e.pdf

a.pdf 处理完成
a.pdf 释放许可证

d.pdf 拿到许可证,开始处理

顺序不一定完全固定。

但只要同时观察“开始处理”的日志,就会发现同一时间最多只有 3 个任务在处理。


替换成真实 PDF 水印处理

现在把模拟处理换成真实的 PDF 水印方法。

需要引入:

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

可以把线程里的 Thread.sleep(3000) 替换成下面这段:

File outputDir = new File("output");

if (!outputDir.exists()) {
    outputDir.mkdirs();
}

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

完整的处理逻辑大概就是:

try {
    System.out.println(file.getName()
            + " 等待许可证,线程:"
            + Thread.currentThread().getName());

    semaphore.acquire();
    acquired = true;

    System.out.println(file.getName()
            + " 拿到许可证,开始处理,线程:"
            + Thread.currentThread().getName());

    File outputDir = new File("output");

    if (!outputDir.exists()) {
        outputDir.mkdirs();
    }

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

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

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

} finally {
    if (acquired) {
        semaphore.release();

        System.out.println(file.getName()
                + " 释放许可证,线程:"
                + Thread.currentThread().getName());
    }
}

这样就是真正的 PDF 并发处理了。


这版代码还有什么不足

这版代码比上一节完整一些,但它仍然不是最终方案。

最大的不足是:线程数量还是不受控制

比如有 1000 个 PDF,我还是会创建 1000 个线程。

虽然 Semaphore(3) 保证同一时间只有 3 个线程真正处理 PDF,但剩下的 997 个线程已经创建出来了,只是卡在 acquire() 等许可证。

这对资源还是不友好。

所以这里要分清楚:

Semaphore 控制的是同时进入处理逻辑的线程数量;
它不控制线程创建数量。

如果要控制线程数量、任务排队、线程复用,就要用线程池。

也就是说,这一节是一个过渡版本。

它帮助我理解“并发数量限制”这件事,但后面真正做批量任务,还是会转到 ThreadPoolExecutor


Semaphore 适合放在哪里

虽然 PDF 批量任务后面更适合用线程池,但 Semaphore 仍然有实际价值。

比如我已经有一个线程池,线程池最多 10 个线程。

但其中某个资源只能同时被 2 个线程访问,比如:

某个外部接口限流;
某个 Python 脚本不支持并发;
某个模型服务一次只能跑一个任务;
某个上传通道最多允许 3 个并发。

这时就可以在线程池任务内部再加一个 Semaphore

也就是说:

线程池控制整体任务执行;
Semaphore 控制某个局部资源的并发访问。

这个组合在真实项目里是有意义的。


这一节小结

这一节我主要记住几点:

1. Semaphore 可以限制同一时间最多几个任务进入处理逻辑;
2. join 可以让 main 等所有子线程处理完成;
3. acquire 和 release 要成对出现;
4. release 要放在 finally 里;
5. acquired 标记可以避免没有拿到许可证却错误释放;
6. Semaphore 不能控制线程创建数量;
7. PDF 批量处理最终更适合用线程池,但 Semaphore 很适合理解并发限制。

用一句话总结:

Semaphore 解决的是“同时干活的人不能太多”,但它不负责管理这些线程从哪里来。

下一节继续看 CountDownLatch

因为现在等待所有线程用的是 join(),任务一多就不太方便。CountDownLatch 更适合表达“等这一批任务全部完成”。