返回文章列表
JUC并发编程
JUC设计模式GuardedObject顺序控制

15同步模式:保护性暂停与顺序控制

前面学习了 synchronized、wait/notify、LockSupport 等基础工具,本篇进入「共享模型之模式」,把这些工具组合成解决特定并发问题的套路。本篇先整理三个同步模式:

  • 保护性暂停(Guarded Suspension):一个线程等待另一个线程的执行结果
  • Balking(犹豫):别人(或自己)已经做过的事,就不用再做了
  • 顺序控制:固定运行顺序、多个线程交替输出

同步模式之保护性暂停

定义

保护性暂停即 Guarded Suspension,用在一个线程等待另一个线程的执行结果。

要点:

  • 有一个结果需要从一个线程传递到另一个线程,让它们关联同一个 GuardedObject
  • 如果有结果不断从一个线程传递到另一个线程,那么可以使用消息队列(见生产者/消费者模式)
  • JDK 中,join 的实现、Future 的实现,采用的就是此模式
  • 因为要等待另一方的结果,因此归类到同步模式

t1 在 GuardedObject 上等待 response 有值,t2 任务完成后为 response 赋值并通知 t1

基本实现

GuardedObject 内部用一个 response 成员变量保存结果:等待方在条件不满足时 wait,生产方算出结果后赋值并 notifyAll。

class GuardedObject {
    private Object response;
    private final Object lock = new Object();
 
    public Object get() {
        synchronized (lock) {
            // 条件不满足则等待
            while (response == null) {
                try {
                    lock.wait();
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
            }
            return response;
        }
    }
 
    public void complete(Object response) {
        synchronized (lock) {
            // 条件满足,通知等待线程
            this.response = response;
            lock.notifyAll();
        }
    }
}

一个线程等待另一个线程执行结果的示例:

public static void main(String[] args) {
    GuardedObject guardedObject = new GuardedObject();
    new Thread(() -> {
        try {
            // 子线程执行下载
            List<String> response = download();
            log.debug("download complete...");
            guardedObject.complete(response);
        } catch (IOException e) {
            e.printStackTrace();
        }
    }).start();
 
    log.debug("waiting...");
    // 主线程阻塞等待
    Object response = guardedObject.get();
    log.debug("get response: [{}] lines", ((List<String>) response).size());
}

运行结果:

08:42:18.568 [main] c.TestGuardedObject - waiting...
08:42:23.312 [Thread-0] c.TestGuardedObject - download complete...
08:42:23.312 [main] c.TestGuardedObject - get response: [3] lines

注意两个细节:

  • 判断条件必须用 while 而不是 if,防止虚假唤醒——线程被意外唤醒时,要重新检查 response == null 是否仍然成立
  • 唤醒要用 notifyAll,因为同一锁上等待的线程可能不止一个

带超时版 GuardedObject

如果不想无限期等待,而是要控制超时时间呢?思路是把「一次性等待」拆成「多轮等待」,每轮记录已经经历的时间,剩余时间作为下一轮 wait 的参数:

  • 记录开始等待的时间 begin 和已经历的时间 timePassed
  • 每轮要等待的时间为 millis - timePassed
  • 若剩余等待时间 waitTime <= 0,说明已经超时,退出循环
class GuardedObjectV2 {
    private Object response;
    private final Object lock = new Object();
 
