跳到正文
hello world

17. CyclicBarrier:让多个线程在某个位置互相等待

发布于阅读量 0

17. CyclicBarrier:让多个线程在某个位置互相等待

前面学了 CountDownLatch

它的使用场景比较好理解:主线程等一批任务全部完成。

比如:

main 线程等待所有 PDF 处理完;
等计数变成 0;
main 再继续统计结果。

这一节看另一个和“等待”有关的工具:CyclicBarrier

它和 CountDownLatch 有点像,但思路不一样。

我现在的理解是:

CountDownLatch 更像是一个线程等一批线程完成;
CyclicBarrier 更像是一批线程互相等待,等人齐了再一起往下走。

这个区别很重要。


CyclicBarrier 解决什么问题

假设有 3 个线程,它们各自先做一段准备工作。

准备完成以后,不能马上继续往下执行,而是要等另外两个线程也准备好。

等 3 个线程都到达同一个位置,再一起继续执行。

这种场景就适合 CyclicBarrier

它的名字里有个 Barrier,可以理解成“栅栏”或者“屏障”。

多个线程执行到屏障位置时,会先停下来。

等达到指定数量后,屏障打开,大家一起继续。


先写一个简单例子

新建类:

com.succos.thread.CyclicBarrierDemo

代码如下:

package com.succos.thread;

import java.util.concurrent.CyclicBarrier;

public class CyclicBarrierDemo {

    public static void main(String[] args) {

        CyclicBarrier barrier = new CyclicBarrier(3, () -> {
            System.out.println("3 个线程都准备好了,开始统一执行下一步");
        });

        for (int i = 1; i <= 3; i++) {

            int taskId = i;

            new Thread(() -> {

                try {
                    System.out.println("线程 " + taskId
                            + " 开始准备数据,线程名:"
                            + Thread.currentThread().getName());

                    Thread.sleep(taskId * 1000L);

                    System.out.println("线程 " + taskId
                            + " 准备完成,等待其他线程");

                    barrier.await();

                    System.out.println("线程 " + taskId
                            + " 开始执行正式任务,线程名:"
                            + Thread.currentThread().getName());

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

            }, "worker-" + taskId).start();
        }
    }
}

这里创建了一个 CyclicBarrier

CyclicBarrier barrier = new CyclicBarrier(3, () -> {
    System.out.println("3 个线程都准备好了,开始统一执行下一步");
});

意思是:

等 3 个线程都调用 barrier.await();
等人齐了以后,先执行一次回调;
然后这 3 个线程一起继续往下走。

await() 是集合点

代码里最关键的是这一行:

barrier.await();

每个线程执行到这里时,都会停下来等。

比如:

线程1 准备 1 秒后到了 await;
线程2 准备 2 秒后到了 await;
线程3 准备 3 秒后到了 await。

线程1 最先到,但它不能继续。

它要等线程2、线程3。

等第 3 个线程也到达 await() 后,屏障才会打开。

然后 3 个线程继续执行:

System.out.println("线程 " + taskId + " 开始执行正式任务");

所以 await() 可以理解成一个集合点:

大家先在这里集合;
人齐了;
再一起往后走。

运行结果大概是什么样

运行后,可能看到类似输出:

线程 1 开始准备数据,线程名:worker-1
线程 2 开始准备数据,线程名:worker-2
线程 3 开始准备数据,线程名:worker-3

线程 1 准备完成,等待其他线程
线程 2 准备完成,等待其他线程
线程 3 准备完成,等待其他线程

3 个线程都准备好了,开始统一执行下一步

线程 3 开始执行正式任务,线程名:worker-3
线程 1 开始执行正式任务,线程名:worker-1
线程 2 开始执行正式任务,线程名:worker-2

最后 3 个线程继续执行的顺序不一定固定。

但关键现象是:

在 3 个线程都到达 await() 之前,没有任何一个线程会执行“正式任务”。

这就是 CyclicBarrier 的作用。


和 CountDownLatch 的区别

这里一定要对比一下,不然后面容易混。

CountDownLatch 的典型场景是:

main 线程等多个子线程完成。

代码大概是:

CountDownLatch latch = new CountDownLatch(3);

latch.await();

每个子线程执行完以后:

latch.countDown();

主线程在等。

子线程一般不会互相等。


CyclicBarrier 的典型场景是:

多个子线程互相等待。

代码大概是:

CyclicBarrier barrier = new CyclicBarrier(3);

barrier.await();

每个线程都要执行到 await()

谁先到谁等。

等大家都到了,再一起继续。

所以我会这样记:

CountDownLatch:一个线程等一批任务结束;
CyclicBarrier:一批线程在某个位置互相等。

CountDownLatch 是倒计时,CyclicBarrier 是凑人数

从感觉上也可以这样区分。

CountDownLatch 像倒计时:

还有 3 个任务没完成;
完成一个,减到 2;
再完成一个,减到 1;
最后减到 0;
等待线程继续。

CyclicBarrier 像凑人数:

需要 3 个人到齐;
来了 1 个,等;
来了 2 个,等;
来了 3 个,人齐;
一起继续。

这两个工具都和等待有关,但表达的业务含义不一样。


为什么叫 Cyclic

CyclicBarrier 里的 Cyclic 是“循环”的意思。

它和 CountDownLatch 不一样。

CountDownLatch 用完一次就结束了,计数变成 0 后不能重置。

CyclicBarrier 可以重复使用。

比如我有 3 个线程,要分 3 个阶段执行。

每个阶段都要等 3 个线程全部完成后,才能进入下一阶段。

可以这样写:

package com.succos.thread;

import java.util.concurrent.CyclicBarrier;

public class CyclicBarrierReuseDemo {

