前言

众所周知,Java 招聘时面试官刀人三板斧:

第一板斧:"你说你在项目中用了xxx,具体怎么用的呢?" —— 这波尚且扛得住;

第二板斧:"那你知道xxx的底层原理吗?/ xxx底层怎么实现的呢?" —— 这波已经有点汗流浃背了,有所准备的可能答得出来,没准备过的就出门右拐了,悲。

但是,还有第三板斧:"那xxx为什么要这么设计呢?直接xx不就可以吗?"

—— 然后苦命的 Java 选手,就直接 g 了。

为了防止再次被问穿导致出门右拐,笔者今天开始开个坑——手写 JUC 系列。这个系列打算从一些比较简洁的 JUC 类开始一步步深入,且主要聚焦于 JUC 类的一些核心功能的手写复现,而非所有功能(比如支持可打断等可能就不会手写)。

手写的目的主要是为了加深笔者自己,以及各位读者对 JUC 设计哲学的了解。从0开始自己实现某个 JUC 的类,你首先想到的是怎么实现?为什么 JUC 不像我们想的那样来实现?笔者认为,只有经过手写的一番对比,我们才能真正从容地应对第二和第三板斧,此所谓知其然,更知其所以然。

最后需要声明的是,笔者本人也对 JUC 的认识也绝对称不上深入,如有谬误也烦请大家指出。

题外话就先说到这里。今天我们先看看 FutureTask —— 阻塞等待异步执行结果的 JUC 类。

一、FutureTask 用法回顾

首先回顾下 FutureTask 的基本用法或者说核心特性:

  1. FutureTask 间接实现了 Runnable 接口,可以作为 Thread 构造函数的参数传入,相当于启动一个线程去执行某个任务。
  2. FutureTask 间接实现了 Future 借口,支持当前线程阻塞等待任务完成(死等/超时等待)、查看任务是否完成、取消任务等。
  3. FutureTask 的构造函数支持包装 Callable 任务,允许线程执行一个有返回值的任务,且能通过 FutureTask 对象拿到任务的返回结果。
  4. 若 FutureTask 包装的 Callable 任务执行异常,将会把该异常传播至阻塞等待结果的线程。

多说无益,直接上一个简单的例子来回顾 FutureTask 的用法:

@Slf4j(topic = "JucTest")
@SpringBootTest
class JucTestApplicationTests {

    public static void main(String[] args) {
        FutureTask<String> futureTask = new FutureTask<>(() -> {
            log.info("任务线程启动,执行任务中...");
            Thread.sleep(2000);
            log.info("任务执行完成...");
            return "hello, FutureTask!";
        });
        new Thread(futureTask, "任务线程").start();


        int waitSeconds = 3;
        new Thread(()->{
            waitFutureTask(futureTask, waitSeconds);
        }, "A").start();

        waitFutureTask(futureTask, waitSeconds);
    }

    static void waitFutureTask(FutureTask<String> futureTask, int waitSeconds) {
        log.info("线程{}等待 futureTask 完成...",  Thread.currentThread().getName());
        long start = System.currentTimeMillis();
        try {
            String futureTaskResult = futureTask.get(waitSeconds, TimeUnit.SECONDS);
            log.info("线程{}等待{} ms 后,任务正常完成,返回值:{}",Thread.currentThread().getName(),
                     System.currentTimeMillis() - start, futureTaskResult);
        } catch (ExecutionException e) {
            log.error("线程{}等待{} ms 后,futureTask 任务内部执行错误,错误信息:{}",Thread.currentThread().getName(),
                      System.currentTimeMillis() - start, e.getMessage());
        } catch (TimeoutException e) {
            log.warn("线程{}等待{} ms 后, futureTask 仍未完成,任务超时", Thread.currentThread().getName(),
                     System.currentTimeMillis() - start);
        } catch (Exception e) {
            log.info("其他情况");
        }
    }

}

简单介绍下这段代码的结构:

任务线程负责执行 FutureTask,并 sleep 2 秒模拟任务耗时;主线程和线程 A 阻塞等待 FutureTask 完成,waitSeconds 为 FutureTask.get 的最大等待时间。

最终输出日志如下:

20:41:37.850  INFO  JucTest [任务线程] - 任务线程启动,执行任务中...
20:41:37.850  INFO  JucTest [A   ] - 线程A等待 futureTask 完成...
20:41:37.850  INFO  JucTest [main] - 线程main等待 futureTask 完成...
20:41:39.863  INFO  JucTest [任务线程] - 任务执行完成...
20:41:39.864  INFO  JucTest [main] - 线程main等待2012 ms 后,任务正常完成,返回值:hello, FutureTask!
20:41:39.864  INFO  JucTest [A   ] - 线程A等待2012 ms 后,任务正常完成,返回值:hello, FutureTask!

同时,在代码的 waitFutureTask 方法中也可以看到,futureTask.get() 这个阻塞等待任务结果的方法抛出多个异常,包括 TimeoutException 和 ExecutionException;