    public Object get(long millis) {
        synchronized (lock) {
            // 1) 记录最初时间
            long begin = System.currentTimeMillis();
            // 2) 已经经历的时间
            long timePassed = 0;
            while (response == null) {
                // 4) 假设 millis 是 1000,结果在 400 时唤醒了,那么还有 600 要等
                long waitTime = millis - timePassed;
                log.debug("waitTime: {}", waitTime);
                if (waitTime <= 0) {
                    log.debug("break...");
                    break;
                }
                try {
                    lock.wait(waitTime);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                // 3) 如果提前被唤醒,这时已经经历的时间假设为 400
                timePassed = System.currentTimeMillis() - begin;
                log.debug("timePassed: {}, object is null {}",
                        timePassed, response == null);
            }
            return response;
        }
    }
 
    public void complete(Object response) {
        synchronized (lock) {
            // 条件满足,通知等待线程
            this.response = response;
            log.debug("notify...");
            lock.notifyAll();
        }
    }
}

测试:子线程先用 null 唤醒一次(模拟虚假唤醒),1 秒后再传入真正的结果,等待方在总时长 2500ms 内可以正确拿到结果:

public static void main(String[] args) {
    GuardedObjectV2 v2 = new GuardedObjectV2();
    new Thread(() -> {
        sleep(1);
        v2.complete(null);
        sleep(1);
        v2.complete(Arrays.asList("a", "b", "c"));
    }).start();
 
    Object response = v2.get(2500);
    if (response != null) {
        log.debug("get response: [{}] lines", ((List<String>) response).size());
    } else {
        log.debug("can't get response");
    }
}
08:49:39.917 [main] c.GuardedObjectV2 - waitTime: 2500
08:49:40.917 [Thread-0] c.GuardedObjectV2 - notify...
08:49:40.917 [main] c.GuardedObjectV2 - timePassed: 1003, object is null true
08:49:40.917 [main] c.GuardedObjectV2 - waitTime: 1497
08:49:41.918 [Thread-0] c.GuardedObjectV2 - notify...
08:49:41.918 [main] c.GuardedObjectV2 - timePassed: 2004, object is null false
08:49:41.918 [main] c.TestGuardedObjectV2 - get response: [3] lines

如果等待时间不足(超时),get 返回 null,等待方不会一直挂住:

08:47:54.963 [main] c.GuardedObjectV2 - waitTime: 1500
08:47:55.963 [Thread-0] c.GuardedObjectV2 - notify...
08:47:55.963 [main] c.GuardedObjectV2 - timePassed: 1002, object is null true
08:47:55.963 [main] c.GuardedObjectV2 - waitTime: 498
08:47:56.461 [main] c.GuardedObjectV2 - timePassed: 1500, object is null true
08:47:56.461 [main] c.GuardedObjectV2 - waitTime: 0
08:47:56.461 [main] c.GuardedObjectV2 - break...
08:47:56.461 [main] c.TestGuardedObjectV2 - can't get response

原理之 join:Thread.join() 的底层正是保护性暂停——它内部判断「被等待线程还活着」就执行 wait(0),而线程终止时 JVM 会自动对该线程对象调用 notifyAll,等待线程被唤醒后再次检查存活状态,直到目标线程结束。

多任务版 GuardedObject

图中 Futures 就好比居民楼一层的信箱(每个信箱有房间编号),左侧的 t0、t2、t4 就好比等待邮件的居民,右侧的 t1、t3、t5 就好比邮递员。

Futures 信箱:t0/t2/t4 凭 id 等待各自 GuardedObject 的结果,t1/t3/t5 按 id 投递结果

如果需要在多个类之间使用 GuardedObject,把它作为参数到处传递很不方便。因此设计一个用来解耦的中间类 Mailboxes,这样不仅能解耦【结果等待者】和【结果生产者】,还能同时支持多个任务的管理。

首先给 GuardedObject 增加唯一 id:

class GuardedObject {
 
    // 标识 Guarded Object
    private int id;
 
    public GuardedObject(int id) {
        this.id = id;
    }
 
    public int getId() {
        return id;
    }
 
    // 结果
    private Object response;
 
    // 获取结果
    // timeout 表示要等待多久,单位毫秒
    public Object get(long timeout) {
        synchronized (this) {
            // 开始时间
            long begin = System.currentTimeMillis();
            // 经历的时间
            long passedTime = 0;
            while (response == null) {
                // 这一轮循环应该等待的时间
                long waitTime = timeout - passedTime;
                // 经历的时间超过了最大等待时间时,退出循环
                if (timeout - passedTime <= 0) {
                    break;
                }
                try {
                    this.wait(waitTime); // 虚假唤醒时进入下一轮继续等
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                // 求得经历时间
                passedTime = System.currentTimeMillis() - begin;
            }
            return response;
        }
    }
 
    // 产生结果
    public void complete(Object response) {
        synchronized (this) {
            // 给结果成员变量赋值
            this.response = response;
            this.notifyAll();
        }
    }
}

中间解耦类 Mailboxes 用 Hashtable(线程安全)维护 id 到 GuardedObject 的映射:

class Mailboxes {
    private static Map<Integer, GuardedObject> boxes = new Hashtable<>();
 