    public static void main(String[] args) {

        CyclicBarrier barrier = new CyclicBarrier(3, () -> {
            System.out.println("所有线程到达屏障,进入下一阶段");
        });

        for (int i = 1; i <= 3; i++) {

            int workerId = i;

            new Thread(() -> {

                try {
                    for (int stage = 1; stage <= 3; stage++) {

                        System.out.println("线程 " + workerId
                                + " 开始第 " + stage + " 阶段");

                        Thread.sleep(workerId * 1000L);

                        System.out.println("线程 " + workerId
                                + " 完成第 " + stage + " 阶段,等待其他线程");

                        barrier.await();
                    }

                    System.out.println("线程 " + workerId + " 所有阶段完成");

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

            }, "worker-" + workerId).start();
        }
    }
}

这里同一个 barrier 会被反复使用。

第 1 阶段等 3 个线程到齐,放行。

第 2 阶段又等 3 个线程到齐,再放行。

第 3 阶段再等一次。

这就是 CyclicBarrier 可以循环使用的意思。


PDF 项目里会不会用到 CyclicBarrier

说实话,普通 PDF 批量加水印场景,一般不太需要 CyclicBarrier

因为 PDF 批量处理更常见的是:

每个 PDF 是一个独立任务;
处理完就结束;
主线程等所有任务完成;
最后统计结果。

这更适合:

CountDownLatch;
ThreadPoolExecutor;
CompletableFuture.allOf()。

CyclicBarrier 更适合这种“阶段同步”的任务。

比如:

多个线程先加载数据;
都加载完成后,再一起计算;
都计算完成后,再一起汇总;
都汇总完成后,再一起输出。

如果 PDF 项目里非要找一个类似场景,可能是这样:

多个线程先准备不同目录下的 PDF 清单;
所有清单都准备好后;
再统一开始处理文件。

但这种场景并不常见。

所以这一节重点不是说 PDF 项目一定要用 CyclicBarrier,而是把它和 CountDownLatch 的区别搞清楚。


一个更接近业务的例子

假设有三个线程,分别加载三类数据:

线程1:加载 PDF 文件列表;
线程2:加载水印配置;
线程3:加载用户任务信息。

这三类数据都准备好以后,才能进入下一步。

可以写成这样:

package com.succos.thread;

import java.util.concurrent.CyclicBarrier;

public class CyclicBarrierBusinessDemo {