在上面的例子中,我们设置了 futureTask.get() 的时间为3秒,如果这个时间设置为小于2秒,则main 线程和线程A会在任务完成前超时,就会捕获 TimeoutException;如果任务线程执行 FutureTask 时内部异常,就会捕获 ExecutionException。

这也是我们希望手写实现的功能,即:手写 get 方法,get 方法包含以下几个功能:

1.若任务在指定时间内正常完成,调用 future.get 的线程能立即停止阻塞并得到任务返回值;

2.若任务执行耗时超过最大等待时间,调用 future.get 的线程能在达到超时时间后停止阻塞,get 方法抛出 TimeoutException 超时异常;

下面我们来分析并实现这些功能。

二、实现思路分析与图解

2.1. FutureTask 如何实现一个带返回值的任务?

我们知道,线程的 run 方法是没有返回值的。那 start 一个 Thread 去执行 FutureTask ,为什么其他线程能获取任务结果呢?

  • 首先,阻塞等待任务完成的多个线程,都是通过同一个 FutureTask 引用来获取返回值,本质上是多个线程读取一个共享的对象。那我们的思路就很清晰了:

    ——在 FutureTask 类中放置一个成员变量 result 用来存储任务的返回结果,这样其他线程只要持有 FutureTask 对象的引用,不就能看到这个结果了吗。

  • 其次,观察 Thread 类 run 方法的逻辑,线程在启动后调用 Thread 的 run 方法,该方法实际上是委托传入的 Runnable 对象调用 run 方法,也就是让 FutureTask 对象调用 run 方法:

    ——因此,在 FutureTask 类的 run 方法中,让封装的 Callable 实例调用 call 方法,再将结果存到成员变量 result 中,这样其他的线程就可以通过 FutureTask 的引用,获取到 result

public class Thread implements Runnable {

    private Runnable target;
    
    @Override
    public void run() {
        if (target != null) {
            target.run();
        }
    }
    
}

回顾 Thread 的构造

2.2.FutureTask 调用 get 如何实现阻塞等待至任务完成?

有了上面的思路,我们可以这样去做:通过 FutureTask 的引用查看成员变量 result 是否为空:

如果 result 不为空,证明任务已完成,无需等待直接返回;

如果 result 为空,开始等待;对于等待的情况,考虑以下两个问题:

  • 什么时候、通过什么方式,让线程结束等待?

    ——既然 FutureTask 的功能是任务完成后,调用 get 的线程停止等待,那么我们完全可以在任务线程将返回结果存入 result 后,新增一步逻辑,唤醒这些等待的线程。

  • 任务线程如何知道哪些线程在等待并需要唤醒呢?

    ——这就需要我们在 FutureTask 对象中维护一个集合,在判断 result 为空,让线程阻塞等待前,将其 Thread 对象存入这个集合,任务线程不就能从这个集合里遍历唤醒线程了吗?

由此,我们已经掌握了 FutureTask 的基本实现思路,总结一下:

  1. FutureTask 中维护一个 result 变量存放 Callable 返回值,这样其他线程只要持有 FutureTask 引用就可以看到结果。
  2. 如果其他线程调用 FutureTask.get,发现结果为空,则将自身 Thread 放入等待队列后进入阻塞;
  3. 任务线程完成 Callable 调用后,将返回值存入 result ,并唤醒等待队列中的线程。

根据目前我们的理解,任务线程生命周期可视化流程图如下:

三、手写实现 FutureTask 的几种方式与优缺点

思路已经明确了,如何实现呢?其实上图的流程并不复杂,主要的实现点在于:

  • 选择什么集合来存放等待的线程的 Thread 对象;
  • 怎么处理线程的等待和唤醒。

我们先考虑第二个实现点:Java 中有哪些方式实现线程的等待和唤醒?可能你会想到一系列的 JUC 类,比如 ReentrantLock Conditon 的 await 和 signal,CountDownLatch 的 await 和 countDown 等。但本质上这些等待和唤醒的机制都可以按下面两种方式分类:

native 方法和 非 native 方法:

  1. native 方法只有两种:Object 类的 wait/notify,和 Unsafe 类的 park/unpark ( LockSupport 的 park/unpark 实质就是 Unsafe 类的 park/unpark)
  2. 非 native 方法:都是对 native 方法的封装。比如各种 JUC 类的同步工具,本质上还是对 park/unpark 的封装;Thread 类的 join,本质上是对 wait/notify 的封装。

依赖独占锁和无锁:

独占锁和阻塞/唤醒的概念不一样。独占锁是同步机制,用于保护临界区资源,而阻塞和唤醒是针对线程状态转换的操作。独占锁的实现依赖线程的阻塞和唤醒,但阻塞和唤醒不依赖独占锁。所以阻塞、唤醒线程的方法也可以根据是否依赖独占锁来划分:(下面依赖独占锁简称依赖锁)

  1. 依赖锁:wait/notify 依赖于同步块 synchronized,Condition 的 await/signal 依赖于 lock
  2. 不依赖锁:park/unpark,以及其他无锁的 JUC 类。