    private static int id = 1;
 
    // 产生唯一 id
    private static synchronized int generateId() {
        return id++;
    }
 
    public static GuardedObject getGuardedObject(int id) {
        return boxes.remove(id);
    }
 
    public static GuardedObject createGuardedObject() {
        GuardedObject go = new GuardedObject(generateId());
        boxes.put(go.getId(), go);
        return go;
    }
 
    public static Set<Integer> getIds() {
        return boxes.keySet();
    }
}

getGuardedObject 取信时直接 remove,结果被消费后信箱就被回收。业务相关的两个角色:居民(等待结果)和邮递员(生产结果)。

class People extends Thread {
    @Override
    public void run() {
        // 收信
        GuardedObject guardedObject = Mailboxes.createGuardedObject();
        log.debug("开始收信 id:{}", guardedObject.getId());
        Object mail = guardedObject.get(5000);
        log.debug("收到信 id:{}, 内容:{}", guardedObject.getId(), mail);
    }
}
 
class Postman extends Thread {
    private int id;
    private String mail;
 
    public Postman(int id, String mail) {
        this.id = id;
        this.mail = mail;
    }
 
    @Override
    public void run() {
        GuardedObject guardedObject = Mailboxes.getGuardedObject(id);
        log.debug("送信 id:{}, 内容:{}", id, mail);
        guardedObject.complete(mail);
    }
}

测试:3 个居民先各自创建信箱等待,1 秒后邮递员按 id 逐个送信:

public static void main(String[] args) throws InterruptedException {
    for (int i = 0; i < 3; i++) {
        new People().start();
    }
    Sleeper.sleep(1);
    for (Integer id : Mailboxes.getIds()) {
        new Postman(id, "内容" + id).start();
    }
}

某次运行结果:

10:35:05.689 c.People [Thread-1] - 开始收信 id:3
10:35:05.689 c.People [Thread-2] - 开始收信 id:1
10:35:05.689 c.People [Thread-0] - 开始收信 id:2
10:35:06.688 c.Postman [Thread-4] - 送信 id:2, 内容:内容2
10:35:06.688 c.Postman [Thread-5] - 送信 id:1, 内容:内容1
10:35:06.688 c.People [Thread-0] - 收到信 id:2, 内容:内容2
10:35:06.688 c.People [Thread-2] - 收到信 id:1, 内容:内容1
10:35:06.688 c.Postman [Thread-3] - 送信 id:3, 内容:内容3
10:35:06.689 c.People [Thread-1] - 收到信 id:3, 内容:内容3

这其实就是 Future 模式的雏形:createGuardedObject 相当于提交任务拿到 Future,get 是等待结果,complete 是任务完成后回填结果。

同步模式之 Balking

定义

Balking(犹豫)模式用在:一个线程发现另一个线程或本线程已经做了某一件相同的事,那么本线程就无需再做了,直接结束返回。

实现

例如前端页面多次点击按钮调用 start,要保证监控线程只被启动一次:

public class MonitorService {
 
    // 用来表示是否已经有线程在执行启动了
    private volatile boolean starting;
 
    public void start() {
        log.info("尝试启动监控线程...");
        synchronized (this) {
            if (starting) {
                return;
            }
            starting = true;
        }
 
        // 真正启动监控线程...
    }
}

多次点击的输出中,只有第一次真正执行了启动:

[http-nio-8080-exec-1] c.MonitorService - 该监控线程已启动?(false)
[http-nio-8080-exec-1] c.MonitorService - 监控线程已启动...
[http-nio-8080-exec-2] c.MonitorService - 该监控线程已启动?(true)
[http-nio-8080-exec-3] c.MonitorService - 该监控线程已启动?(true)
[http-nio-8080-exec-4] c.MonitorService - 该监控线程已启动?(true)

它还经常用来实现线程安全的单例——发现实例已经存在就直接返回,不再创建:

public final class Singleton {
    private Singleton() {
    }
 
