ServerSocketChannel.open()创建一个服务器的socket,并且建立管道Channel。这个方式跟我们一个开始使用new ServerSocket(8080);没有太大区别。
public class NioTcpDemo01 {
public static void main(String[] args) throws Exception{
ServerSocketChannel ssc = ServerSocketChannel.open();
ssc.bind(new InetSocketAddress(8080));
SocketChannel accept = ssc.accept();
System.out.println(accept);
}
}
⬇接着是接收客户端发来的消息,跟前面没有用Channel是一个性质的。
public class NioTcpDemo01 {
public static void main(String[] args) throws Exception{
ByteBuffer buffer = ByteBuffer.allocate(1024);
ServerSocketChannel ssc = ServerSocketChannel.open();
ssc.bind(new InetSocketAddress(8080));
SocketChannel accept = ssc.accept();
System.out.println(accept);
accept.read(buffer);
System.out.println(new String(buffer.array()));
}
}
⬇设置为非阻塞模式,充分利用系统的函数进行优化
public class NioTcpDemo01 {
//设置为非阻塞
public static void main(String[] args) throws Exception {
ServerSocketChannel ssc = ServerSocketChannel.open();
ssc.bind(new InetSocketAddress(8080));
ssc.configureBlocking(false);
while(true){
SocketChannel accept = ssc.accept();//重点为非阻塞
if(accept != null)
System.out.println(accept);
}
}
}
⬆设置为非阻塞之后accept()不再阻塞,如果没有连接进来,会返回null。通过主线程while循环进行不断轮询可以知道是否有连接进来。
将SocketChannel也设置为非阻塞。即当客户端连接进来时不阻塞read事件.
public static void main(String[] args) throws Exception {
ByteBuffer buffer = ByteBuffer.allocate(1024);
ServerSocketChannel ssc = ServerSocketChannel.open();
ssc.bind(new InetSocketAddress(8080));
ssc.configureBlocking(false);
List<SocketChannel> socketChannels = new ArrayList<>();
while(true){
SocketChannel accept = ssc.accept();//重点为非阻塞
if(accept != null){
System.out.println(accept);
accept.configureBlocking(false); //设置SocketChannel为非阻塞
socketChannels.add(accept);
}
for (SocketChannel socketChannel : socketChannels){
int read = socketChannel.read(buffer);//这里不阻塞了,如果读到数据,则返回数据长度,否则返回0
if (read > 0){
buffer.flip();
byte[] data = new byte[read];
buffer.get(data);
System.out.println(new String(data));
buffer.clear();
socketChannel.write(ByteBuffer.wrap(data));
}
if (read == -1) { //read==-1就断开连接了
// 客户端断开连接
System.out.println("客户端断开连接: " + socketChannel);
socketChannel.close();
socketChannels.remove(socketChannel);
continue;
}
}
}
}
selector和channel建立连接
public class NioTcpDemo02 {
public static void main(String[] args) throws Exception {
ServerSocketChannel ssc = ServerSocketChannel.open();
ssc.bind(new InetSocketAddress(8080));
ssc.configureBlocking(false);//设置为非阻塞
//创建多路复用
Selector selector = Selector.open();
//ssc注册到多路复用器,监听连接事件,selector管理和监听这个Channel
ssc.register(selector, SelectionKey.OP_ACCEPT);
while(true){
int select = selector.select(); //阻塞,等待事件触发
if(select > 0){
Set<SelectionKey> selectionKeys = selector.selectedKeys(); //获取所有事件
Iterator<SelectionKey> iterator = selectionKeys.iterator();
while (iterator.hasNext()){
SelectionKey key = iterator.next();
if(key.isAcceptable()){//判断是否是连接事件
//连接事件的channel必定是ServerSocketChannel
ServerSocketChannel channel = (ServerSocketChannel) key.channel();
//接受连接
SocketChannel socketChannel = channel.accept();
System.out.println(socketChannel);
socketChannel.configureBlocking(false);
//socketChannel也交给Selector管理和监听
socketChannel.register(selector,SelectionKey.OP_READ);
System.out.println("连接成功");
iterator.remove(); // 时间处理完成移除事件
}
if(key.isReadable()){
SocketChannel socketChannel = (SocketChannel) key.channel();
byte[] buffer = new byte[1024];
int read = socketChannel.read(ByteBuffer.wrap(buffer));
if(read > 0){
System.out.println(new String(buffer));
}
socketChannel.write(ByteBuffer.wrap(buffer));
iterator.remove();
}
}
}
}
}
}
boss-worker
但实际上我们只用到了一个线程,一个selector,不能充分使用多核cpu的优势。现在我们创建多个selector。我们规定boss-selector只负责建立连接,worker-selector负责业务处理。每次连接创建一个线程,这个线程中创建worker-selector,这个线程中的selector负责处理其他事件。而且规定创建的worker数量。
连接事件->boss
读写事件->worker
假设worker最大数量=4个
当一个连接进来,我们创建一个worker,有4个连接进来,创建4个worker后,第5个连接进来时会复用原来的worker