我们从依赖锁和不依赖锁两种方式来实现 FutureTask。

3.1依赖锁:synchronized wait/notify 或 Condition await/signal

我们尽量避免用一个 JUC 类去实现另一个 JUC 类。所以这里我就用 wait/notify 实现。我们回顾一下主要的实现点:

  • 选择什么集合来存放等待的线程的 Thread 对象;
  • 怎么处理线程的等待和唤醒。

再联想 wait / notify 的功能,是否一目了然了呢?我们甚至无需自己去创建存放 Thread 的集合,只需要调用锁对象的 wait/notify 方法,线程等待集合的维护,等待 / 唤醒的处理,就都由 wait/notify 帮我们处理好了( synchronized 内部利用锁对象的 monitor 的 waitSet 维护等待线程集合)。

代码如下:

@Slf4j
public class MyFutureTask3 <V> implements Runnable, Future {

    @Getter
    private volatile V result;

    private final Callable<V> callable;

    public MyFutureTask3(Callable<V> callable) {
        this.callable = callable;
    }

    @Override
    public void run() {
        try {
            this.result  = callable.call();
        } catch (Exception e) {
            throw new RuntimeException(e);
        }
        // 任务完成后唤醒等待线程
        synchronized (this) {
            this.notifyAll();
        }

    }

    // get 和 run 方法以外的方法暂不手写
    @Override
    public boolean cancel(boolean mayInterruptIfRunning) {
        return false;
    }

    @Override
    public boolean isCancelled() {
        return false;
    }

    @Override
    public boolean isDone() {
        return false;
    }

    @Override
    public Object get() throws InterruptedException, ExecutionException {
        if (result != null) {
            return result;
        }
        // 如果任务未完成,Thread 串行入队并阻塞,等待唤醒
        synchronized (this) {
            if (result != null) {
                return result;
            }
            this.wait();
            return result;
        }
    }

    @Override
    public Object get(long timeout, TimeUnit unit) throws InterruptedException, TimeoutException, ExecutionException {
        if (result != null) {
            return result;
        }
        synchronized (this) {
            if (result != null) {
                return result;
            }
            this.wait(unit.toMillis(timeout));
            // 判断是任务完成正常唤醒还是超时
            if (result != null) {
                return result;
            }
            // 如果 result 为空,证明任务还未完成,wait是超时返回(此处暂不处理虚假唤醒,后面会讲)
            throw new TimeoutException();
        }
    }
  }

以上的实现看着好像挺完美的?synchronized 也能保证串行入队/唤醒,防止并发安全问题。

可是坏就坏在这个 synchronized 上。虽然代码逻辑的确一点问题也没有,但是在并发场景下性能却低了。在多个线程都阻塞等待 FutureTask 结果时,我们且不谈多线程同时执行 get 导致的锁竞争问题,就说他们在被唤醒时这一段代码吧:

下面的例子中,线程 A、B、C 都在执行到 this.wait 时阻塞在此处,在被任务线程 notifyAll 后,尝试竞争 this 对象锁继续执行 return result 的代码:

可以看到,多个等待线程被唤醒后,本应该可以用无阻塞的方式读取 result 并返回;

——可是因为我们使用了 wait / notify 这类依赖锁的方式,所以某个等待线程被唤醒后,如果竞争对象锁失败,则可能又要陷入阻塞。一来二去的,开销可就大了。

所以不管是 synchronized 块中的 wait / notify,还是 ReentrantLock lock / unlock 之间代码 Condition 的 await / signal ,由于都依赖独占锁,就会导致多线程等待 FutureTask 时的性能问题。

那么我们能不能开发更高效的 FutureTask 呢?我们希望避免锁竞争导致的阻塞,直接无锁。——当然是可以的,这就是我们上面说的另外一种分类,用 park / unpark 这个无锁的线程协作方式来实现。

但无锁的设计对并发安全提出了更高的挑战,这也没办法,哪有两全其美的事?我们先来看看无锁 FutureTask 的基础版怎么实现。

3.2不依赖锁:park / unpark 处理线程阻塞唤醒 + 无锁数据结构存放等待线程

同样是思考这两点的实现:

  • 选择什么集合来存放等待的线程的 Thread 对象;
  • 怎么处理线程的等待和唤醒。

其中,第二点已经很明确了,用 park 和 unpark 挂起 / 唤醒线程;而对于第一点,在有锁时我们可以借助 synchronized monitor 自带的 waitSet 来帮我们维护线程的等待集合,而不依赖锁时,我们需要自己定义集合存放等待的线程。

可这个集合就比较难以选择了。也许你会说,就选一个 ArrayList<Thread> 来存放嘛,线程在进入阻塞等待前,先把 Thread 丢到这个 ArrayList 中;如果考虑到多线程同时调用 get 的场景,选用 Vector 这个线程安全集合不就好了吗?这个方法看着挺美好的,我们按着这个思路来实现下,看看有什么问题。