    private static Singleton INSTANCE = null;
 
    public static synchronized Singleton getInstance() {
        if (INSTANCE != null) {
            return INSTANCE;
        }
 
        INSTANCE = new Singleton();
        return INSTANCE;
    }
}

对比一下保护性暂停模式:保护性暂停用在一个线程等待另一个线程的执行结果,条件不满足时线程等待;Balking 模式条件不满足(事情已被做过)时线程直接返回,并不等待。

同步模式之顺序控制

固定运行顺序

比如要求必须先打印 2、后打印 1。

wait/notify 版

用一个「运行标记」记录 t2 是否执行过:t1 没等到标记就 wait,t2 打印完修改标记并 notifyAll。

// 用来同步的对象
static Object obj = new Object();
// t2 运行标记,代表 t2 是否执行过
static boolean t2runed = false;
 
public static void main(String[] args) {
    Thread t1 = new Thread(() -> {
        synchronized (obj) {
            // 如果 t2 没有执行过
            while (!t2runed) {
                try {
                    // t1 先等一会
                    obj.wait();
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
            }
        }
        System.out.println(1);
    });
 
    Thread t2 = new Thread(() -> {
        System.out.println(2);
        synchronized (obj) {
            // 修改运行标记
            t2runed = true;
            // 通知 obj 上等待的线程(可能有多个,因此需要用 notifyAll)
            obj.notifyAll();
        }
    });
 
    t1.start();
    t2.start();
}

可以看到,实现上很麻烦:

  • 首先,需要保证先 wait 再 notify,否则 wait 线程永远得不到唤醒,因此使用了『运行标记』来判断该不该 wait
  • 第二,如果有些干扰线程错误地 notify 了等待线程,条件不满足时还要重新等待,使用了 while 循环来解决
  • 最后,唤醒对象上的 wait 线程需要使用 notifyAll,因为『同步对象』上的等待线程可能不止一个

Park/Unpark 版

可以使用 LockSupport 类的 park 和 unpark 来简化上面的题目:

Thread t1 = new Thread(() -> {
    try { Thread.sleep(1000); } catch (InterruptedException e) { }
    // 当没有『许可』时,当前线程暂停运行;有『许可』时,用掉这个『许可』,当前线程恢复运行
    LockSupport.park();
    System.out.println("1");
});
 
Thread t2 = new Thread(() -> {
    System.out.println("2");
    // 给线程 t1 发放『许可』(多次连续调用 unpark 只会发放一个『许可』)
    LockSupport.unpark(t1);
});
 
t1.start();
t2.start();

park 和 unpark 方法比较灵活,它俩谁先调用、谁后调用无所谓,并且是以线程为单位进行『暂停』和『恢复』,不需要『同步对象』和『运行标记』。

交替输出

线程 1 输出 a 5 次,线程 2 输出 b 5 次,线程 3 输出 c 5 次,要求输出 abcabcabcabcabc,分别用三种手段实现。

wait/notify 版

用一个 flag(1/2/3)表示轮到哪个线程。每个线程在 print 中传入「自己的等待标记」和「下一个线程的标记」,打印后修改标记并唤醒全体等待者,没轮到的线程醒来后继续 wait。

class SyncWaitNotify {
    private int flag;
    private int loopNumber;
 
    public SyncWaitNotify(int flag, int loopNumber) {
        this.flag = flag;
        this.loopNumber = loopNumber;
    }
 
    public void print(int waitFlag, int nextFlag, String str) {
        for (int i = 0; i < loopNumber; i++) {
            synchronized (this) {
                while (this.flag != waitFlag) {
                    try {
                        this.wait();
                    } catch (InterruptedException e) {
                        e.printStackTrace();
                    }
                }
                System.out.print(str);
                flag = nextFlag;
                this.notifyAll();
            }
        }
    }
}
SyncWaitNotify syncWaitNotify = new SyncWaitNotify(1, 5);
new Thread(() -> {
    syncWaitNotify.print(1, 2, "a");
}).start();
new Thread(() -> {
    syncWaitNotify.print(2, 3, "b");
}).start();
new Thread(() -> {
    syncWaitNotify.print(3, 1, "c");
}).start();

Lock 条件变量版

ReentrantLock 可以创建多个 Condition,每个条件变量对应一个线程的「等待休息室」,这样可以精确唤醒下一个线程(signal),不必像 notifyAll 那样叫醒所有人。

class AwaitSignal extends ReentrantLock {
    public void start(Condition first) {
        this.lock();
        try {
            log.debug("start");
            first.signal();
        } finally {
            this.unlock();
        }
    }
 
    public void print(String str, Condition current, Condition next) {
        for (int i = 0; i < loopNumber; i++) {
            this.lock();
            try {
                current.await();
                log.debug(str);
                next.signal();
            } catch (InterruptedException e) {
                e.printStackTrace();
            } finally {
                this.unlock();
            }
        }
    }
 
    // 循环次数
    private int loopNumber;
 
    public AwaitSignal(int loopNumber) {
        this.loopNumber = loopNumber;
    }
}
AwaitSignal as = new AwaitSignal(5);
Condition aWaitSet = as.newCondition();
Condition bWaitSet = as.newCondition();
Condition cWaitSet = as.newCondition();
 
new Thread(() -> {
    as.print("a", aWaitSet, bWaitSet);
}).start();
new Thread(() -> {
    as.print("b", bWaitSet, cWaitSet);
}).start();
new Thread(() -> {
    as.print("c", cWaitSet, aWaitSet);
}).start();
 
as.start(aWaitSet);

注意:该实现没有考虑 a、b、c 三个线程都就绪后再开始,主线程调用 start 发令时如果某个线程还没执行到 await,可能错过信号。实际使用时需要额外的就绪同步。

Park/Unpark 版

每个线程打印前先 park 等许可,打印后 unpark 数组中的下一个线程;最后一个线程的下一个回到第一个,形成循环:

class SyncPark {
    private int loopNumber;
    private Thread[] threads;
 
    public SyncPark(int loopNumber) {
        this.loopNumber = loopNumber;
    }
 
    public void setThreads(Thread... threads) {
        this.threads = threads;
    }
 
    public void print(String str) {
        for (int i = 0; i < loopNumber; i++) {
            LockSupport.park();
            System.out.print(str);
            LockSupport.unpark(nextThread());
        }
    }
 
    private Thread nextThread() {
        Thread current = Thread.currentThread();
        int index = 0;
        for (int i = 0; i < threads.length; i++) {
            if (threads[i] == current) {
                index = i;
                break;
            }
        }
        if (index < threads.length - 1) {
            return threads[index + 1];
        } else {
            return threads[0];
        }
    }
 
    public void start() {
        for (Thread thread : threads) {
            thread.start();
        }
        LockSupport.unpark(threads[0]);
    }
}
SyncPark syncPark = new SyncPark(5);
Thread t1 = new Thread(() -> {
    syncPark.print("a");
});
Thread t2 = new Thread(() -> {
    syncPark.print("b");
});
Thread t3 = new Thread(() -> {
    syncPark.print("c\n");
});
syncPark.setThreads(t1, t2, t3);
syncPark.start();

小结

  • 保护性暂停用于一个线程等待另一个线程的结果,核心是 while 判条件 + wait/notifyAll;带超时版通过「开始时间 - 已历时间」拆成多轮等待,同时处理了虚假唤醒与超时边界
  • 多任务版用 Mailboxes 在 id 与 GuardedObject 之间做映射,解耦结果等待者与生产者,这就是 Future/FutureTask 的思想
  • Balking 用于「只做一次」的场景:发现事情已被做过就直接返回,典型应用是只初始化一次、线程安全单例
  • 固定顺序控制中,wait/notify 需要锁对象与运行标记,park/unpark 以线程为单位、先后次序无所谓,更简洁
  • 交替输出三种实现:wait/notify 靠标记位 + 全体唤醒;Lock 条件变量靠多个 Condition 精确唤醒;park/unpark 靠线程间传递许可