你好,我是彤哥,本篇是 netty 系列的第四篇。
欢迎来我的公从号 彤哥读源码 系统地学习 源码 & 架构 的知识。
简介
上一章我们一起学习了 Java 中的 BIO/NIO/AIO 的故事,本章将带着大家一起使用纯纯的 NIO 实现一个越聊越上瘾的“群聊系统”。
业务逻辑分析
首先,我们先来分析一下群聊的功能点:
(1)加入群聊,并通知其他人;
(2)发言,并通知其他人;
(3)退出群聊,并通知其他人;
一个简单的群聊系统差不多这三个功能足够了,为了方便记录用户信息,当用户加入群聊的时候自动给他分配一个用户 ID。
业务实现
上代码:
// 这是一个内部类
private static class ChatHolder {
// 我们只用了一个线程,用普通的 HashMap 也可以
static final Map<SocketChannel, String> USER_MAP = new ConcurrentHashMap<>();
/**
* 加入群聊
* @param socketChannel
*/
static void join(SocketChannel socketChannel) {
// 有人加入就给他分配一个 id,本文来源于公从号“彤哥读源码”String userId = "用户"+ ThreadLocalRandom.current().nextInt(Integer.MAX_VALUE);
send(socketChannel, "您的 id 为:" + userId + "\n\r");
for (SocketChannel channel : USER_MAP.keySet()) {send(channel, userId + "加入了群聊" + "\n\r");
}
// 将当前用户加入到 map 中
USER_MAP.put(socketChannel, userId);
}
/**
* 退出群聊
* @param socketChannel
*/
static void quit(SocketChannel socketChannel) {String userId = USER_MAP.get(socketChannel);
send(socketChannel, "您退出了群聊" + "\n\r");
USER_MAP.remove(socketChannel);
for (SocketChannel channel : USER_MAP.keySet()) {if (channel != socketChannel) {send(channel, userId + "退出了群聊" + "\n\r");
}
}
}
/**
* 扩散说话的内容
* @param socketChannel
* @param content
*/
public static void propagate(SocketChannel socketChannel, String content) {String userId = USER_MAP.get(socketChannel);
for (SocketChannel channel : USER_MAP.keySet()) {if (channel != socketChannel) {send(channel, userId + ":" + content + "\n\r");
}
}
}
/**
* 发送消息
* @param socketChannel
* @param msg
*/
static void send(SocketChannel socketChannel, String msg) {
try {ByteBuffer writeBuffer = ByteBuffer.allocate(1024);
writeBuffer.put(msg.getBytes());
writeBuffer.flip();
socketChannel.write(writeBuffer);
} catch (Exception e) {e.printStackTrace();
}
}
}
服务端代码
服务端代码直接使用上一章 NIO 的实现,只不过这里要把上面实现的业务逻辑适时地插入到相应的事件中。
(1)accept 事件,即连接建立的时候,说明加入了群聊;
(2)read 事件,即读取数据的时候,说明有人说话了;
(3)连接断开的时候,说明退出了群聊;
OK,直接上代码,为了与上一章的代码作区分,彤哥特意加入了一些标记:
public class ChatServer {public static void main(String[] args) throws IOException {Selector selector = Selector.open();
ServerSocketChannel serverSocketChannel = ServerSocketChannel.open();
serverSocketChannel.bind(new InetSocketAddress(8080));
serverSocketChannel.configureBlocking(false);
// 将 accept 事件绑定到 selector 上
serverSocketChannel.register(selector, SelectionKey.OP_ACCEPT);
while (true) {
// 阻塞在 select 上
selector.select();
Set<SelectionKey> selectionKeys = selector.selectedKeys();
// 遍历 selectKeys
Iterator<SelectionKey> iterator = selectionKeys.iterator();
while (iterator.hasNext()) {SelectionKey selectionKey = iterator.next();
// 如果是 accept 事件
if (selectionKey.isAcceptable()) {ServerSocketChannel ssc = (ServerSocketChannel) selectionKey.channel();
SocketChannel socketChannel = ssc.accept();
System.out.println("accept new conn:" + socketChannel.getRemoteAddress());
socketChannel.configureBlocking(false);
socketChannel.register(selector, SelectionKey.OP_READ);
// 加入群聊,本文来源于公从号“彤哥读源码”ChatHolder.join(socketChannel);
} else if (selectionKey.isReadable()) {
// 如果是读取事件
SocketChannel socketChannel = (SocketChannel) selectionKey.channel();
ByteBuffer buffer = ByteBuffer.allocate(1024);
// 将数据读入到 buffer 中
int length = socketChannel.read(buffer);
if (length > 0) {buffer.flip();
byte[] bytes = new byte[buffer.remaining()];
// 将数据读入到 byte 数组中
buffer.get(bytes);
// 换行符会跟着消息一起传过来
String content = new String(bytes, "UTF-8").replace("\r\n", "");
if (content.equalsIgnoreCase("quit")) {
// 退出群聊,本文来源于公从号“彤哥读源码”ChatHolder.quit(socketChannel);
selectionKey.cancel();
socketChannel.close();} else {
// 扩散,本文来源于公从号“彤哥读源码”ChatHolder.propagate(socketChannel, content);
}
}
}
iterator.remove();}
}
}
}
测试
打开四个 XSHELL 客户端,分别连接telnet 127.0.0.1 8080
,然后就可以开始群聊了。
彤哥发现,自己跟自己聊天也是会上瘾的,完全停不下来,不行了,我再去自聊一会儿 ^^
总结
本文彤哥跟着大家一起实现了“群聊系统”,去掉注释也就 100 行左右的代码,是不是非常简单?这就是 NIO 网络编程的魅力,我发现写网络编程也上瘾了 ^^
问题
这两章我们都没有用 NIO 实现客户端,你知道怎么实现吗?
提示:服务端需要监听 accept 事件,所以需要有一个 ServerSocketChannel,而客户端是直接去连服务器了,所以直接用 SocketChannel 就可以了,一个 SocketChannel 就相当于一个 Connection。
最后,也欢迎来我的公从号 彤哥读源码 系统地学习 源码 & 架构 的知识。