@Slf4j
public class MyFutureTask2<V> implements Runnable, Future {

    @Getter
    private volatile V result;

    private final Callable<V> callable;

    // 存放等待线程 Thread 对象的集合
    private final Vector<Thread> waitThreads = new Vector<>();

    public MyFutureTask2(Callable<V> callable) {
        this.callable = callable;
    }

    @Override
    public void run() {
        try {
            this.result  = callable.call();
            // 任务线程完成后,遍历存放等待线程的集合,一个个唤醒。
            for (Thread waitThread : waitThreads) {
                LockSupport.unpark(waitThread);
            }
        } catch (Exception e) {
            throw new RuntimeException(e);
        }
    }

    @Override
    public boolean cancel(boolean mayInterruptIfRunning) {
        return false;
    }

    @Override
    public boolean isCancelled() {
        return false;
    }

    @Override
    public boolean isDone() {
        return false;
    }

    @Override
    public Object get() throws InterruptedException, ExecutionException {
        // 任务未完成时,线程加入Vector 并 park,等待任务完成时唤醒。
        // 被唤醒后,将 Thread 从集合中移除。
        if (result == null) {
            waitThreads.add(Thread.currentThread());
            LockSupport.park();
            waitThreads.remove(Thread.currentThread());
            return result;
        }
        return result;
    }

    @Override
    public Object get(long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException {
        long waitNanos = unit.toNanos(timeout);
        if (result == null) {
            waitThreads.add(Thread.currentThread());
            LockSupport.parkNanos(waitNanos);
            waitThreads.remove(Thread.currentThread());
            
            // 判断 parkNanos 是超时导致的返回还是正常唤醒导致的返回
            // 如果是超时返回则抛超时异常。
            if (result == null) {
                throw new TimeoutException();
            }
            return result;
        }
        return result;
    }
}

看着挺合理的吧?——在"正常"情况下,确实没问题。可是放到并发场景下,这段实现却可能错误百出。我们来剖析下这段代码问题出在哪,无锁的实现要考虑哪些问题:

  1. 由于代码未被锁保护,可能出现调用 futureTask.get 的线程 park 阻塞后,没有线程来 unpark,导致调用 futureTask.get 的线程永久阻塞;
  2. 线程可能在 parkNanos 返回后仍被任务线程调用 unpark,这将导致线程下次调用 某个 FutureTask 的 get 方法时,出现调用 park,不阻塞、直接返回的情况;
  3. 使用 Vector 将重新引入锁竞争的性能问题,且对比 ArrayList ,仍无法保证线程安全,任务线程遍历 Vector 进行 unpark 时可能抛出 ConcurrentModificationException;

我们来一个个看这些问题,以及在不使用锁的场景下,是怎么解决这些问题的。

问题一:永久阻塞问题

看上面的代码,你可能会疑惑:我在 park 前不是已经判断 result == null 了么,这说明任务线程还未完成任务呀,既然如此,我 park 之后应该预期总是能被后续的任务线程 unpark 才对。可这就是无锁编程的麻烦之处:

在判断 result == null 时,任务的确未完成,可这并不代表,在执行到 LockSupport.park 时,任务还是未完成。

这就可能导致,在判断 result == null 之后,ThreadA 加入等待队列之前,如果在这个空隙内,任务线程完成任务并将集合中的所有线程(此时集合中线程不包括 ThreadA) unpark 完毕:那么 ThreadA 的 LockSupport.park 将无人唤醒并永久阻塞。示意图如下所示:

这个问题的根因在于:在无锁情况下调用 park,但没有保证 Thread 对象加入集合的操作发生在任务线程遍历集合之前;

比如上图所示的例子中,如果我们能够保证 waitThreads.add 发生在 任务线程唤醒等待队列之前,那么,即使线程A的 park 操作发生在任务线程唤醒等待队列之后,我们仍能保证此时线程A已被 unpark,不会出现永久阻塞的问题。

示意如下:

那么我们如何保证 waitThreads.add 发生在 任务线程唤醒等待队列之前呢?

答案是:—— Double Check。

在 waitThreads.add 将 Thread 加入等待队列后,再判断一次 result == null:

