Netty里面的Boss和Worker【Server篇】

#Netty里面的Boss和Worker【Server篇】 最近在总结Dubbo关于Netty通信方面的实现,于是也就借此机会深入体会了一下Netty。一般启动Netty的Server端时都会设置两个ExecutorService对象,我们都习惯用boss,worker两个变量来引用这两个对象,于是从我一开始接触Netty就有了boss和worker的概念。这篇博客将对boss和worker进行介绍,但并不是涉及Netty其他部分介绍。

在Netty的里面有一个Boss,他开了一家公司(开启一个服务端口)对外提供业务服务,它手下有一群做事情的workers。Boss一直对外宣传自己公司提供的业务,并且接受(accept)有需要的客户(client),当一位客户找到Boss说需要他公司提供的业务,Boss便会为这位客户安排一个worker,这个worker全程为这位客户服务(read/write)。如果公司业务繁忙,一个worker可能会为多个客户进行服务。这就是Netty里面Boss和worker之间的关系。下面看看Netty是如何让Boss和Worker进行协助的。

1<!--lang:java--> 2protected void doOpen() throws Throwable { 3 NettyHelper.setNettyLoggerFactory(); 4 ExecutorService boss = Executors.newCachedThreadPool(new NamedThreadFactory("NettyServerBoss", true)); 5 ExecutorService worker = Executors.newCachedThreadPool(new NamedThreadFactory("NettyServerWorker", true)); 6 ChannelFactory channelFactory = new NioServerSocketChannelFactory(boss, worker, getUrl().getPositiveParameter(Constants.IO_THREADS_KEY, Constants.DEFAULT_IO_THREADS)); 7 bootstrap = new ServerBootstrap(channelFactory); 8 9 final NettyHandler nettyHandler = new NettyHandler(getUrl(), this); 10 channels = nettyHandler.getChannels(); 11 bootstrap.setPipelineFactory(new ChannelPipelineFactory() { 12 public ChannelPipeline getPipeline() { 13 NettyCodecAdapter adapter = new NettyCodecAdapter(getCodec() ,getUrl(), NettyServer.this); 14 ChannelPipeline pipeline = Channels.pipeline(); 15 pipeline.addLast("decoder", adapter.getDecoder()); 16 pipeline.addLast("encoder", adapter.getEncoder()); 17 pipeline.addLast("handler", nettyHandler); 18 return pipeline; 19 } 20 }); 21 // bind 22 channel = bootstrap.bind(getBindAddress()); 23}

上面这段代码是Dubbo用来开启服务的,也是大部分使用Netty进行服务端开发常用的方式启动服务端。首先是设置boss和worker的线程池,以能够让它们在各自的线程池里面异步执行。当调用bootstrap.bind(getBindAddress())的时候最终受理绑定操作的是NioServerSocketPipelineSinkeventSunk方法,看类名和方法签名就应该知道是处理IO事件的。方法eventSunk实现如下:

1<!--lang:java--> 2public void eventSunk( 3 ChannelPipeline pipeline, ChannelEvent e) throws Exception { 4 Channel channel = e.getChannel(); 5 if (channel instanceof NioServerSocketChannel) { 6 handleServerSocket(e); 7 } else if (channel instanceof NioSocketChannel) { 8 handleAcceptedSocket(e); 9 } 10}

由于这个时候Server还处于bind阶段,所以channel肯定不是NioSocketChannel,于是就到了方法handleServerSocket里面,最后将会调用bind方法来绑定某个端口启动服务。下面是bind方法实现:

1<!--lang:java--> 2private void bind( 3 NioServerSocketChannel channel, ChannelFuture future, 4 SocketAddress localAddress) { 5 6 boolean bound = false; 7 boolean bossStarted = false; 8 try { 9 channel.socket.socket().bind(localAddress, channel.getConfig().getBacklog()); 10 bound = true; 11 12 future.setSuccess(); 13 fireChannelBound(channel, channel.getLocalAddress()); 14 15 Executor bossExecutor = 16 ((NioServerSocketChannelFactory) channel.getFactory()).bossExecutor; 17 DeadLockProofWorker.start( 18 bossExecutor, 19 new ThreadRenamingRunnable( 20 new Boss(channel), 21 "New I/O server boss #" + id + " (" + channel + ')')); 22 bossStarted = true; 23 } catch (Throwable t) { 24 future.setFailure(t); 25 fireExceptionCaught(channel, t); 26 } finally { 27 if (!bossStarted && bound) { 28 close(channel, future); 29 } 30 } 31}

可以看到socket的绑定以及设置异步的future成功,已通知服务启动成功,同时将绑定成功事件通知出去。接下来我看的重点来了,就是bossExecutor,可以看到它是通过NioServerSocketChannelFactory里面去获取的,NioServerSocketChannelFactory里面的boss就是之前我们设置进去的,可以确定我们之前设置boss的异步线程池是在这里被使用了。紧接下来的是启动我们的异步线程池,到这里进入了Boss该做的事情,Boss其实是实现了Runnable接口,从而可以交给boss的线程池运行,接下来的关注点就是Boss的run方法,这里才是Boss做事情的地方。再此之前先看看Boss初始化做了什么事情:

1<!--lang:java--> 2Boss(NioServerSocketChannel channel) throws IOException { 3 this.channel = channel; 4 5 selector = Selector.open(); 6 7 boolean registered = false; 8 try { 9 channel.socket.register(selector, SelectionKey.OP_ACCEPT); 10 registered = true; 11 } finally { 12 if (!registered) { 13 closeSelector(); 14 } 15 } 16 17 channel.selector = selector; 18 }

Boss初始化过程中其实就是将serversocket注册到一个selector里面,从而可以实现NIO的异步IO处理。

1<!--lang:java--> 2public void run() { 3 final Thread currentThread = Thread.currentThread(); 4 5 channel.shutdownLock.lock(); 6 try { 7 for (;;) { 8 try { 9 if (selector.select(1000) > 0) { 10 selector.selectedKeys().clear(); 11 } 12 13 SocketChannel acceptedSocket = channel.socket.accept(); 14 if (acceptedSocket != null) { 15 registerAcceptedChannel(acceptedSocket, currentThread); 16 } 17 } catch (SocketTimeoutException e) { 18 // Thrown every second to get ClosedChannelException 19 // raised. 20 } catch (CancelledKeyException e) { 21 // Raised by accept() when the server socket was closed. 22 } catch (ClosedSelectorException e) { 23 // Raised by accept() when the server socket was closed. 24 } catch (ClosedChannelException e) { 25 // Closed as requested. 26 break; 27 } catch (Throwable e) { 28 logger.warn( 29 "Failed to accept a connection.", e); 30 try { 31 Thread.sleep(1000); 32 } catch (InterruptedException e1) { 33 // Ignore 34 } 35 } 36 } 37 } finally { 38 channel.shutdownLock.unlock(); 39 closeSelector(); 40 } 41 }

run方法里面是一个死循环,里面在不间断的等待客户端的连接,如果有客户端的连接,那么将会调用方法registerAcceptedChannel进行后续的处理。

1<!--lang:java--> 2 private void registerAcceptedChannel(SocketChannel acceptedSocket, Thread currentThread) { 3 try { 4 ChannelPipeline pipeline = 5 channel.getConfig().getPipelineFactory().getPipeline(); 6 NioWorker worker = nextWorker(); 7 worker.register(new NioAcceptedSocketChannel( 8 channel.getFactory(), pipeline, channel, 9 NioServerSocketPipelineSink.this, acceptedSocket, 10 worker, currentThread), null); 11 } catch (Exception e) { 12 logger.warn( 13 "Failed to initialize an accepted socket.", e); 14 try { 15 acceptedSocket.close(); 16 } catch (IOException e2) { 17 logger.warn( 18 "Failed to close a partially accepted socket.", 19 e2); 20 } 21 } 22 }

方法registerAcceptedChannel就是将客户端的channle分配给一个worker,而这个worker是通过方法nextWorker获取 <!--lang:java--> NioWorker nextWorker() { return workers[Math.abs( workerIndex.getAndIncrement() % workers.length)]; }

可以看到方法nextWorker是一个让worker里面的客户端channel保持平衡的作用,可能你会疑问这个workers是哪里来的,其实是在上面初始化NioServerSocketChannelFactory的时候,NioServerSocketChannelFactory再去初始化NioServerSocketPipelineSink时候构造出来的,默认情况下workers的数量是我们初始化NioServerSocketChannelFactory设置进去的。可以看到是调用worker的register方法将客户端的channel注册到worker里面的。

1<!--lang:java--> 2void register(NioSocketChannel channel, ChannelFuture future) { 3 4 boolean server = !(channel instanceof NioClientSocketChannel); 5 Runnable registerTask = new RegisterTask(channel, future, server); 6 Selector selector; 7 8 synchronized (startStopLock) { 9 if (!started) { 10 ..... 11 this.selector = selector = Selector.open(); 12 ..... 13 DeadLockProofWorker.start( 14 executor, new ThreadRenamingRunnable(this, threadName)); 15 success = true; 16 ..... 17 } else { 18 selector = this.selector; 19 } 20 21 assert selector != null && selector.isOpen(); 22 23 started = true; 24 boolean offered = registerTaskQueue.offer(registerTask); 25 assert offered; 26 } 27 28 if (wakenUp.compareAndSet(false, true)) { 29 selector.wakeup(); 30 } 31}

上面对worker有一个started状态的检测,如果没启动,则启动worker,这个额一般都是将第一个客户端的channel注册到worker里面才进行的。由于worker也是实现了Rannable接口,所以启动的主要工作就是让worker在某个线程里面跑起来,并且为这个worker分配一个selector,用来进行监控IO事件。下面便是这个过程实现:

1<!--lang:java--> 2 DeadLockProofWorker.start( 3 executor, new ThreadRenamingRunnable(this, threadName)); 4 success = true;

其中的executor便是我们一开始设置的workerExecutor。 worker启动成功之后,接下来要做的便是让worker管理器客户端的channel

1<!--lang:java--> 2 Runnable registerTask = new RegisterTask(channel, future, server); 3 ....... 4 boolean offered = registerTaskQueue.offer(registerTask); 5 assert offered;

worker是将客户端包装成一个RegisterTask,然后放入队列,可见RegisterTask也实现了Runnable接口。那放入队列以后谁去取这个队列里面的数据呢?当然,肯定是worker去取。上面介绍启动worker的时候是让worker在某个线程里面跑起来,并且worker是实现了Rannable方法,于是运行worker的线程肯定是调用worker的run方法。

1<!--lang:java--> 2 public void run() { 3 thread = Thread.currentThread(); 4 boolean shutdown = false; 5 Selector selector = this.selector; 6 for (;;) { 7 ..... 8 try { 9 SelectorUtil.select(selector); 10 ..... 11 12 cancelledKeys = 0; 13 processRegisterTaskQueue(); 14 processWriteTaskQueue(); 15 processSelectedKeys(selector.selectedKeys()); 16 ..... 17 } catch (Throwable t) { 18 19 try { 20 Thread.sleep(1000); 21 } catch (InterruptedException e) { 22 23 } 24 } 25 } 26}

可以看到run方法里面也是一个死循环,在不断的轮询调用selector的select IO的事件。接下来会调用三个方法processRegisterTaskQueue,processWriteTaskQueueprocessSelectedKeys。通过方法签名就应该知道这个三个方法具体是做什么事情的,第一个是处理上面registerTaskQueue的,并且queue里面对象的run方法,而第二个processWriteTaskQueue是处理写任务的,而processSelectedKeys是处理selector匹配的IO事件。我们先看看registerTaskQueue是做了什么?

1<!--lang:java--> 2 private void processRegisterTaskQueue() throws IOException { 3 for (;;) { 4 final Runnable task = registerTaskQueue.poll(); 5 if (task == null) { 6 break; 7 } 8 9 task.run(); 10 cleanUpCancelledKeys(); 11 } 12}

上面介绍过registerTaskQueue里面的元素是RegisterTask。所以需要去看看RegisterTask的run方法实现,其中RegisterTaskNioWorker里面的内部类,所以RegisterTask是可以访问NioWorker的元素信息。

1<!--lang:java--> 2 public void run() { 3 SocketAddress localAddress = channel.getLocalAddress(); 4 SocketAddress remoteAddress = channel.getRemoteAddress(); 5 if (localAddress == null || remoteAddress == null) { 6 if (future != null) { 7 future.setFailure(new ClosedChannelException()); 8 } 9 close(channel, succeededFuture(channel)); 10 return; 11 } 12 13 try { 14 if (server) { 15 channel.socket.configureBlocking(false); 16 } 17 18 synchronized (channel.interestOpsLock) { 19 channel.socket.register( 20 selector, channel.getRawInterestOps(), channel); 21 } 22 if (future != null) { 23 channel.setConnected(); 24 future.setSuccess(); 25 } 26 } catch (IOException e) { 27 if (future != null) { 28 future.setFailure(e); 29 } 30 close(channel, succeededFuture(channel)); 31 .... 32 } 33 34 if (!server) { 35 if (!((NioClientSocketChannel) channel).boundManually) { 36 fireChannelBound(channel, localAddress); 37 } 38 fireChannelConnected(channel, remoteAddress); 39 } 40 }

可以看到这里面主要做的事情是将Boss分配给worker的客户端channel和worker的selector关联上,从而worker可以处理该客户端channel的IO事件。

到这里就完成了由Boss接收到一个客户端连接,到分配给某个worker,以及worker是怎么去和客户端的channel关联的,其中由于worker有可能为多个客户端channel服务,所以worker并不会直接和某个channel产生引用,而是将客户端的channel注册在该worker的selector上面,worker的run方法里面通过不断对selector的select轮询,以达到对channel进行处理。接下来看看worker怎么处理selector的io事件的

1<!--java:lang--> 2private void processSelectedKeys(Set<SelectionKey> selectedKeys) throws IOException { 3 for (Iterator<SelectionKey> i = selectedKeys.iterator(); i.hasNext();) { 4 SelectionKey k = i.next(); 5 i.remove(); 6 try { 7 int readyOps = k.readyOps(); 8 if ((readyOps & SelectionKey.OP_READ) != 0 || readyOps == 0) { 9 if (!read(k)) { 10 // Connection already closed - no need to handle write. 11 continue; 12 } 13 } 14 if ((readyOps & SelectionKey.OP_WRITE) != 0) { 15 writeFromSelectorLoop(k); 16 } 17 } catch (CancelledKeyException e) { 18 close(k); 19 } 20 21 if (cleanUpCancelledKeys()) { 22 break; // break the loop to avoid ConcurrentModificationException 23 } 24 } 25}

上面的方法完成的是处理selector产生的io事件,其中如果当前IO时间是读,那么将SelectionKey中的channel流进行读出,并且向上交给Netty的Handler。如果是当前某个channel的写满足条件,则触发writeFromSelectorLoop查看是否有待写出的内容。

对于写数据Netty在worker提供了三种入口

1<!--lang:java--> 2void writeFromUserCode(final NioSocketChannel channel) { 3 if (!channel.isConnected()) { 4 cleanUpWriteBuffer(channel); 5 return; 6 } 7 8 if (scheduleWriteIfNecessary(channel)) { 9 return; 10 } 11 12 if (channel.writeSuspended) { 13 return; 14 } 15 16 if (channel.inWriteNowLoop) { 17 return; 18 } 19 20 write0(channel); 21} 22 23void writeFromTaskLoop(final NioSocketChannel ch) { 24 if (!ch.writeSuspended) { 25 write0(ch); 26 } 27} 28 29void writeFromSelectorLoop(final SelectionKey k) { 30 NioSocketChannel ch = (NioSocketChannel) k.attachment(); 31 ch.writeSuspended = false; 32 write0(ch); 33}

其中writeFromUserCode是提供外部直接写出的,writeFromTaskLoop是在worker的run方法调用processWriteTaskQueue时候会触发。

点赞
收藏

评论区

加载中...

相关推荐

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

Reactor模式的.net版本简单实现

    近期在学习DotNetty,遇到不少的问题。由于dotnetty是次netty的.net版本的实现。导致在网上叙述dotnetty的原理,以及实现技巧方面的东西较少,这还是十分恼人的。在此建议学习和使用Dotnetty的和位小伙伴,真心阅读下netty的相关书籍,如《netty权威指南》。    闲话少说,进入正题。netty的性能之所以能够

Netty源码分析(二):服务端启动

上一篇粗略的介绍了一下netty,本篇将详细介绍Netty的服务器的启动过程。ServerBootstrap看过上篇事例的人,可以知道ServerBootstrap是Netty服务端启动中扮演着一个重要的角色。它是Netty提供的一个服务端引导类,继承自AbstractBootstrap。Serv

Netty创建服务器与客户端

Netty创建Server服务端Netty创建全部都是实现自AbstractBootstrap。客户端的是Bootstrap,服务端的则是ServerBootstrap。创建一个HelloServerpackageorg.examp

Netty之大名鼎鼎的EventLoop

EventLoopGroup与Reactor:前面的章节中我们已经知道了,一个Netty程序启动时,至少要指定一个EventLoopGroup(如果使用到的是NIO,通常是指NioEventLoopGroup),那么,这个NioEventLoopGroup在Netty中到底扮演着什么角色呢?我们知道,Netty是Reactor模型的