Netty学习——BootStrap

在这里插入图片描述
netty的BootStrap类是一个辅助类,其暴露出的接口有助于我们更方便的建立服务器和客户端,下面从源码角度分析Bootstrap是如何引导建立客户端的(使用断点的方式跟踪源码)

1、首先找到io.netty.example中的EchoClient,打上断点,Debug模式开始跟踪源码(断点如图)
在这里插入图片描述

2、b.group方法

class AbastactBootStrap

public B group(EventLoopGroup group) {
    if (group == null) {
        throw new NullPointerException("group");
    }
    if (this.group != null) {
        throw new IllegalStateException("group set already");
    }
    this.group = group;
    return self();
}

这里就有了一个重要的类AbastactBootStrap,其中包含了重要的属性和方法的实现
AbastactBootStrap
属性

//保存EventLoopGroup的引用
volatile EventLoopGroup group;
//保存channel的工厂类,用于创建channel实例
private volatile ChannelFactory<? extends C> channelFactory;
//handler的引用,多为ChannelInitializer的实现类,可添加多个handler
private volatile ChannelHandler handler;

方法之后再介绍,

3、b.channel(channel.class)

public B channel(Class<? extends C> channelClass) {
    if (channelClass == null) {
        throw new NullPointerException("channelClass");
    }
    return channelFactory(new ReflectiveChannelFactory<C>(channelClass));
}

public ReflectiveChannelFactory(Class<? extends T> clazz) {
    ObjectUtil.checkNotNull(clazz, "clazz");
    try {
        //重点在这个,将反射构造方法交给factory
        this.constructor = clazz.getConstructor();
    } catch (NoSuchMethodException e) {
        throw new IllegalArgumentException("Class " + StringUtil.simpleClassName(clazz) +
                " does not have a public non-arg constructor", e);
    }
}

4、其他就不一一写了,太多了,总之

 b.group(workgroup)
            .channel(NioSocketChannel.class)//客户端 -->NioSocketChannel
            .option(ChannelOption.SO_KEEPALIVE, true)
            .handler(new ChannelInitializer<SocketChannel>() {//handler
                @Override
                protected void initChannel(SocketChannel sc) throws Exception {
                    sc.pipeline().addLast(new ClientHandler());
                }
            });

这一串的方法级联调用就是 对AbastactBootStrap的属性进行赋值

5、然后追踪connect()方法,这是重头戏
在这里插入图片描述

走到 initAndRegister() 方法,这个方法是先初始化一个channel,然后再将这个channel注册到eventLoop上

final ChannelFuture initAndRegister() {
    //工厂类反射新建channel类,注意其中已经创建好了相应的pipeline
    Channel channel = channelFactory.newChannel();
    //初始化channel,将handler绑定到channel的pipeline中
    init(channel);
    //将channel注册到eventloop上
    ChannelFuture regFuture = config().group().register(channel);
}

channelFactory.newChannel() 我们就不跟了,就是使用反射构造方法创建实例

5.1 init(channel)方法

void init(Channel channel) throws Exception {
    ChannelPipeline p = channel.pipeline();
    //将handler加入到pipeline中
    p.addLast(config.handler());
}

其中pipeline的addLast()方法,我会在pipeline这一节详细将,这里就不做赘述了

5.2 register(channel)方法

public ChannelFuture register(Channel channel) {
    return next().register(channel);
}

public EventExecutor next() {
    return chooser.next();
}

public EventExecutor next() {
    return executors[idx.getAndIncrement() & executors.length - 1];
}

这里就是我们在EventLoopGroup中所讲的,EventLoopGroup是一个线程池,这里我们需要新建一个EventLoop时就直接从中取出一个就行了

然后就会追到这儿

class SingleThreadEventLoop

public ChannelFuture register(final ChannelPromise promise) {
    //channel调用了register方法将自己注册到eventloop上
    promise.channel().unsafe().register(this, promise);
    return promise;
}

class AbstractChannel 

public final void register(EventLoop eventLoop, final ChannelPromise promise) {

AbstractChannel.this.eventLoop = eventLoop;

if (eventLoop.inEventLoop()) {
    register0(promise);
} else {//追的是这一个分支,因为当前线程不是eventloop对应的线程

    //eventloop就是SingleThreadEventLoop,是一个单线程线程池
    //我们将register方法加入到线程池的执行队列中
    eventLoop.execute(new Runnable() {
        public void run() {
            register0(promise);
        }
    }); 
}
}

到这里就跟不下去了,因为注册的线程不是当前执行的线程,所以我们手动跟一下register0(promise)方法

protected void doRegister() throws Exception {
    boolean selected = false;
    for (;;) {
        try {
            selectionKey = javaChannel().register(eventLoop().unwrappedSelector(), 0, this);
            return;
        }
    }
}

//这里就跟到了package java.nio.channels.spi.AbstractSelectableChannel
//是Java的NIO的实现
public final SelectionKey register(Selector sel, int ops,
                                   Object att)
    throws ClosedChannelException
{
    synchronized (regLock) {
        if (!isOpen())
            throw new ClosedChannelException();
        if ((ops & ~validOps()) != 0)
            throw new IllegalArgumentException();
        if (blocking)
            throw new IllegalBlockingModeException();
        SelectionKey k = findKey(sel);
        if (k != null) {
            k.interestOps(ops);
            k.attach(att);
        }
        if (k == null) {
            // New registration
            synchronized (keyLock) {
                if (!isOpen())
                    throw new ClosedChannelException();
                k = ((AbstractSelector)sel).register(this, ops, att);
                addKey(k);
            }
        }
        return k;
    }
}

再往底层追register 的实现

class SelectorImpl
protected final SelectionKey register(AbstractSelectableChannel var1, int var2, Object var3) {
    if (!(var1 instanceof SelChImpl)) {
        throw new IllegalSelectorException();
    } else {
        SelectionKeyImpl var4 = new SelectionKeyImpl((SelChImpl)var1, this);
        var4.attach(var3);
        synchronized(this.publicKeys) {
            this.implRegister(var4);
        }

        var4.interestOps(var2);
        return var4;
    }
}

class WindowsSelectorImpl
//这里就是Java最后的封装了,再往底层就是将channel的文件描述符和感兴趣的事件写入到selector的轮询数组中
//由底层完成IO操作,返回结果
protected void implRegister(SelectionKeyImpl var1) {
    synchronized(this.closeLock) {
        if (this.pollWrapper == null) {
            throw new ClosedSelectorException();
        } else {
            this.growIfNeeded();
            this.channelArray[this.totalChannels] = var1;
            var1.setIndex(this.totalChannels);
            this.fdMap.put(var1);
            this.keys.add(var1);
            this.pollWrapper.addEntry(this.totalChannels, var1);
            ++this.totalChannels;
        }
    }
}

至此,一个完整的eventloop创建了出来,然后就是selector的死循环执行逻辑,直到底层操作完成
最后还有一个
ChannelFuture cf1 = b.connect(HOST, PORT1).sync();

这个sync()是一个同步阻塞等待结果的方法,我们跟一下其底层实现

public Promise<V> sync() throws InterruptedException {
    await();
    rethrowIfFailed();
    return this;
}

public Promise<V> await() throws InterruptedException 

    synchronized (this) {
        //这是一个死循环,直到promise中操作完成才返回
        while (!isDone()) {
            incWaiters();
            try {
                wait();
            } finally {
                decWaiters();
            }
        }
    }
    return this;
}

版权声明:本文为weixin_44601714原创文章,遵循CC 4.0 BY-SA版权协议,转载请附上原文出处链接和本声明。