  • 如果 result 仍为 null,则能够保证 waitThreads.add 发生在任务线程唤醒等待队列之前,此时,线程A可以放心 park,因为能够确保有线程会来唤醒它;
  • 如果 result 不为 null,则代表 Callable 任务已经完成,则无需调用 park,直接返回 result;

再考虑上面永久阻塞的情况,通过 Double Check ,完美预防了这个问题,示意图如下:

从这个问题中,我们可以认识到 JUC 设计的某个哲学:

无锁编程中,调用 park 阻塞线程前,慎之又慎地检查条件,确保一定有其他线程会 unpark 它。

问题二:parkNanos 返回后仍被任务线程调用 unpark

我们知道,park / unpark 是基于许可证机制的。如果 unpark(threadA) 先于 threadA 调用 park,则 threadA 的 park 不会阻塞,而是消耗许可证后直接返回。

那么,如果在上述的代码中,出现了这种情况:

那么下次在进行另一个 futureTask.get 的方法时,其内部的 park 操作将会直接返回。

如何解决?也许你会说,我们可以在任务线程 unpark 等待队列的时候,判断一下线程的状态,如果线程是 Waiting 状态,才执行 unpark,否则不执行 unpark;我们思考下,这样是否可行?

—— 答案是不可行。park / unpark 的设计,本身就应该能允许 unpark 一个 Runnable 的线程,然后后续该线程 park 就可以直接返回。我们不能因为上述的这种情况,就去改变 unpark 的逻辑,否则一些先 unpark 再 park 的情况,反而因为新增的这个判断而无法处理了。

以及,LockSupport 类的作者同样有注明,park 可能在一些情况下,即便没有被 unpark,也会直接返回;也就是说,即便我们处理了上述这种 parkNanos 返回后仍被 unpark 的情况,也无法 100% 防止下次的 park 操作直接返回。

那究竟怎么做?

答案是:总是使用 while(!condition) {LockSupport.park;} 的代码结构,在 park 返回后检查一个能判断任务是否完成的"条件"(condition),如果条件不满足则再次 park,直到条件满足才退出。

这代表什么呢?这代表我们在编写 get 方法时,不需要去关心每一次 parkNanos 返回后是否刚好又被另一个线程 unpark ,而是在每一次调用 parkNanos 后,反复检查一个退出条件 condition,如果条件不满足则再次 park。

说检查条件可能有点抽象,什么是这里的"条件"?

  • 对于 park,我们希望的是"等待到任务完成为止",所以条件就是"任务完成",即 result != null;
  • 对于 parkNanos,我们希望的是"等待到任务完成或者超时",所以条件就是"任务完成或达到超时时间",即 result != null || System.currentTime > startWaitingTime + maxWaitingTime

由此你可能已经明白了:在上述的特殊情况中,如果出现线程A parkNanos 超时返回后仍被任务线程 unpark 的情况,这的确会导致下次线程A调用其他的 FutureTask.get 时,在 park 操作处直接返回;但是,返回后,线程A会去检查条件,result 是否仍然为空?超时时间是否未达到?如果是,则代表退出条件尚未满足,那么线程A会再次 调用 park;这就完美避免了这个问题。

这也是处理"虚假唤醒问题"的经典方式。总结一下:JUC 关注的不是"unpark线程如何处理 park 线程的意外返回",而是"park 线程如何处理自身的意外返回";park 线程处理自身意外返回的方式是,在 while 循环中检查条件,不满足则继续 park。

该方式对于处理 wait / notifyAll 的虚假唤醒同样适用;代码有以下两种写法,当 contition 比较复杂时,一般选第二种:

// 条件不满足则再次 park,条件满足则退出循环
while(!condition){
    LockSupport.park();
}

// 死循环,条件满足时break
while(true){
    if(condition){
        break;
    }
    LockSupport.park();
}

问题三:Vector 无法保证遍历安全且有性能问题

首先需要明确的是,如果一个线程(任务线程)在遍历 Vector 的过程中,另一个线程向 Vector 加入或删除元素,将会导致 ConcurrentModificationException。

可这是为什么?Vector 不是线程安全吗?这个以后我会单出一篇博客来讲,简单理解就是,Vector 迭代器的 next 方法中,会检查 Vector 集合是否在该迭代器创建后发生过 add / remove,如果有,next 方法就抛出异常;

我们利用迭代器遍历 Vector 时,虽然 next 方法的确用 synchronized 修饰,每次 next 操作都有 synchronized 保护,但多个 next 操作的组合,即遍历操作整体,是无法保障整体的原子性的。这就会导致,如果在多个 next 操作间穿插其他线程对 Vector 的 add 或 remove,则下一次 next 会检测到迭代器创建后有 add / remove,而抛出异常。

同时,Vector 的入队、出队操作,即 add 和 remove,由于用 synchronized 修饰,也同样有竞争情况下线程阻塞的性能问题。

那怎么改进呢?首先,要确保性能,我们得使用一个无锁的数据结构;其次,无锁数据结构又不能选择 ArrayList 这种线程不安全的。那又要线程安全又要无锁怎么办?

其实,看到安全 + 无锁这两个词,答案就应该呼之欲出:—— CAS。我们可以自定义一个使用 CAS 进行入队的数据结构;

同时,考虑到多线程的入队,插入操作多,存放等待线程的集合应该使用链表而非数组。

——综上所述,选用自定义链表 + CAS 入队来实现一个高性能的无锁数据结构,用于存放等待线程的 Thread 对象。

下面我们开始构想这个 基于 CAS 的链表怎么实现。

说起来简单实现起来难。我们知道,CAS 是针对于单个变量进行 CAS;可我们的入队操作是将线程节点加入链表,链表又不是单个变量,而是多个节点的组合,这怎么 CAS 呢?

但我们想想,链表真的不是变量吗。标识一个链表的是什么?——在写算法时都见过,是头结点。我们遍历链表,总是从头结点开始遍历,确定了头结点的引用就确定了链表。

那么基于这个思想,我们可以想到,头结点的引用标识链表,可以对头结点变量进行 CAS。而 既然 你 CAS 头结点,至少说明头结点应该有变化。那么是不是可以这么做:

线程入队时,线程节点利用头插法加入链表,此时原链表和新链表的头结点就不同了;然后,将链表的头结点引用 CAS 替换为新链表的头结点

我们首先定义链表节点的结构:

class WaitNode{

