Netty: интерпретация исходного кода DefaultPromise

задняя часть

1. Зачем вам нужен io.netty.util.concurrent.Promise?

Если у вас есть метод блокировки, такой как Thread.sleep(1000), и вы не хотите блокировать текущий поток A, просто оберните метод в задачу, которая будет выполняться другим потоком B.

ExecutorService pool = Executors.newFixedThreadPool(3);
Future<Integer> future = pool.submit(() -> {
    Thread.sleep(1000);
    return 1;
});

Если вам нужно выполнить другую логику после завершения задачи, одним из способов является первый вызов потока A.future.get()Получите значение, а затем выполните другой код, но сам метод get также является методом блокировки, во время которого блокируется поток A.

Другой метод продолжит выполнение последующей логики после завершения потока B, выполняющих задачу. Будущее в Netty, io.netty.util.concurrent.future, через метод обратного вызоваFuture<V> addListener(GenericFutureListener<? extends Future<? super V>> listener);Эта функция реализована.

Интерфейс Promise наследует интерфейс Future и в случае добавления слушателя предоставляетPromise<V> setSuccess(V result)метод, вы можете вручную установить возвращаемое значение в задаче и немедленно уведомить слушателей.

Во-вторых, пример программы

private static NioEventLoopGroup loopGroup = new NioEventLoopGroup(8);

public void methodA() {
    Promise promise = methodA("ceee...eeeb");
    promise.addListener(future -> {		// 1
        Object ret = future.get();      // 4. 此时可以直接拿到结果
      	// 后续逻辑由 B 线程执行
        System.out.println(ret);
    });
  	// A 线程不阻塞,继续执行其他代码...
}

public Promise<ResponsePacket> methodB(String name) {
    Promise<ResponsePacket> promise = new DefaultPromise<>(loopGroup.next());
    loopGroup.schedule(() -> {		// 2
        try {
            Thread.sleep(1000);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        System.out.println("scheduler thread: " + Thread.currentThread().getName());
        promise.setSuccess("hello " + name);	// 3
    }, 0, TimeUnit.SECONDS);

    return promise;
}

Простое использование промисов включает в себя:

  1. Добавьте слушателя к обещанию,promise.addListener();
  2. Назначьте потоки для выполнения задач,loopGroup.schedule();
  3. Во время выполнения задачи установить результат,promise.setSuccess();

3. Анализ исходного кода

1. addListener

// class: DefaultPromise
public Promise<V> addListener(GenericFutureListener<? extends Future<? super V>> listener) {
    checkNotNull(listener, "listener");

    synchronized (this) {  // 1. 增加 listener
        addListener0(listener);	
    }

    if (isDone()) {  // 2. 如果任务执行完了,通知所有 listener
        notifyListeners();
    }

    return this;
}

Продолжайте смотреть на addListener0:

private Object listeners;
private void addListener0(GenericFutureListener<? extends Future<? super V>> listener) {
 	// 1. 添加第 1 个 listener 时,直接赋值即可
    if (listeners == null) {
        listeners = listener;
    } 
  	// 3. 添加第 3 个以及更多 listener 时,直接加入数组即可
  	else if (listeners instanceof DefaultFutureListeners) {
        ((DefaultFutureListeners) listeners).add(listener);
    } 
  	// 2. 添加第 2 个 listener 时,listeners 类型更改为 DefaultFutureListeners,内部实现为一个数组
  	else {
        listeners = new DefaultFutureListeners((GenericFutureListener<?>) listeners, listener);
    }
}

Поскольку можно добавить несколько слушателей, легко представить, что все слушатели хранятся в массиве. Тип слушателей в классе реализации — Object, вероятно, потому, что у большинства из них есть только один слушатель, что экономит место в памяти.

2. schedule

Задача добавляется в очередь для выполнения пулом потоков.

3. setSuccess

// class: DefaultPromise
public Promise<V> setSuccess(V result) {
    if (setSuccess0(result)) {	// 如果设置成功,返回;否则抛异常
        return this;
    }
    throw new IllegalStateException("complete already: " + this);
}

private boolean setSuccess0(V result) {
  	// 设置 result
    return setValue0(result == null ? SUCCESS : result);
}

private boolean setValue0(Object objResult) {
    // cas 操作
    if (RESULT_UPDATER.compareAndSet(this, null, objResult) ||
        RESULT_UPDATER.compareAndSet(this, UNCANCELLABLE, objResult)) {
        if (checkNotifyWaiters()) {
            notifyListeners();
        }
        return true;
    }
    return false;
}

private synchronized boolean checkNotifyWaiters() {
      /**
       * 有些线程不是通过增加 listener 的方式获取结果,而是通过 promise.get() 方法获取,
       * 那么这些线程为阻塞状态;当设置了 result 后,需要唤醒这些线程
       */
    if (waiters > 0) {
        notifyAll();
    }
    return listeners != null;  // 只要存在 listener,就返回 true
}

Продолжить просмотр notifyListeners:

// class: DefaultPromise
private void notifyListeners() {
    EventExecutor executor = executor();
    if (executor.inEventLoop()) {
        final InternalThreadLocalMap threadLocals = InternalThreadLocalMap.get();
        final int stackDepth = threadLocals.futureListenerStackDepth();
      	// TODO 嵌套监听
        if (stackDepth < MAX_LISTENER_STACK_DEPTH) {
            threadLocals.setFutureListenerStackDepth(stackDepth + 1);
            try {
              	// 1. 如果是 promise 绑定的线程,直接执行
                notifyListenersNow();
            } finally {
                threadLocals.setFutureListenerStackDepth(stackDepth);
            }
            return;
        }
    }

    // 2. 否则,加入任务调度, 因此 listener 方法最终还是由 promise 绑定的线程执行的
    safeExecute(executor, new Runnable() {
        @Override
        public void run() {
            notifyListenersNow();
        }
    });
}

private void notifyListenersNow() {
    Object listeners;
    synchronized (this) {
        if (notifyingListeners || this.listeners == null) {
            return;
        }
        notifyingListeners = true;
        listeners = this.listeners;
        this.listeners = null;
    }
    for (;;) {
      	// 依次通知所有 listener
        if (listeners instanceof DefaultFutureListeners) {
            notifyListeners0((DefaultFutureListeners) listeners);
        } else {
            notifyListener0(this, (GenericFutureListener<?>) listeners);
        }
        synchronized (this) {
            if (this.listeners == null) {
                notifyingListeners = false;
                return;
            }
            // 通知原先的 listeners 时,有可能有新的 listener 在此期间注册, 也需要通知到
            listeners = this.listeners;
            this.listeners = null;
        }
    }
}

private static void notifyListener0(Future future, GenericFutureListener l) {
    try {
        l.operationComplete(future);	// 执行 listener 中的方法
    } catch (Throwable t) {
        if (logger.isWarnEnabled()) {
            logger.warn("An exception was thrown by " + l.getClass().getName() + ".operationComplete()", t);
        }
    }
}