一、前言
最近看了netty原始碼,打算寫個博客記下來,方便后面再復習,同時希望也能方便看到的人,在研究netty的時候,多少能方便點,
二、環境搭建
git clone netty的代碼下來,或者可以fork到自己的git 倉庫,然后git clone下來,
后面的版本統一用
<dependency>
<groupId>io.netty</groupId>
<artifactId>netty-all</artifactId>
<version>4.1.6.Final</version>
</dependency>
三、例子研究
如下是服務端標準的代碼案例,bossgroup主要是用來接收連接請求的,workergroup主要是用來處理讀寫請求的
1 EventLoopGroup bossGroup = new NioEventLoopGroup(1); 2 EventLoopGroup workerGroup = new NioEventLoopGroup(); 3 final EchoServerHandler serverHandler = new EchoServerHandler(); 4 try { 5 ServerBootstrap b = new ServerBootstrap(); 6 b.group(bossGroup, workerGroup) 7 .channel(NioServerSocketChannel.class) 8 .option(ChannelOption.SO_BACKLOG, 100) 9 .handler(new LoggingHandler(LogLevel.INFO)) 10 .childHandler(new ChannelInitializer<SocketChannel>() { 11 @Override 12 public void initChannel(SocketChannel ch) throws Exception { 13 ChannelPipeline p = ch.pipeline(); 14 if (sslCtx != null) { 15 p.addLast(sslCtx.newHandler(ch.alloc())); 16 } 17 //p.addLast(new LoggingHandler(LogLevel.INFO)); 18 p.addLast(serverHandler); 19 } 20 }); 21 22 // Start the server. 23 ChannelFuture f = b.bind(PORT).sync();
前面5-20都是初始化,我們先看23行,bind方法,一路跟下去,分為三部分
1 private ChannelFuture doBind(final SocketAddress localAddress) { 2 final ChannelFuture regFuture = initAndRegister(); 3 final Channel channel = regFuture.channel(); 4 if (regFuture.cause() != null) { 5 return regFuture; 6 } 7 8 if (regFuture.isDone()) { 9 // At this point we know that the registration was complete and successful. 10 ChannelPromise promise = channel.newPromise(); 11 doBind0(regFuture, channel, localAddress, promise); 12 return promise; 13 } else { 14 // Registration future is almost always fulfilled already, but just in case it's not. 15 final PendingRegistrationPromise promise = new PendingRegistrationPromise(channel); 16 regFuture.addListener(new ChannelFutureListener() { 17 @Override 18 public void operationComplete(ChannelFuture future) throws Exception { 19 Throwable cause = future.cause(); 20 if (cause != null) { 21 // Registration on the EventLoop failed so fail the ChannelPromise directly to not cause an 22 // IllegalStateException once we try to access the EventLoop of the Channel. 23 promise.setFailure(cause); 24 } else { 25 // Registration was successful, so set the correct executor to use. 26 // See https://github.com/netty/netty/issues/2586 27 promise.registered(); 28 29 doBind0(regFuture, channel, localAddress, promise); 30 } 31 } 32 }); 33 return promise; 34 } 35 }
第一是 initAndRegister方法

在這里,channelFactory啥時候初始化的?我們回到標準案例那里

跟進去看看
public B channel(Class<? extends C> channelClass) { if (channelClass == null) { throw new NullPointerException("channelClass"); } return channelFactory(new ReflectiveChannelFactory<C>(channelClass)); }
跟到底會發現下面這段
public B channelFactory(ChannelFactory<? extends C> channelFactory) { if (channelFactory == null) { throw new NullPointerException("channelFactory"); } if (this.channelFactory != null) { throw new IllegalStateException("channelFactory set already"); } //初始化cannelFactory this.channelFactory = channelFactory; return self(); }
所以很明顯,cannelFactory是ReflectiveChannelFactory,我們繼續看ReflectiveChannelFactory的newChannel方法
public T newChannel() { try { return clazz.getConstructor().newInstance(); } catch (Throwable t) { throw new ChannelException("Unable to create Channel from class " + clazz, t); } }
這就可以看出,channel的創建是通過工廠模式,反射創建無參建構式的,實體就是我們初始化傳進去的 NioServerSocketChannel,我們把channel的創建看完,繼續跟它的建構式
private static final SelectorProvider DEFAULT_SELECTOR_PROVIDER = SelectorProvider.provider();
public NioServerSocketChannel() { this(newSocket(DEFAULT_SELECTOR_PROVIDER)); }
DEFAULT_SELECTOR_PROVIDER 是根據作業系統選擇的provider,而newSocket其實就是根據provider到jdk底層去獲取對應的serversockerchannel,我們繼續this,
public NioServerSocketChannel(ServerSocketChannel channel) { super(null, channel, SelectionKey.OP_ACCEPT); config = new NioServerSocketChannelConfig(this, javaChannel().socket()); }
呼叫父類的構造方法,注意引數 SelectionKey.OP_ACCEPT,繼續
protected AbstractNioChannel(Channel parent, SelectableChannel ch, int readInterestOp) { super(parent); this.ch = ch; this.readInterestOp = readInterestOp; try { ch.configureBlocking(false); } catch (IOException e) { try { ch.close(); } catch (IOException e2) { if (logger.isWarnEnabled()) { logger.warn( "Failed to close a partially initialized socket.", e2); } } throw new ChannelException("Failed to enter non-blocking mode.", e); } }
這個建構式就是把前面生成的channel和accept事件保存起來,并設定該channel為非阻塞模式,是不是就是nio的代碼方式,我們繼續看super
protected AbstractChannel(Channel parent) { this.parent = parent; id = newId(); unsafe = newUnsafe(); pipeline = newChannelPipeline(); }
id我們暫時不管,這里會生成一個unsafe 來操作bytebuffer的,還生成了pipeline,這個主要是為了執行我們初始化設定的一些handdler,我們后面分析;到這里把channel的初始化分析完了,回到之前的initAndRegister方法,我們繼續往下看有個init方法,它有兩個實作,一個是客戶端的BootStrap,一個是服務端的ServerBootStrap,做的事情都差不多,我們看下ServerBootStrap的,
@Override void init(Channel channel) throws Exception { // 把代碼啟動的時候設定的引數放到它該有的位置上 final Map<ChannelOption<?>, Object> options = options0(); synchronized (options) { // options設定到channel上 setChannelOptions(channel, options, logger); } final Map<AttributeKey<?>, Object> attrs = attrs0(); synchronized (attrs) { // 遍歷attr事件,設定到channel上 for (Entry<AttributeKey<?>, Object> e: attrs.entrySet()) { @SuppressWarnings("unchecked") AttributeKey<Object> key = (AttributeKey<Object>) e.getKey(); channel.attr(key).set(e.getValue()); } } ChannelPipeline p = channel.pipeline(); ...... // 把所有handler組裝成pipeline p.addLast(new ChannelInitializer<Channel>() { @Override public void initChannel(final Channel ch) throws Exception { final ChannelPipeline pipeline = ch.pipeline(); ChannelHandler handler = config.handler(); if (handler != null) { pipeline.addLast(handler); } ch.eventLoop().execute(new Runnable() { @Override public void run() { pipeline.addLast(new ServerBootstrapAcceptor( ch, currentChildGroup, currentChildHandler, currentChildOptions, currentChildAttrs)); } }); } }); }
省略了些代碼,主要還是把之前一開始初始化保存的物件綁到對應的channel上,然后放到一個 inboundhandler型別-ServerBootstrapAcceptor物件上,并給放到pipeline 鏈上,
我們繼續看initAndRegister的另一行代碼
ChannelFuture regFuture = config().group().register(channel);
config().group(),在這里是 NioEventLoopGroup,register執行的是它的父類 MultithreadEventLoopGroup
public ChannelFuture register(Channel channel) { return next().register(channel); }
next方法有兩種實作,在NioEventLoopGroup初始化的時候會呼叫它的父類建構式,如果執行緒數是2的次方就實體化 PowerOfTwoEventExecutorChooser,否則就是GenericEventExecutorChooser,我們看看兩個實作的有啥區別
PowerOfTwoEventExecutorChooser: public EventExecutor next() { return executors[idx.getAndIncrement() & executors.length - 1]; } GenericEventExecutorChooser: public EventExecutor next() { return executors[Math.abs(idx.getAndIncrement() % executors.length)]; }
一個按位與,一個是取模運算,明顯按位與快一點,所以推薦設定2的n次方
講完chooser選擇后,繼續看register,因為我們是 NioEventLoop,它繼承于 SingleThreadEventLoop,所以我們看它的register
public ChannelFuture register(final ChannelPromise promise) { ObjectUtil.checkNotNull(promise, "promise"); promise.channel().unsafe().register(this, promise); return promise; }
從DefaultChannelPromise 拿到niosocketchannel,拿到對應的unsafe,而 AbstractUnsafe 是 AbstractChannel的內部類,
public final void register(EventLoop eventLoop, final ChannelPromise promise) { ...... AbstractChannel.this.eventLoop = eventLoop; if (eventLoop.inEventLoop()) { register0(promise); } else { try { eventLoop.execute(new Runnable() { @Override public void run() { register0(promise); } }); } } }
繼續看 register0
private void register0(ChannelPromise promise) { ....... doRegister(); ....... // Ensure we call handlerAdded(...) before we actually notify the promise. This is needed as the // user may already fire events through the pipeline in the ChannelFutureListener. pipeline.invokeHandlerAddedIfNeeded(); safeSetSuccess(promise); pipeline.fireChannelRegistered(); // Only fire a channelActive if the channel has never been registered. This prevents firing // multiple channel actives if the channel is deregistered and re-registered. ....... }
繼續 doRegister
protected void doRegister() throws Exception { boolean selected = false; for (;;) { try { selectionKey = javaChannel().register(eventLoop().unwrappedSelector(), 0, this); return; } } }
呼叫底層的jdk channel來注冊selector,拿到一個selectionKey,initAndRegister方法算是結束了,主要是初始化channel和注冊selector的,接下來看看 doBind0 方法,
private static void doBind0( final ChannelFuture regFuture, final Channel channel, final SocketAddress localAddress, final ChannelPromise promise) { // This method is invoked before channelRegistered() is triggered. Give user handlers a chance to set up // the pipeline in its channelRegistered() implementation. channel.eventLoop().execute(new Runnable() { @Override public void run() { if (regFuture.isSuccess()) { // 通過pipeline 從tail到head節點執行,最終在對應的NioServerSocketChannel執行bind方法 channel.bind(localAddress, promise).addListener(ChannelFutureListener.CLOSE_ON_FAILURE); } else { promise.setFailure(regFuture.cause()); } } }); }
我們繼續往下跟會發現,扔給執行緒池的任務異步執行,從AbstractChannel開始做bind操作,通過每個channel對應的 DefaultChannelPipeline 來執行bind,最侄訓通過 pipeline 從tail執行到head節點,最終跑到NioServerSocketChannel 類
protected void doBind(SocketAddress localAddress) throws Exception { // 最終執行到jdk底層的 bind方法 if (PlatformDependent.javaVersion() >= 7) { javaChannel().bind(localAddress, config.getBacklog()); } else { javaChannel().socket().bind(localAddress, config.getBacklog()); } }
最終bind 方法執行完畢,
轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/568.html
標籤:其他