    private volatile Thread thread;
    private volatile WaitNode next;

    public WaitNode(Thread thread) {
        this.thread = thread;
    }
}

接着,CAS 入队流程如下:

  1. 将线程包装为一个Node,记为 Node current;
  2. 头插法入队得到新的链表:即读取当前的链表头结点 Node head,将当前节点 Node current 的 next 指向链表头结点 head;

     2.5. 理解为,旧的链表就是 head,新的链表就是 current;

     但此时链表的头结点变量 head 还没有真正变过来,需要进行 CAS。

     3.执行 compareAndSet(head变量,旧 head 值, current)替换链表

     如果 CAS 失败则重复2、3步再次 CAS 直到成功。

上述的这个流程,解决了多线程并发入队的竞争问题;那还有第二个问题:如何解决任务线程遍历链表与其他线程入队的竞争问题呢?

答案是:理论上其实可以不用管,任务线程直接读取当前链表头结点,往后遍历 unpark 就可以。

这是由于链表节点都是头插法入队的,就算任务线程读取当前链表头结点后,又有其他线程在头结点前插入了新节点,也不影响任务线程往后遍历的安全性。

那你可能要问了,照这种做法,那这些其他线程已经入队了,但是不会被 unpark 呀。其实这个问题也不必担心:——我们前面已经提到过,在线程入队之后,执行 park 之前,可以进行一次 result 的 Double Check;那么,如果在任务线程读取链表头结点之后,有其他线程来入队,那么虽然这些线程的确不会被 unpark,但此时任务显然已经完成,result 已设置,这时线程入队后进行 Double Check 就会判断任务已完成并直接返回。

示意如下:

也许你还要问,在 parkNanos 超时返回后,不是要移除链表中的该线程节点吗,这也还是会导致和任务线程竞争呀?—— 但其实 parkNanos 超时返回后也可以不移除节点,允许任务线程 unpark 它,而依赖 while(!condition){ LockSupport.park(); } 来保证后续 park 的逻辑正确性。

总结无锁 FutureTask get 方法实现要点及终版代码

好了,解决上面三个问题之后,我们总结一下编写无锁 FutureTask 的 get 方法的注意要点:

——判断 result 是否为 null,如果为 null(代表任务未完成),则接下来的流程可以分为入队、park 前检查、park、park 返回后检查这四个部分

  1. 入队:CAS 头插法入队链表,如果 CAS 失败则重复尝试
  2. park 前检查:CAS 入队成功后,不立即 park,而是再次检查 result 或超时条件;
  3. park:只有再次检查发现任务仍未完成且仍未超时,才让线程 park
  4. park 后检查:检查条件是否满足,如果任务仍未完成且仍未超时,则再次 park,直到满足条件。

这个代码怎么写比较优雅一些?我们想到,既然第二步和第四步都要检查条件,同时 CAS 也可能要重复尝试,其实 1-4 步都可以合并到一个 while 循环中写。同时,条件是要反复检查的,但入队操作只执行一次,所以我们可以维护一个是否入队的变量,在 while 循环中判断入队操作是否已执行,如果是,则不再重复执行。

以超时等待为例,代码如下:

@Slf4j
public class MyFutureTask4<V> implements Runnable, Future {

    /**
     * 维护等待链表的头部。
     */
    private volatile WaitNode waiterHead;

    @Getter
    private volatile V result;

    private final Callable<V> callable;

    // UNSAFE 类用于执行 CAS; Java 原子类执行 CAS,底层就是 unsafe 类完成的
    private static final sun.misc.Unsafe UNSAFE;

    // UNSAFE 类通过对象字段偏移量来定位 CAS 的字段
    private static final long waiterHeadOffset;

    // 只负责初始化 UNSAFE,没其他用
    static {
        try {
            Field theUnsafe = Unsafe.class.getDeclaredField("theUnsafe");
            theUnsafe.setAccessible(true);
            UNSAFE = (Unsafe) theUnsafe.get(null);
            waiterHeadOffset = UNSAFE.objectFieldOffset
                    (MyFutureTask4.class.getDeclaredField("waiterHead"));
        } catch (Exception e) {
            throw new Error(e);
        }
    }

    public MyFutureTask4(Callable<V> callable) {
        this.callable = callable;
    }

    // 等待线程节点,作为内部类
    class WaitNode{
        private volatile Thread thread;
        private volatile WaitNode next;

        public WaitNode(Thread thread) {
            this.thread = thread;
        }
    }

