NOTE

11.6 源代码层面

时序图 源码追踪 由启动流程.md tomcat启动由两个主要步骤,一个是init,另一个是start init主要的工作是绑定端口号(bind),start主要的工作是监听客户端链接(accept) 其中监听客户端链接的代码在这段逻辑 NioEndpoint - startInternal Acc

Java创建于 更新于 historical

这是历史学习笔记,可能存在过时或不完整的理解。

时序图

源码追踪

启动流程.md tomcat启动由两个主要步骤,一个是init,另一个是start init主要的工作是绑定端口号(bind),start主要的工作是监听客户端链接(accept) 其中监听客户端链接的代码在这段逻辑

NioEndpoint

  • startInternal
public void startInternal() throws Exception {

    if (!running) {
        running = true;
        paused = false;

        processorCache = new SynchronizedStack<>(SynchronizedStack.DEFAULT_SIZE,
                socketProperties.getProcessorCache());
        eventCache = new SynchronizedStack<>(SynchronizedStack.DEFAULT_SIZE,
                        socketProperties.getEventCache());
        nioChannels = new SynchronizedStack<>(SynchronizedStack.DEFAULT_SIZE,
                socketProperties.getBufferPool());

        // 创建线程池作为worker,处理已经链接的请求
        if ( getExecutor() == null ) {
            createExecutor();
        }

        initializeConnectionLatch();

        // 需要搞懂poller和accptor的关系
        //acceptor接受的链接封装好后丢给poller,poller调用线程池的worker处理链接
        pollers = new Poller[getPollerThreadCount()];
        for (int i=0; i<pollers.length; i++) {
            pollers[i] = new Poller();
            Thread pollerThread = new Thread(pollers[i], getName() + "-ClientPoller-"+i);
            pollerThread.setPriority(threadPriority);
            pollerThread.setDaemon(true);
            pollerThread.start();
        }
		//Acceptor线程用于接收客户端的链接请求
		//org.apache.tomcat.util.net.AbstractEndpoint#startAcceptorThreads
        startAcceptorThreads();
    }
}

protected final void startAcceptorThreads() {
    int count = getAcceptorThreadCount();
    acceptors = new Acceptor[count];

    for (int i = 0; i < count; i++) {
    	//创建的Acceptor是个Runnable实例
    	//org.apache.tomcat.util.net.NioEndpoint.Acceptor
        acceptors[i] = createAcceptor();
        String threadName = getName() + "-Acceptor-" + i;
        acceptors[i].setThreadName(threadName);
        Thread t = new Thread(acceptors[i], threadName);
        t.setPriority(getAcceptorThreadPriority());
        t.setDaemon(getDaemon());
        t.start();//启动这个Runnable
    }
}

Acceptor

  • Acceptor run
protected class Acceptor extends AbstractEndpoint.Acceptor {

    @Override
    public void run() {

        int errorDelay = 0;

        // Loop until we receive a shutdown command
        while (running) {

            // Loop if endpoint is paused
            while (paused && running) {
                state = AcceptorState.PAUSED;
                try {
                    Thread.sleep(50);
                } catch (InterruptedException e) {
                    // Ignore
                }
            }

            if (!running) {
                break;
            }
            state = AcceptorState.RUNNING;

            try {
                //if we have reached max connections, wait
                countUpOrAwaitConnection();

                SocketChannel socket = null;
                try {
                    //阻塞等待客户端链接
                    socket = serverSock.accept();
                } catch (IOException ioe) {
                    // We didn't get a socket
                    countDownConnection();
                    if (running) {
                        // Introduce delay if necessary
                        errorDelay = handleExceptionWithDelay(errorDelay);
                        // re-throw
                        throw ioe;
                    } else {
                        break;
                    }
                }
                // Successful accept, reset the error delay
                errorDelay = 0;

                // Configure the socket
                if (running && !paused) {
                    // 把socket交给poller处理
                    //getPoller0().register(channel);

                    if (!setSocketOptions(socket)) {
                        closeSocket(socket);
                    }
                } else {
                    closeSocket(socket);
                }
            } catch (Throwable t) {
                ExceptionUtils.handleThrowable(t);
                log.error(sm.getString("endpoint.accept.fail"), t);
            }
        }
        state = AcceptorState.ENDED;
    }



}

我们在org.apache.tomcat.util.net.NioEndpoint.Acceptor的run方法上打个断点,任何一个客户端链接过来都要经过这段逻辑。 之后会将这个socket转交给poller处理,我们继续研究Poller

Poller

poller也是一个Runnable,不停的检查是否由新事件到达,有的话进行处理

  • Poller run
