首先我們來看看future和promise介面整體設計

最頂層的future是jdk的,第二個是netty自定義的future,兩個同名,繼承關系
看看jdk的future介面
public interface Future<V> { // 取消任務 boolean cancel(boolean mayInterruptIfRunning); // 任務是否取消 boolean isCancelled(); // 任務是否完成 boolean isDone(); // 阻塞的獲取執行的結果 V get() throws InterruptedException, ExecutionException; // 在一定時間內-超時-阻塞的獲取執行的結果 V get(long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException; }
瞅瞅netty的future介面
public interface Future<V> extends java.util.concurrent.Future<V> { // 是否成功 boolean isSuccess(); // 是否取消 boolean isCancellable(); Throwable cause(); // 添加listener進行回呼 Future<V> addListener(GenericFutureListener<? extends Future<? super V>> listener); Future<V> addListeners(GenericFutureListener<? extends Future<? super V>>... listeners); Future<V> removeListener(GenericFutureListener<? extends Future<? super V>> listener); Future<V> removeListeners(GenericFutureListener<? extends Future<? super V>>... listeners); // 阻塞的等待任務執行,如果失敗則拋出失敗原因的例外 Future<V> sync() throws InterruptedException; // 不回應中斷等待例外 Future<V> syncUninterruptibly(); // 阻塞等待任務執行,失敗不拋例外 Future<V> await() throws InterruptedException; Future<V> awaitUninterruptibly(); boolean await(long timeout, TimeUnit unit) throws InterruptedException; boolean await(long timeoutMillis) throws InterruptedException; boolean awaitUninterruptibly(long timeout, TimeUnit unit); boolean awaitUninterruptibly(long timeoutMillis); // 馬上獲取到任務的結果,不阻塞,而jdk的future是阻塞的 V getNow(); // 取消任務執行,如果取消成功,任務會因為 CancellationException 例外而導致失敗 // 也就是 isSuccess()==false,同時上面的 cause() 方法回傳 CancellationException 的實體, // mayInterruptIfRunning 說的是:是否對正在執行該任務的執行緒進行中斷(這樣才能停止該任務的執行), // 似乎 Netty 中 Future 介面的各個實作類,都沒有使用這個引數 @Override boolean cancel(boolean mayInterruptIfRunning); }
netty的future在jdk的基礎上擴展了它需要的方法,sync和await的區別我們放到下面看實作類的時候說
同時我們也可以看到,這個future介面跟io操作是無關的
接下來我們看看ChannelFuture介面,介面注釋上寫的很清楚,我們來看看
* The result of an asynchronous {@link Channel} I/O operation.
* <p>
* All I/O operations in Netty are asynchronous. It means any I/O calls will
* return immediately with no guarantee that the requested I/O operation has
* been completed at the end of the call. Instead, you will be returned with
* a {@link ChannelFuture} instance which gives you the information about the
* result or status of the I/O operation.
* <p>
* A {@link ChannelFuture} is either <em>uncompleted</em> or <em>completed</em>.
* When an I/O operation begins, a new future object is created. The new future
* is uncompleted initially - it is neither succeeded, failed, nor cancelled
* because the I/O operation is not finished yet. If the I/O operation is
* finished either successfully, with failure, or by cancellation, the future is
* marked as completed with more specific information, such as the cause of the
* failure. Please note that even failure and cancellation belong to the
* completed state.
所有io操作都是異步的,一個io操作的呼叫會立即回傳一個帶有結果或者狀態的io實體,
io操作要么是未完成的,要么是完成的,當它開始時,future會被創建,一開始是未完成的,未完成的時候沒有成功、失敗或者取消狀態
當它是完成的時候,可以是失敗或者取消的,失敗或者取消原因會被附加到future上,
* <pre> * +---------------------------+ * | Completed successfully | * +---------------------------+ * +----> isDone() = true | * +--------------------------+ | | isSuccess() = true | * | Uncompleted | | +===========================+ * +--------------------------+ | | Completed with failure | * | isDone() = false | | +---------------------------+ * | isSuccess() = false |----+----> isDone() = true | * | isCancelled() = false | | | cause() = non-null | * | cause() = null | | +===========================+ * +--------------------------+ | | Completed by cancellation | * | +---------------------------+ * +----> isDone() = true | * | isCancelled() = true | * +---------------------------+ * </pre>
上面那個狀態遷移圖很清楚了,在兩種程序的時候會有什么狀態,我們看看介面
public interface ChannelFuture extends Future<Void> { // 回傳future關聯的channel Channel channel(); // 重寫下面幾個方法,修改回傳值為channelfuture @Override ChannelFuture addListener(GenericFutureListener<? extends Future<? super Void>> listener); @Override ChannelFuture addListeners(GenericFutureListener<? extends Future<? super Void>>... listeners); @Override ChannelFuture removeListener(GenericFutureListener<? extends Future<? super Void>> listener); @Override ChannelFuture removeListeners(GenericFutureListener<? extends Future<? super Void>>... listeners); @Override ChannelFuture sync() throws InterruptedException; @Override ChannelFuture syncUninterruptibly(); @Override ChannelFuture await() throws InterruptedException; @Override ChannelFuture awaitUninterruptibly(); /** * Returns {@code true} if this {@link ChannelFuture} is a void future and so not allow to call any of the * following methods: * <ul> * <li>{@link #addListener(GenericFutureListener)}</li> * <li>{@link #addListeners(GenericFutureListener[])}</li> * <li>{@link #await()}</li> * <li>{@link #await(long, TimeUnit)} ()}</li> * <li>{@link #await(long)} ()}</li> * <li>{@link #awaitUninterruptibly()}</li> * <li>{@link #sync()}</li> * <li>{@link #syncUninterruptibly()}</li> * </ul> 標記該future是void的,使不能使用上面的方法 */ boolean isVoid(); }
netty其實是強烈建議直接通過添加監聽器的方式來獲取io操作結果,或者進行后續操作的,ChannelFuture可以增加或者洗掉一個多個 GenericFutureListener,它定義如下
public interface GenericFutureListener<F extends Future<?>> extends EventListener { void operationComplete(F future) throws Exception; }
執行完后會回呼 operationComplete方法
注意一點,不要在ChannelHandler中呼叫ChannelFuture的await方法,會導致死鎖,這是因為發起io操作后,由io執行緒負責異步通知發起io操作的用戶執行緒,如果io執行緒和用戶執行緒是同一個的話,就會導致io執行緒等待自己通知操作完成,這就會導致死鎖,自己掛死自己,
我們繼續看promise介面,
public interface Promise<V> extends Future<V> { // 標記該future成功及設定結果,并通知所有listener // 如果失敗的話拋例外 Promise<V> setSuccess(V result); // 和setsuccess一樣,只是失敗的話回傳false boolean trySuccess(V result); // 標記future失敗,然后通知listener Promise<V> setFailure(Throwable cause); boolean tryFailure(Throwable cause); // 標記該future 不可被取消 boolean setUncancellable(); // 下面跟ChannelFuture一樣,都是覆寫重寫方法 @Override Promise<V> addListener(GenericFutureListener<? extends Future<? super V>> listener); @Override Promise<V> addListeners(GenericFutureListener<? extends Future<? super V>>... listeners); @Override Promise<V> removeListener(GenericFutureListener<? extends Future<? super V>> listener); @Override Promise<V> removeListeners(GenericFutureListener<? extends Future<? super V>>... listeners); @Override Promise<V> await() throws InterruptedException; @Override Promise<V> awaitUninterruptibly(); @Override Promise<V> sync() throws InterruptedException; @Override Promise<V> syncUninterruptibly(); }
Promise是可寫的future,Future本身并沒有寫操作相關的介面,netty通過Promise對其進行擴展,用于設定io操作的結果,Promise 實體內部是一個任務,任務的執行往往是異步的,通常是一個執行緒池來處理任務,Promise 提供的 setSuccess(V result) 或 setFailure(Throwable t) 將來會被某個執行任務的執行緒在執行完成以后呼叫,同時那個執行緒在呼叫 setSuccess(result) 或 setFailure(t) 后會回呼 listeners 的回呼函式(當然,回呼的具體內容不一定要由執行任務的執行緒自己來執行,它可以創建新的執行緒來執行,也可以將回呼任務提交到某個執行緒池來執行),而且,一旦 setSuccess(...) 或 setFailure(...) 后,那些 await() 或 sync() 的執行緒就會從等待中回傳,
接下來我們看看ChannelPromise
public interface ChannelPromise extends ChannelFuture, Promise<Void> { @Override Channel channel();
@Override ChannelPromise setSuccess(Void result); ChannelPromise setSuccess(); boolean trySuccess(); @Override ChannelPromise setFailure(Throwable cause);
@Override ChannelPromise addListener(GenericFutureListener<? extends Future<? super Void>> listener); @Override ChannelPromise addListeners(GenericFutureListener<? extends Future<? super Void>>... listeners); @Override ChannelPromise removeListener(GenericFutureListener<? extends Future<? super Void>> listener); @Override ChannelPromise removeListeners(GenericFutureListener<? extends Future<? super Void>>... listeners);
@Override ChannelPromise sync() throws InterruptedException; @Override ChannelPromise syncUninterruptibly(); @Override ChannelPromise await() throws InterruptedException; @Override ChannelPromise awaitUninterruptibly(); /** * Returns a new {@link ChannelPromise} if {@link #isVoid()} returns {@code true} otherwise itself. */ ChannelPromise unvoid(); }
看方法其實很清楚,基本都是覆寫綜合了ChannelFuture和Promise介面的,就回傳值變了,看看我們一開始的類繼承圖,ChannelPromise 介面同時繼承了 ChannelFuture 和 Promise,最終繼承的都是Future介面,接下來我們看看具體的實作類DefaultPromise吧
public class DefaultPromise<V> extends AbstractFuture<V> implements Promise<V> { // 為了后面操作成功后通過cas來保存結果到result欄位 @SuppressWarnings("rawtypes") private static final AtomicReferenceFieldUpdater<DefaultPromise, Object> RESULT_UPDATER = AtomicReferenceFieldUpdater.newUpdater(DefaultPromise.class, Object.class, "result"); // result為null的時候默認值 private static final Object SUCCESS = new Object(); // 操作成功后cas比對的值 private static final Object UNCANCELLABLE = new Object(); // 保存執行的結果 private volatile Object result; // 執行緒執行器 private final EventExecutor executor; // 監聽者 private Object listeners; /** * Threading - synchronized(this). We are required to hold the monitor to use Java's underlying wait()/notifyAll(). */ // 等待這個 promise 的執行緒數(呼叫sync()/await()進行等待的執行緒數量) private short waiters; /** * Threading - synchronized(this). We must prevent concurrent notification and FIFO listener notification if the * executor changes. */ // 是否喚醒正在等待執行緒,用于防止重復執行喚醒,不然會重復執行 listeners 的回呼方法 private boolean notifyingListeners; .... }
屬性看完了,我們可以看看它主要的方法
@Override public Promise<V> setSuccess(V result) { if (setSuccess0(result)) { notifyListeners(); return this; } throw new IllegalStateException("complete already: " + this); } @Override public boolean trySuccess(V result) { if (setSuccess0(result)) { notifyListeners(); return true; } return false; } @Override public Promise<V> setFailure(Throwable cause) { if (setFailure0(cause)) { notifyListeners(); return this; } throw new IllegalStateException("complete already: " + this, cause); } @Override public boolean tryFailure(Throwable cause) { if (setFailure0(cause)) { notifyListeners(); return true; } return false; }
set和try的區別就是回傳值不一樣而已,我們看看底層的方法 setSuccess0
private boolean setSuccess0(V result) { return setValue0(result == null ? SUCCESS : result); } private boolean setValue0(Object objResult) { if (RESULT_UPDATER.compareAndSet(this, null, objResult) || RESULT_UPDATER.compareAndSet(this, UNCANCELLABLE, objResult)) { checkNotifyWaiters(); return true; } return false; }
就是通過cas來把objResult保存到result屬性上,然后Notify其他執行緒,其他方法都差不多,可以比對看看
我們再看個await方法
public Promise<V> await() throws InterruptedException { if (isDone()) { return this; } if (Thread.interrupted()) { throw new InterruptedException(toString()); } checkDeadLock(); synchronized (this) { while (!isDone()) { incWaiters(); try { wait(); } finally { decWaiters(); } } } return this; }
如果當前Promise已被設定,則回傳;如果碰到執行緒中斷則回應中斷;檢查死鎖,由于在IO執行緒中呼叫Promise的await方法或者sync方法會導致死鎖,前面說過的,所以需要檢驗保護,判定當前執行緒是否是io執行緒;同步鎖定當前Promise物件,回圈判定是否設定完成,使用回圈是避免偽喚醒,防止執行緒 被意外喚醒導致功能例外,
接下來我們順便也看下sync方法
@Override public Promise<V> sync() throws InterruptedException { await(); rethrowIfFailed(); return this; }
首先呼叫await方法,然后看是否需要拋出例外,如果任務失敗的話就重新拋出例外,這也是兩方法區別了,
DefaultChannelPromise實作我們就不看了,基本都是基于DefaultPromise的,只是回傳值都是 ChannelPromise而已,
下面我們來寫個例子吧
public class ChannelPromiseExample extends Thread{ private static final Object object = new Object(); public static void main(String[] args) { final DefaultEventExecutor executor = new DefaultEventExecutor(); final Promise<Integer> promise = executor.newPromise(); // 任務seccess或者failure來回呼operationComplete 方法 promise.addListener(new GenericFutureListener<Future<? super Integer>>() { @Override public void operationComplete(Future<? super Integer> future) throws Exception { System.out.println(Thread.currentThread().getName() + " 第一個監聽器"); if (future.isSuccess()) { System.out.println("任務成功,result:" + future.get()); } else { System.out.println("任務失敗,result:" + future.cause()); } } }).addListener(new GenericFutureListener<Future<? super Integer>>() { @Override public void operationComplete(Future<? super Integer> future) throws Exception { System.out.println(Thread.currentThread().getName() + " 第二個監聽器"); } }); // 提交任務 executor.execute(new Runnable() { @Override public void run() { try { Thread.sleep(2000); } catch (InterruptedException e) { e.printStackTrace(); } //可以設定成功或者失敗 //promise.setSuccess(1); promise.setFailure(new Throwable("FAILURE")); } }); try { System.out.println("promise wait begin"); //promise.sync(); promise.await(); System.out.println("promise wait end"); } catch (InterruptedException e) { e.printStackTrace(); } finally { executor.shutdown(); } } }
可以體會下await和sync的區別
轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/569.html
標籤:其他
下一篇:Yii2原始碼分析(一):入口