    @Override
    public Object get(long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException {
        // deadline 为超时时刻对应的时间戳
        long waitNanos = unit.toNanos(timeout);
        final long deadline = System.nanoTime() + waitNanos;
        boolean queued = false;
        
        while (true) {
            if (result != null) {
                return result;
            }
            if (!queued) {
                // 1.入队和CAS:读取当前头结点值,并用头插法加入链表
                // 执行 CAS 替换链表,如果失败则continue,在下次循环再次尝试
                // compareAndSwapObject 的那一步就是将waiterHead变量由 head CAS 为 current
                WaitNode current = new WaitNode(Thread.currentThread());
                WaitNode head = this.waiterHead;
                current.next = head;
                boolean isCASSuccess = UNSAFE.compareAndSwapObject(this, waiterHeadOffset, head, current);
                if (!isCASSuccess) {
                    continue;
                }
                // 2.如果 CAS 成功,进入第二步:park前检查
                // 标记入队变量为 true,之后的循环不会执行入队代码块
                // continue 进入下一轮循环,进行Double Check
                queued = true;
                continue;
            }
            
            long remainTime =  deadline - System.nanoTime();
            if (remainTime <= 0L) {
                throw new TimeoutException();
            } else {
                // 3.park前检查完成,进行 park 操作,remainTime 为剩余等待时间
                LockSupport.parkNanos(remainTime);
                // 4.park操作返回,进入下一轮循环,进行 park 后检查
            }
        }

    }

我们来模拟一个过程:某个线程等待 FutureTask 完成,CAS 尝试两次才成功入队,且 park 操作意外返回,利用这段代码能否保证逻辑正常。

整个执行流程如下所示:

可以看到:重复尝试 CAS 入队、入队完成后 Double Check 再 park、park 意外返回后检查条件发现不满足,再次 park 直到满足条件,在这个例子中都有体现。

任务线程的代码则比较简单:

  @Override
  public void run() {
      try {
          this.result  = callable.call();
          // 任务完成;读取当前的 head 节点直接开始遍历 unpark
          WaitNode node = this.waiterHead;
          while (node != null) {
              LockSupport.unpark(node.thread);
              node = node.next;
          }
      } catch (Exception e) {
          throw new RuntimeException(e);
      }
  }

至此,我们终于写完了这个不依赖锁的 FutureTask。但这也只是一个简陋版,主要在于思路的理清,比较耗费时间。我们接着看看 JUC FutureTask 源码,这次,你会发现许多不懂的地方也能迎刃而解了,同时还会有一些额外的收获。

四、对比 JUC 源码(JDK 1.8),一次看懂 JUC 设计

我们对比下 JUC 源码,看看除了解决虚假唤醒、永久阻塞、无锁入队等问题外,JUC 是否还做了哪些比较精妙的处理。

首先我们看 get 方法,其他的先不用看,关注这个 awaitDone:

awaitDone 源码如下:

private int awaitDone(boolean timed, long nanos) throws InterruptedException {
    final long deadline = timed ? System.nanoTime() + nanos : 0L;
    WaitNode q = null;
    boolean queued = false;
    for (;;) {
        if (Thread.interrupted()) {
            removeWaiter(q);
            throw new InterruptedException();
        }

        int s = state;
        if (s > COMPLETING) {
            if (q != null)
                q.thread = null;
            return s;
        }
        else if (s == COMPLETING) // cannot time out yet
            Thread.yield();
        else if (q == null)
            q = new WaitNode();
        else if (!queued)
            queued = UNSAFE.compareAndSwapObject(this, waitersOffset,
                                                 q.next = waiters, q);
        else if (timed) {
            nanos = deadline - System.nanoTime();
            if (nanos <= 0L) {
                removeWaiter(q);
                return state;
            }
            LockSupport.parkNanos(this, nanos);
        }
        else
            LockSupport.park(this);
    }
}

这里的 int s = state 可以理解为和 result 类似,是一个判断任务是否完成的标志,s > Completing 则可以认为任务完成(后面会讲为什么要搞一个 state)

如果我们把判断是否打断的代码,以及 removeWaiter 的代码忽略的话,你会发现,其实和我们上面的手写异曲同工!在 CAS 替换链表头结点成功后,if - else if - else 的代码块直接结束,进入下一轮循环再次判断 state(任务完成状态),而不是立刻 park,实现 Double Check;LockSupport.park 返回后,又在下一轮 while 循环中再次检查 state,如果任务未完成且未超时则再次 park;

—— 一切都和我们已知的相同。

那么,什么不同呢?

首先是 state 的设计。也许你会疑惑,为什么要搞一个 state?一个 result 感觉也够用了。实际上,我们觉得 result 够用,是因为我们的手写中没有实现 get 方法感知任务执行异常、任务打断、任务取消的功能,而完整的 FutureTask 调用 get,能知道任务执行是异常还是被打断、取消、超时等更细致的状态;

为了维护任务的各种状态,JUC 的 FutureTask 用一个 state 变量维护任务状态,用一个 outcome(即 result)变量维护任务结果;

state 和 outcome 通常一起写或一起读,所以 outcome 没加 volitile

其次是 removeWaiter 没有在我们的代码中出现,但我们之前也已经分析过,其实也可以不用 removeWaiter,让他自然被 unpark。这里 remove 可能是为了防止链表过长导致任务线程唤醒较慢;

最后不同的是,JUC FutureTask 的 get 方法在 return 任务完成状态前搞了一个设置该线程 WaitNode 的 Thread 为 null 的情况。这是为了让任务线程遍历链表时,只对非空的 thread 进行 unpark,其实也没太大必要。

再看看 JUC 的 FutureTask 的 run 方法:

public void run() {
    if (state != NEW ||
        !UNSAFE.compareAndSwapObject(this, runnerOffset,
                                     null, Thread.currentThread()))
        return;
    try {
        Callable<V> c = callable;
        if (c != null && state == NEW) {
            V result;
            boolean ran;
            try {
                result = c.call();
                ran = true;
            } catch (Throwable ex) {
                result = null;
                ran = false;
                setException(ex);
            }
            if (ran)
                set(result);
        }
    } finally {
        // runner must be non-null until state is settled to
        // prevent concurrent calls to run()
        runner = null;
        // state must be re-read after nulling runner to prevent
        // leaked interrupts
        int s = state;
        if (s >= INTERRUPTING)
            handlePossibleCancellationInterrupt(s);
    }
}

方法开始有个 CAS,是考虑了多线程都去跑 FutureTask 的 run 方法,只允许一个执行,这一般也不会用到,手写中也就没写;

然后 Callable 任务完成后,还有个 set 和 setException 的区别:这其实就是根据 Callable 任务正常执行和异常执行情况,决定怎么设置 state 和 outcome(result);

我们可以看下 set 和 setException 都干了什么:

可以发现,如果是正常执行,那么 Callable 返回值设置到 outcome(即result)中,标记 state 为正常完成;如果执行异常,那么把 exception 设置到 outcome(result)中,然后标记 state 为异常执行。最后的 finishCompletion 就是 unpark 所有等待线程;总的来说,和我们自己实现的 run 方法区别不大。

可以细看一下 finishCompletion:

private void finishCompletion() {
    // assert state > COMPLETING;
    for (WaitNode q; (q = waiters) != null;) {
        if (UNSAFE.compareAndSwapObject(this, waitersOffset, q, null)) {
            for (;;) {
                Thread t = q.thread;
                if (t != null) {
                    q.thread = null;
                    LockSupport.unpark(t);
                }
                WaitNode next = q.next;
                if (next == null)
                    break;
                q.next = null; // unlink to help gc
                q = next;
            }
            break;
        }
    }

    done();

    callable = null;        // to reduce footprint
}

可以看到,任务线程遍历 unpark 的过程又搞了个 CAS 把 waiter 的 head 换成 null,这又是为什么?其实这个 CAS 我个人认为没有太大必要,之前也分析过直接读取的可行性;感觉这里用 CAS 主要是有助于垃圾回收,同时让后续的 removeWaiter 遍历到 null 值直接结束。

也可以看到在 CAS 成功后,就是循环遍历并 unpark 的流程,这个过程就是简单遍历,和我们的代码趋向一致。

最后回到调用 awaitDone 的地方,再来看这个 report(s),它就是根据任务的状态是正常、异常还是打断、取消来决定怎么返回:

report 代码如下图:

分析下 report 的逻辑:

如果 state 是正常,那么把 outcome 用泛型转换一下返回;

如果 state 是异常,那么 outcome 其实就是 Exception 对象,将其包装成 ExecutionException 返回;(不记得这里outcome 为什么可以是 exception 对象的话,可以再看下上面的 run 方法的 setException)

如果 state 是取消,那么就抛一个取消的异常;

取消是怎么做的呢?可以看到其实也简单,当前线程 CAS 尝试把 FutureTask 的 state 从 NEW 改为已取消;CAS成功后,如果不考虑打断任务线程,就直接 finishCompletion 唤醒等待线程;如果考虑打断任务线程,就执行一个 runner.interrupt() ,然后更改 state 为打断。

终于——FutureTask 的蓝图完全展开了!

总结

最后总结一下吧:无锁 FutureTask,最核心的思想仍然是 get 方法中的那个循环,为什么要在入队后还做一步判断再 park,以及对于链表 CAS 入队和虚假唤醒的处理。

其次的一些点就是,通过 state 来维护任务的各种执行状态,并在 future.get 的最后有一个 report(state),根据任务状态来包装返回值或者抛出对应的异常。

最后感叹一下,终于写完了。。最开始没打算写这么久,看源码和打字画图都很耗时。。但是身为技术人员,有时候也不光是功利地准备面试,看源码的过程中,有些对技术的热情也在心里升腾吧。

完整版的代码可以见我的 GitHub:https://github.com/walkovertheworld/handWriteJUC

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