public void run() {
    // Loop until destroy() is called
    while (true) {

        boolean hasEvents = false;

        try {
            if (!close) {
                hasEvents = events();
                if (wakeupCounter.getAndSet(-1) > 0) {
                    //if we are here, means we have other stuff to do
                    //do a non blocking select
                    keyCount = selector.selectNow();
                } else {
                    keyCount = selector.select(selectorTimeout);
                }
                wakeupCounter.set(0);
            }
            if (close) {
                events();
                timeout(0, false);
                try {
                    selector.close();
                } catch (IOException ioe) {
                    log.error(sm.getString("endpoint.nio.selectorCloseFail"), ioe);
                }
                break;
            }
        } catch (Throwable x) {
            ExceptionUtils.handleThrowable(x);
            log.error("",x);
            continue;
        }
        //either we timed out or we woke up, process events first
        if ( keyCount == 0 ) hasEvents = (hasEvents | events());

		//遍历所有事件
        Iterator<SelectionKey> iterator =
            keyCount > 0 ? selector.selectedKeys().iterator() : null;
        // Walk through the collection of ready keys and dispatch
        // any active event.
        while (iterator != null && iterator.hasNext()) {
            SelectionKey sk = iterator.next();
            NioSocketWrapper attachment = (NioSocketWrapper)sk.attachment();
            // Attachment may be null if another thread has called
            // cancelledKey()
            if (attachment == null) {
                iterator.remove();
            } else {
                iterator.remove();
            	//进行处理
                processKey(sk, attachment);
            }
        }//while

        //process timeouts
        timeout(keyCount,hasEvents);
    }//while

    getStopLatch().countDown();
}
  • processKey
protected void processKey(SelectionKey sk, NioSocketWrapper attachment) {
    try {
        if ( close ) {
            cancelledKey(sk);
        } else if ( sk.isValid() && attachment != null ) {
            if (sk.isReadable() || sk.isWritable() ) {
                if ( attachment.getSendfileData() != null ) {
                    processSendfile(sk,attachment, false);
                } else {
                    unreg(sk, attachment, sk.readyOps());
                    boolean closeSocket = false;
                    // Read goes before write
                    if (sk.isReadable()) {
                    	//处理读事件
                        if (!processSocket(attachment, SocketEvent.OPEN_READ, true)) {
                            closeSocket = true;
                        }
                    }
                    if (!closeSocket && sk.isWritable()) {
                    	//处理写事件
                        if (!processSocket(attachment, SocketEvent.OPEN_WRITE, true)) {
                            closeSocket = true;
                        }
                    }
                    if (closeSocket) {
                        cancelledKey(sk);
                    }
                }
            }
        } else {
            //invalid key
            cancelledKey(sk);
        }
    } catch ( CancelledKeyException ckx ) {
        cancelledKey(sk);
    } catch (Throwable t) {
        ExceptionUtils.handleThrowable(t);
        log.error("",t);
    }
}
  • processSocket
public boolean processSocket(SocketWrapperBase<S> socketWrapper,
        SocketEvent event, boolean dispatch) {
    try {
        if (socketWrapper == null) {
            return false;
        }
        SocketProcessorBase<S> sc = processorCache.pop();
        if (sc == null) {
        	//把socket事件封装成SocketProcessorBase
            sc = createSocketProcessor(socketWrapper, event);
        } else {
            sc.reset(socketWrapper, event);
        }
        Executor executor = getExecutor();
        if (dispatch && executor != null) {
        	//调用线程池处理SocketProcessorBase
            executor.execute(sc);
        } else {
            sc.run();
        }
    } catch (RejectedExecutionException ree) {
        getLog().warn(sm.getString("endpoint.executor.fail", socketWrapper) , ree);
        return false;
    } catch (Throwable t) {
        ExceptionUtils.handleThrowable(t);
        // This means we got an OOM or similar creating a thread, or that
        // the pool and its queue are full
        getLog().error(sm.getString("endpoint.process.fail"), t);
        return false;
    }
    return true;
}

SocketProcessorBase

这个也是一个Runnable,我们看org.apache.tomcat.util.net.NioEndpoint.SocketProcessor#doRun