    public static void main(String[] args) {

        CyclicBarrier barrier = new CyclicBarrier(3, () -> {
            System.out.println("文件列表、水印配置、任务信息都准备好了,可以进入下一步");
        });

        new Thread(() -> {
            loadData("PDF 文件列表", 1000, barrier);
        }, "file-loader").start();

        new Thread(() -> {
            loadData("水印配置", 2000, barrier);
        }, "config-loader").start();

        new Thread(() -> {
            loadData("任务信息", 3000, barrier);
        }, "task-loader").start();
    }

    private static void loadData(String name, long sleepMillis, CyclicBarrier barrier) {

        try {
            System.out.println("开始加载:" + name
                    + ",线程:"
                    + Thread.currentThread().getName());

            Thread.sleep(sleepMillis);

            System.out.println(name + " 加载完成,等待其他数据");

            barrier.await();

            System.out.println(name + " 所在线程继续执行下一步");

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

这个例子更能体现 CyclicBarrier 的味道:

不是主线程在等子线程;
而是几个工作线程之间互相等。

如果有线程异常会怎样

CyclicBarrier 有一个需要注意的地方:如果其中一个线程出问题,没有正常到达 await(),其他已经到达的线程就可能一直等。

比如 3 个线程约好到屏障集合。

结果线程3 在到达 await() 之前就抛异常退出了。

那线程1、线程2 可能已经在 await() 处等着,但永远等不到第 3 个线程。

所以在真实项目里,用 CyclicBarrier 要特别注意异常处理。

有时也会使用带超时时间的 await()

barrier.await(5, TimeUnit.SECONDS);

意思是最多等 5 秒。

如果等不到,就抛异常,不要无限等下去。

当然,这一节先不展开超时版本,先把基本概念理解清楚。


CyclicBarrier 适合什么场景

我现在会把 CyclicBarrier 归到“阶段同步”场景里。

比如:

多线程并行计算,每一轮计算完后再进入下一轮;
多个玩家加载完成后,游戏统一开始;
多个模块初始化完成后,系统进入下一阶段;
多个数据源都准备好以后,再一起执行后续流程。

它强调的是:

大家都到同一个点,再一起继续。

而不是简单地让主线程等结果。


和 Semaphore 再区分一下

前面刚学过 Semaphore,这里也顺便区分一下。

Semaphore 是控制数量:

同一时间最多允许几个线程进去。

比如:

最多 3 个 PDF 同时处理。

CyclicBarrier 是等待到齐:

必须等指定数量的线程都到达某个点,才能继续。

比如:

3 个线程都准备好后,再一起开始下一阶段。

一个是“限流”,一个是“集合”。

这两个不要混。


这一节小结

这一节我主要记住几点:

1. CyclicBarrier 可以让多个线程在某个位置互相等待;
2. barrier.await() 是集合点,谁先到谁等;
3. 等指定数量的线程都到达后,屏障打开,大家继续执行;
4. CyclicBarrier 可以重复使用;
5. CountDownLatch 更像主线程等一批任务完成;
6. CyclicBarrier 更像一批线程互相等,等人齐了再继续;
7. 普通 PDF 批量处理一般更常用 CountDownLatch 或线程池,不太需要 CyclicBarrier。

用一句话总结:

CyclicBarrier 解决的是“几个线程先在这里集合,人齐了再一起往后走”的问题。

下一节开始进入线程池。

前面一直在手动 new Thread,但真实项目里不可能一直这么写。线程池才是批量任务更常用的做法。

17. CyclicBarrier:让多个线程在某个位置互相等待 - Java并发之从Thread到CompletableFuture - 上下文网