Использование Runnable и Callable в ThreadPoolExecutor
При использовании ExecutorService метод submit() обычно используется для отправки запускаемых или вызываемых
Взгляните на реализацию в AbstractExecutorService.
//java.util.concurrent.AbstractExecutorService
public Future<?> submit(Runnable task) {
if (task == null) throw new NullPointerException();
RunnableFuture<Void> ftask = newTaskFor(task, null);
execute(ftask);
return ftask;
}
public <T> Future<T> submit(Callable<T> task) {
if (task == null) throw new NullPointerException();
RunnableFuture<T> ftask = newTaskFor(task);
execute(ftask);
return ftask;
}
Вы можете видеть, что объект RunnableFuture генерируется методом newTaskFor().
protected <T> RunnableFuture<T> newTaskFor(Runnable runnable, T value) {
return new FutureTask<T>(runnable, value);
}
protected <T> RunnableFuture<T> newTaskFor(Callable<T> callable) {
return new FutureTask<T>(callable);
}
Взгляните на конструкцию RunnableFuture
//java.util.concurrent.FutureTask
public FutureTask(Callable<V> callable) {
if (callable == null)
throw new NullPointerException();
this.callable = callable;
this.state = NEW; // ensure visibility of callable
}
public FutureTask(Runnable runnable, V result) {
//包装runnable->callable
this.callable = Executors.callable(runnable, result);
this.state = NEW; // ensure visibility of callable
}
//java.util.concurrent.Executors
public static <T> Callable<T> callable(Runnable task, T result) {
if (task == null)
throw new NullPointerException();
return new RunnableAdapter<T>(task, result);
}
//java.util.concurrent.Executors.RunnableAdapter
//适配器模式 runnable->callable
static final class RunnableAdapter<T> implements Callable<T> {
final Runnable task;
final T result;
RunnableAdapter(Runnable task, T result) {
this.task = task;
this.result = result;
}
public T call() {
task.run();
return result;
}
}
Из приведенного выше кода вы можете увидеть шаги реализации методов submit(Runnable) и submit(Callable).
- Пустой
- Создать объект RunnableFuture
- вызвать метод execute()
//java.util.concurrent.ThreadPoolExecutor
//任务队列
private final BlockingQueue<Runnable> workQueue;
public void execute(Runnable command) {
...
workQueue.offer(command)
...
}
Из сигнатуры метода execute(Runnable) можно сделать вывод, что RunnableFuture является классом реализации Runnable.
FutureTask是RunnableFuture的实现类
Взгляните на отношения наследования FutureTask
RunnableFuture — эторежим адаптераРеализация адаптирует Future к Runable
Посмотрите, как FutureTask реализует интерфейс Runnable.
public void run() {
...
try {
Callable<V> c = callable;
if (c != null && state == NEW) {
V result;
boolean ran;
try {
//执行callable.call()方法
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);
}
}
Суммировать:
- Объект Runnable, хранящийся в очереди задач в ThreadPoolExecutor.
- Запускаемые объекты, отправленные через метод execute, будут напрямую помещены в очередь задач.
- Объект Callable, отправленный с помощью метода submit(Callable), будет обернут объектом FutureTask, а затем помещен в очередь задач.
- Runnable, отправленный с помощью метода submit(Runnable), будет упакован в объект RunnableAdapter (реализация Callable), затем упакован в объект FutureTask, а затем помещен в очередь задач.
Future
Будущее: асинхронно вычисляемый объект-заполнитель для получения результата, который будет вычисляться
Будущие функции:
- Получить результат, который будет вычисляться с помощью get()
- Отмените вычисление результата с помощью метода cancel().
Реализация FutureTask.cancel()
Несколько состояний FutureTask
private static final int NEW = 0;
private static final int COMPLETING = 1;
private static final int NORMAL = 2;
private static final int EXCEPTIONAL = 3;
private static final int CANCELLED = 4;
private static final int INTERRUPTING = 5;
private static final int INTERRUPTED = 6;
//可能出现的状态变化过程
* NEW -> COMPLETING -> NORMAL
* NEW -> COMPLETING -> EXCEPTIONAL
* NEW -> CANCELLED
* NEW -> INTERRUPTING -> INTERRUPTED
отменить процесс реализации
public boolean cancel(boolean mayInterruptIfRunning) {
// 1.状态判断
// 只有state==New且通过cas修改state值成功 才往下执行 否则return false
if (!(state == NEW &&
//状态变化
// mayInterruptIfRunning? NEW->INTERRUPTING:NEW->CANCELLED
UNSAFE.compareAndSwapInt(this, stateOffset, NEW,
mayInterruptIfRunning ? INTERRUPTING : CANCELLED)))
return false;
try { // in case call to interrupt throws exception
if (mayInterruptIfRunning) {//2.打断运行
try {
Thread t = runner;
if (t != null)
t.interrupt();
} finally { // INTERRUPTING->INTERRUPTED
UNSAFE.putOrderedInt(this, stateOffset, INTERRUPTED);
}
}
} finally {
//3.结束
finishCompletion();
}
return true;
}
Анализ процесса:
-
государственный приговор
Выражение if можно разделить на состояние == NEW и UNSAFE.compareAndSwapInt(this, stateOffset, NEW,mayInterruptIfRunning? ПРЕРЫВАНИЕ: ОТМЕНА)
-
state == NEW, чтобы определить, является ли текущее состояние NEW
-
UNSAFE.compareAndSwapInt(this, stateOffset, NEW,mayInterruptIfRunning ? INTERRUPTING : CANCELLED)
Это модификация cas, эквивалентная state=mayInterruptIfRunning ? INTERRUPTING : CANCELED
Способ cas в основном для атомных соображений
unsafe的用法请自行查阅资料
Резюме: если статус NEW и блокировка получена, может быть выполнена реальная операция отмены.
-
-
прерывание операции
mayInterruptIfRunning?thread.interrupt(): нет операции;
-
конец
finishCompletion();
Суммировать:
- Отмена может быть выполнена только в том случае, если состояние == НОВОЕ.
- Изменить состояние с помощью небезопасной операции cas
- Изменение состояния две строки
- NEW->CANCLE
- NEW->INTERRUPTING->INTERRUPTED
Реализация FutureTask.run()
public void run() {
// 1.状态判断
// 只有state==New且通过cas设置runner值成功 才往下执行 否则return false
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 {
//执行callable.call()
result = c.call();
ran = true;
} catch (Throwable ex) {
result = null;
ran = false;
//NEW->COMPLETING->EXCEPTIONAL
setException(ex);
}
if (ran)//NEW->COMPLETING->NORMAL
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);
}
}
protected void setException(Throwable t) {
// 防止cancle()方法修改state
if (UNSAFE.compareAndSwapInt(this, stateOffset, NEW, COMPLETING)) {
outcome = t;
UNSAFE.putOrderedInt(this, stateOffset, EXCEPTIONAL); // final state
finishCompletion();
}
}
protected void set(V v) {
// 防止cancle()方法修改state
if (UNSAFE.compareAndSwapInt(this, stateOffset, NEW, COMPLETING)) {
outcome = v;
UNSAFE.putOrderedInt(this, stateOffset, NORMAL); // final state
finishCompletion();
}
}
Суммировать:
- Метод выполняет процесс run() только тогда, когда state==NEW
- Установить бегун (текущий поток) через небезопасный cas
- Изменение состояния Две строки:
- NEW->COMPLETING->EXCEPTIONAL
- NEW->COMPLETING->NORMAL
cancel()Сводка
- еслиrun() не был выполненЗатем очистите вызываемый объект и измените состояние на не-НОВОЕ (поэтому метод run() не будет выполняться)
- еслиrun() выполняется, а callable.call() не завершил выполнениеЗатем вызовите thread.interrpt(), чтобы уведомить поток об остановке (просто заметьтеНет гарантии, что поток будет прерван. Пожалуйста, обратитесь к информации о прерывании () для конкретных причин. Поскольку состояние состояния изменяется с помощью отмены, setException () и set () не могут сохранить результат.
- еслиrun() завершен или callable.call() завершенCancle() не продолжает выполняться, потому что state!=NEW возвращает ошибку
Реализация FutureTask.get()
public V get() throws InterruptedException, ExecutionException {
int s = state;
if (s <= COMPLETING)//阻塞等待
s = awaitDone(false, 0L);
return report(s);
}
private int awaitDone(boolean timed, long nanos)
throws InterruptedException {
final long deadline = timed ? System.nanoTime() + nanos : 0L;
WaitNode q = null;
boolean queued = false;
for (;;) {
//cancle()过程中调用thread.interrupt()则退出循环 并抛异常
if (Thread.interrupted()) {
removeWaiter(q);
throw new InterruptedException();
}
int s = state;
if (s > COMPLETING) {//执行完成(包括执行异常) 或被取消 返回当前status
if (q != null)
q.thread = null;
return s;
}
else if (s == COMPLETING) //执行完成但尚未修改状态 则Thread.yield()让出cpu资源
Thread.yield();
else if (q == null)//尚未执行完成 则加入生成等待节点
q = new WaitNode();
else if (!queued)//当前等待节点尚未加入等待队列 则cas方式加入等待队列
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);
}
}
//根据state 返回不同结果
private V report(int s) throws ExecutionException {
Object x = outcome;
if (s == NORMAL)
return (V)x;
if (s >= CANCELLED)
throw new CancellationException();
throw new ExecutionException((Throwable)x);
}
//run()和cancle()最终都会调用finishCompletion() 分析是如何唤醒等待队列中的节点的
private void finishCompletion() {
// assert state > COMPLETING;
for (WaitNode q; (q = waiters) != null;) {
// cas方式修改将等待队列置空
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; // callable置空
}
get()Сводка:
- Генерировать ожидающие узлы и присоединяться к ожидающей очереди
- Вращение через LockSupport.park() Ожидание пробуждения
- Обтекание результатов выполнения задачи по состоянию