protected void doRun() {
    NioChannel socket = socketWrapper.getSocket();
    SelectionKey key = socket.getIOChannel().keyFor(socket.getPoller().getSelector());

    try {
        int handshake = -1;

        try {
            if (key != null) {
            	//tcp握手相关
                if (socket.isHandshakeComplete()) {
                    // No TLS handshaking required. Let the handler
                    // process this socket / event combination.
                    handshake = 0;
                } else if (event == SocketEvent.STOP || event == SocketEvent.DISCONNECT ||
                        event == SocketEvent.ERROR) {
                    // Unable to complete the TLS handshake. Treat it as
                    // if the handshake failed.
                    handshake = -1;
                } else {
                    handshake = socket.handshake(key.isReadable(), key.isWritable());
                    // The handshake process reads/writes from/to the
                    // socket. status may therefore be OPEN_WRITE once
                    // the handshake completes. However, the handshake
                    // happens when the socket is opened so the status
                    // must always be OPEN_READ after it completes. It
                    // is OK to always set this as it is only used if
                    // the handshake completes.
                    event = SocketEvent.OPEN_READ;
                }
            }
        } catch (IOException x) {
            handshake = -1;
            if (log.isDebugEnabled()) log.debug("Error during SSL handshake",x);
        } catch (CancelledKeyException ckx) {
            handshake = -1;
        }
        if (handshake == 0) {
            SocketState state = SocketState.OPEN;
            // Process the request from this socket
            if (event == null) {
            	//重要的是这段调用org.apache.coyote.AbstractProtocol.ConnectionHandler#process方法
            	//接着调用org.apache.coyote.AbstractProcessorLight#process
                state = getHandler().process(socketWrapper, SocketEvent.OPEN_READ);
            } else {
                state = getHandler().process(socketWrapper, event);
            }
            if (state == SocketState.CLOSED) {
                close(socket, key);
            }
        } else if (handshake == -1 ) {
            getHandler().process(socketWrapper, SocketEvent.CONNECT_FAIL);
            close(socket, key);
        } else if (handshake == SelectionKey.OP_READ){
            socketWrapper.registerReadInterest();
        } else if (handshake == SelectionKey.OP_WRITE){
            socketWrapper.registerWriteInterest();
        }
    } catch (CancelledKeyException cx) {
        socket.getPoller().cancelledKey(key);
    } catch (VirtualMachineError vme) {
        ExceptionUtils.handleThrowable(vme);
    } catch (Throwable t) {
        log.error("", t);
        socket.getPoller().cancelledKey(key);
    } finally {
        socketWrapper = null;
        event = null;
        //return to cache
        if (running && !paused) {
            processorCache.push(this);
        }
    }
}

Processor

public SocketState process(SocketWrapperBase<?> socketWrapper, SocketEvent status)
        throws IOException {

		//...
		//调用org.apache.coyote.http11.Http11Processor#service
		//将请求解析成http request和response
		//接着调用getAdapter().service(request, response);
		//即通过Adapter转交给Container处理request和response
        state = service(socketWrapper);
		//...
}

CoyoteAdapter

service

public void service(org.apache.coyote.Request req, org.apache.coyote.Response res)
            throws Exception {
	//关键代码就这一段。先调用org.apache.catalina.core.StandardEngineValve#invoke
	//然后org.apache.catalina.core.StandardHostValve#invoke
	//接着org.apache.catalina.core.StandardContextValve#invoke
	//最后org.apache.catalina.core.StandardWrapperValve#invoke
	connector.getService().getContainer().getPipeline().getFirst().invoke(
		                    request, response);
}

StandardWrapperValve

 public final void invoke(Request request, Response response)
    throws IOException, ServletException {


    //终于看到熟悉的Servlet了
    Servlet servlet = null;

        if (!unavailable) {
        	//通过反射创建servlet实例
        	//同时servlet.init(facade);
            servlet = wrapper.allocate();

    //将servlet和filter封装到一起
    ApplicationFilterChain filterChain =
            ApplicationFilterFactory.createFilterChain(request, wrapper, servlet);

	//同时调用servlet和filter
	//最终调用org.apache.catalina.core.ApplicationFilterChain#internalDoFilter
    filterChain.doFilter(request.getRequest(),
            response.getResponse());

	//放入Servlet池子中唤醒??不是调用destoy?
    wrapper.deallocate(servlet);

}

ApplicationFilterChain

  • internalDoFilter
private void internalDoFilter(ServletRequest request,
                              ServletResponse response)
    throws IOException, ServletException {

   	//调用filter
    filter.doFilter(request, response, this);


        	//终于调用到servlet的service方法
        	//会调用javax.servlet.http.HttpServlet#service(javax.servlet.ServletRequest, javax.servlet.ServletResponse)
        	//javax.servlet.http.HttpServlet#service(javax.servlet.http.HttpServletRequest, javax.servlet.http.HttpServletResponse)
        	//接着调用doXXX方法
            servlet.service(request, response);


}