WebSocket реализует мгновенную связь на стороне Интернета

Java

BST

предисловие

WebSocket — это протокол, предоставляемый HTML5 для полнодуплексной связи между браузерами и серверами. В настоящее время многие веб-приложения, которые не используют WebSocket для связи в реальном времени между клиентом и сервером, в основном используют опрос для установки регулярного времени или используют более длительный опрос для обработки push-сообщений в реальном времени. Это неизбежно приведет к значительной трате ресурсов сервера и полосы пропускания, и WebSocket, о котором мы сейчас поговорим, похоже, решает эту проблему, так что приложения в архитектуре B/S имеют те же возможности связи в реальном времени, что и C/. С архитектура.

Сравнение HTTP и WebSocket

HTTP

Протокол HTTP является полудуплексным протоколом, что означает, что одновременно может обрабатываться только одно направление передачи данных.В то же время сообщение HTTP слишком велико и содержит много данных заголовка сообщения.На самом деле, при обработке сообщений не требуется много данных, это также пустая трата ресурсов.

  • Голосование по времени: Опрос по времени означает, что клиент регулярно отправляет HTTP-запросы на сервер, чтобы узнать, есть ли данные.После получения запроса сервер возвращает данные клиенту, и соединение с ним закрывается. Эта реализация является самой простой, но будут задержки сообщений и много потраченных впустую ресурсов сервера и полосы пропускания.
  • долгий опрос: Долгий опрос, как и обычный опрос, также реализуется через HTTP-запросы, но это не синхронизированный запрос. Клиент отправляет запрос на сервер, в это время сервер удерживает запрос и возвращается к запрашивающему клиенту, когда есть данные или тайм-аут, и начинает следующий раунд запросов.

WebSocket

WebSocket нужен только один запрос между клиентом и сервером, а между клиентом и сервером устанавливается канал связи, который может передавать данные друг другу в режиме реального времени, и не несет много заголовков запросов и другой информации вроде HTTP . Поскольку WebSocket — это протокол, основанный на двунаправленной полнодуплексной связи TCP, он поддерживает обработку отправки и получения сообщений в один и тот же момент времени для обеспечения обработки сообщений в реальном времени.

  • Установить соединение WebSocket: чтобы установить соединение WebSocket, клиент должен сначала отправить специальный HTTP-запрос на сервер.Используемый протокол неhttpилиhttps, но используетсяwsилиwss(один незащищенный, один безопасный, аналогично разнице между первыми двумя), запрос информации об обновлении протокола должен быть прикреплен к заголовку запросаUpgrade: websocket, и случайным образом сгенерироватьSec-WebSocket-Keyзначение и информация о версииSec-WebSocket-Versionи Т. Д. После того, как сервер получит запрос клиента, он проанализирует информацию запроса, включая запрос на обновление протокола, проверку версии иSec-WebSocket-Keyпосле шифрованияsec-websocket-acceptЗначение возвращается клиенту, чтобы установить соединение между клиентом и сервером.
  • Закройте соединение WebSocket: и клиент, и сервер могут отправить кадр управления закрытием, а другой конец активно закрывает соединение.

Опрос HTTP и диаграмма жизненного цикла WebSocket

HTTP轮询和WebSocket生命周期示意图

Сервер

Здесь сервер разработан с использованием Netty WebSocket. Здесь сначала реализуется класс запуска сервера, а затем используется пользовательский процессор для обработки сообщений WebSocket.

package com.ytao.websocket;

import io.netty.bootstrap.ServerBootstrap;
import io.netty.channel.*;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioServerSocketChannel;
import io.netty.handler.codec.http.HttpObjectAggregator;
import io.netty.handler.codec.http.HttpServerCodec;
import io.netty.handler.logging.LogLevel;
import io.netty.handler.logging.LoggingHandler;
import io.netty.handler.stream.ChunkedWriteHandler;

/**
 * Created by YANGTAO on 2019/11/17 0017.
 */
public class WebSocketServer {

    public static String HOST = "127.0.0.1";
    public static int PORT = 8806;

    public static void startUp() throws Exception {
        // 监听端口的线程组
        EventLoopGroup bossGroup = new NioEventLoopGroup();
        // 处理每一条连接的数据读写的线程组
        EventLoopGroup workerGroup = new NioEventLoopGroup();
        // 启动的引导类
        ServerBootstrap serverBootstrap = new ServerBootstrap();
        try {
            serverBootstrap.group(bossGroup, workerGroup)
                    .channel(NioServerSocketChannel.class)
                    .childHandler(new ChannelInitializer<SocketChannel>() {
                        @Override
                        protected void initChannel(SocketChannel ch) throws Exception{
                            ChannelPipeline pipeline = ch.pipeline();
                            pipeline.addLast("logger", new LoggingHandler(LogLevel.INFO));
                            // 将请求和返回消息编码或解码成http
                            pipeline.addLast("http-codec", new HttpServerCodec());
                            // 使http的多个部分组合成一条完整的http
                            pipeline.addLast("aggregator", new HttpObjectAggregator(65536));
                            // 向客户端发送h5文件,主要是来支持websocket通信
                            pipeline.addLast("http-chunked", new ChunkedWriteHandler());
                            // 服务端自定义处理器
                            pipeline.addLast("handler", new WebSocketServerHandler());
                        }
                    })
                    // 开启心跳机制
                    .childOption(ChannelOption.SO_KEEPALIVE, true)
                    .handler(new ChannelInitializer<NioServerSocketChannel>() {
                        protected void initChannel(NioServerSocketChannel ch) {
                            System.out.println("WebSocket服务端启动中...");
                        }
                    });

            Channel ch = serverBootstrap.bind(HOST, PORT).sync().channel();
            System.out.println("WebSocket host: "+ch.localAddress().toString().replace("/",""));
            ch.closeFuture().sync();
        }catch (Exception e){
            e.printStackTrace();
        }finally {
            bossGroup.shutdownGracefully();
            workerGroup.shutdownGracefully();
        }

    }

    public static void main(String[] args) throws Exception {
        startUp();
    }
}

Приведенный выше класс запуска похож на протокол HTTP, поэтому его легче понять. После запуска класса запуска нам нужно обработать запрос WebSocket, который настраивается здесьWebSocketServerHandler. Бизнес-логика, которую мы разработали при обработке, заключается в том, что если есть только одно соединение для отправки информации в чат, мы автоматически ответим серверу, если их несколько, мы отправим информацию другим.

package com.ytao.websocket;

import io.netty.channel.Channel;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelPromise;
import io.netty.channel.SimpleChannelInboundHandler;
import io.netty.handler.codec.http.FullHttpRequest;
import io.netty.handler.codec.http.websocketx.*;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.Date;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;

/**
 * Created by YANGTAO on 2019/11/17 0017.
 */
public class WebSocketServerHandler extends SimpleChannelInboundHandler<Object> {

    private WebSocketServerHandshaker handshaker;

    private static Map<String, ChannelHandlerContext> channelHandlerContextConcurrentHashMap = new ConcurrentHashMap<>();

    private static final Map<String, String> replyMap = new ConcurrentHashMap<>();
    static {
        replyMap.put("博客", "https://ytao.top");
        replyMap.put("公众号", "ytao公众号");
        replyMap.put("在吗", "在");
        replyMap.put("吃饭了吗", "吃了");
        replyMap.put("你好", "你好");
        replyMap.put("谁", "ytao");
        replyMap.put("几点", "现在本地时间:"+LocalDateTime.now().format(DateTimeFormatter.ofPattern("HH:mm:ss")));
    }

    @Override
    public void messageReceived(ChannelHandlerContext channelHandlerContext, Object msg) throws Exception{
        channelHandlerContextConcurrentHashMap.put(channelHandlerContext.channel().toString(), channelHandlerContext);
        // http
        if (msg instanceof FullHttpRequest){
            handleHttpRequest(channelHandlerContext, (FullHttpRequest) msg);
        }else if (msg instanceof WebSocketFrame){ // WebSocket
            handleWebSocketFrame(channelHandlerContext, (WebSocketFrame) msg);
        }
    }

    @Override
    public void channelReadComplete(ChannelHandlerContext channelHandlerContext) throws Exception{
        if (channelHandlerContextConcurrentHashMap.size() > 1){
            for (String key : channelHandlerContextConcurrentHashMap.keySet()) {
                ChannelHandlerContext current = channelHandlerContextConcurrentHashMap.get(key);
                if (channelHandlerContext == current)
                    continue;
                current.flush();
            }
        }else {
            // 单条处理
            channelHandlerContext.flush();
        }
    }

    private void handleHttpRequest(ChannelHandlerContext channelHandlerContext, FullHttpRequest request) throws Exception{
        // 验证解码是否异常
        if (!"websocket".equals(request.headers().get("Upgrade")) || request.decoderResult().isFailure()){
            // todo send response bad
            System.err.println("解析http信息异常");
            return;
        }

        // 创建握手工厂类
        WebSocketServerHandshakerFactory factory = new WebSocketServerHandshakerFactory(
          "ws:/".concat(channelHandlerContext.channel().localAddress().toString()),
                null,
                false
        );
        handshaker = factory.newHandshaker(request);

        if (handshaker == null)
            WebSocketServerHandshakerFactory.sendUnsupportedVersionResponse(channelHandlerContext.channel());
        else
            // 响应握手消息给客户端
            handshaker.handshake(channelHandlerContext.channel(), request);

    }

    private void handleWebSocketFrame(ChannelHandlerContext channelHandlerContext, WebSocketFrame webSocketFrame){
        // 关闭链路
        if (webSocketFrame instanceof CloseWebSocketFrame){
            handshaker.close(channelHandlerContext.channel(), (CloseWebSocketFrame) webSocketFrame.retain());
            return;
        }

        // Ping消息
        if (webSocketFrame instanceof PingWebSocketFrame){
            channelHandlerContext.channel().write(
              new PongWebSocketFrame(webSocketFrame.content().retain())
            );
            return;
        }

        // Pong消息
        if (webSocketFrame instanceof PongWebSocketFrame){
            // todo Pong消息处理
        }

        // 二进制消息
        if (webSocketFrame instanceof BinaryWebSocketFrame){
            // todo 二进制消息处理
        }

        // 拆分数据
        if (webSocketFrame instanceof ContinuationWebSocketFrame){
            // todo 数据被拆分为多个websocketframe处理
        }

        // 文本信息处理
        if (webSocketFrame instanceof TextWebSocketFrame){
            // 推送过来的消息
            String  msg = ((TextWebSocketFrame) webSocketFrame).text();
            System.out.println(String.format("%s 收到消息 : %s", new Date(), msg));

            String responseMsg = "";
            if (channelHandlerContextConcurrentHashMap.size() > 1){
                responseMsg = msg;
                for (String key : channelHandlerContextConcurrentHashMap.keySet()) {
                    ChannelHandlerContext current = channelHandlerContextConcurrentHashMap.get(key);
                    if (channelHandlerContext == current)
                        continue;
                    Channel channel = current.channel();
                    channel.write(
                            new TextWebSocketFrame(responseMsg)
                    );
                }
            }else {
                // 自动回复
                responseMsg = this.answer(msg);
                if(responseMsg == null)
                    responseMsg = "暂时无法回答你的问题 ->_->";
                System.out.println("回复消息:"+responseMsg);
                Channel channel = channelHandlerContext.channel();
                channel.write(
                        new TextWebSocketFrame("【服务端】" + responseMsg)
                );
            }
        }

    }

    private String answer(String msg){
        for (String key : replyMap.keySet()) {
            if (msg.contains(key))
                return replyMap.get(key);
        }
        return null;
    }

    @Override
    public void exceptionCaught(ChannelHandlerContext channelHandlerContext, Throwable throwable){
        throwable.printStackTrace();
        channelHandlerContext.close();
    }

    @Override
    public void close(ChannelHandlerContext channelHandlerContext, ChannelPromise promise) throws Exception {
        channelHandlerContextConcurrentHashMap.remove(channelHandlerContext.channel().toString());
        channelHandlerContext.close(promise);
    }

}

Когда соединение только что установлено, первое рукопожатие обрабатывается протоколом HTTP, поэтомуWebSocketServerHandler#messageReceivedОн будет судить, является ли это HTTP или WebSocket.Если это HTTP, он будет переданWebSocketServerHandler#handleHttpRequestОбработка, которая проверит запрос и вернет сообщение клиенту после обработки рукопожатия. Если это не протокол HTTP, а протокол WebSocket, то обработка передается наWebSocketServerHandler#handleWebSocketFrameОбработка, после входа в обработку WebSocket происходит суждение о том, к какому типу относится сообщение, в том числеCloseWebSocketFrame,PingWebSocketFrame,PongWebSocketFrame,BinaryWebSocketFrame,ContinuationWebSocketFrame,TextWebSocketFrame,Они всеWebSocketFrameподкласс , иWebSocketFrameунаследовано отDefaultByteBufHolder.

channelHandlerContextConcurrentHashMapЭто кеширование подключенной информации WebSocket, потому что нам нужно записать количество подключений.Когда соединение закрывается, нам нужно удалить кешированное соединение, поэтому вWebSocketServerHandler#closeдля удаления кеша.

Последний текст, отправленный клиенту, оценивается по количеству подключений. Если количество подключений не больше 1, то мы «стоим 100 миллионов кода ядра ИИ».WebSocketServerHandler#answerдля ответа на сообщения клиентов. В противном случае, кроме соединения, полученного в этот раз, сообщение будет отправлено всем остальным подключенным клиентам.

клиент

Клиент использует JS для реализации операций WebSocket, и большинство основных браузеров в основном поддерживают WebSocket. Опора показана на рисунке:

支持WebSocket的浏览器

Реализация кода клиента H5:

<!DOCTYPE html>
<html lang="en">
<head>
    <meta charset="UTF-8">
    <meta name="viewport" content="width=device-width, initial-scale=1" />
    <title>ytao-websocket</title>
    <script src="http://libs.baidu.com/jquery/2.0.0/jquery.min.js"></script>
    <style type="text/css">
        #msgContent{
            line-height:200%;
            width: 500px;
            height: 300px;
            resize: none;
            border-color: #FF9900;
        }
        .clean{
            background-color: white;
        }
        .send{
            border-radius: 10%;
            background-color: #2BD56F;
        }
        @media screen and (max-width: 600px) {
            #msgContent{
                line-height:200%;
                width: 100%;
                height: 300px;
            }
        }
    </style>
</head>
<script>
    var socket;
    var URL = "ws://127.0.0.1:8806/ytao";

    connect();

    function connect() {
        $("#status").html("<span>连接中.....</span>");
        window.WebSocket = !window.WebSocket == true? window.MozWebSocket : window.WebSocket;
        if(window.WebSocket){
            socket = new WebSocket(URL);
            socket.onmessage = function(event){
                var msg = event.data + "\n";
                addMsgContent(msg);
            };

            socket.onopen = function(){
                $("#status").html("<span style='background-color: #44b549'>WebSocket已连接</span>");
            };

            socket.onclose = function(){
                $("#status").html("<span style='background-color: red'>WebSocket已断开连接</span>");
                setTimeout("connect()", 3000);
            };
        }else{
            $("#status").html("<span style='background-color: red'>该浏览器不支持WebSocket协议!</span>");
        }
    }

    function addMsgContent(msg) {
        var contet = $("#msgContent").val() + msg;
        $("#msgContent").val(contet)
    }

    function clean() {
        $("#msgContent").val("");
    }

    function getUserName() {
        var n = $("input[name=userName]").val();
        if (n == "")
            n = "匿名";
        return n;
    }

    function send(){
        var message = $("input[name=message]").val();
        if(!window.WebSocket) return;
        if ($.trim(message) == ""){
            alert("不能发送空消息!");
            return;
        }
        if(socket.readyState == WebSocket.OPEN){
            var msg = "【我】" + message + "\n";
            this.addMsgContent(msg);
            socket.send("【"+getUserName()+"】"+message);
            $("input[name=message]").val("");
        }else{
            alert("无法建立WebSocket连接!");
        }
    }

    $(document).keyup(function(){
        if(event.keyCode ==13){
            send()
        }
    });
</script>
<body>
    <div style="text-align: center;">
        <div id="status">
            <span>连接中.....</span>
        </div>
        <div>
            <h2>信息面板</h2>
            <textarea id="msgContent" readonly="readonly"></textarea>
        </div>
        <div>
            <input class="clean" type="button" value="清除聊天纪录" onclick="clean()" />
            <input type="text" name="userName" value="" placeholder="用户名"/>
        </div>
        <hr>
        <div>
            <form onsubmit="return false">
                <input type="text" name="message" value="" placeholder="请输入消息"/>
                <input class="send" type="button" name="msgBtn" value="send" onclick="send()"/>
            </form>
        </div>
        <div>
            <br><br>
            <img src="http://yangtao.ytao.top/ytao%E5%85%AC%E4%BC%97%E5%8F%B7.jpg">
        </div>
    </div>
</body>
</html>

JS здесь относительно просто реализовать, в основном используется:

  • new WebSocket(URL)Создайте объект веб-сокета
  • onopen()открытое соединение
  • onclose()закрыть соединение
  • onmessageполучить сообщение
  • send()Отправить сообщение

Когда соединение разрывается, клиентская сторона повторно инициирует соединение до тех пор, пока соединение не будет установлено успешно.

запускать

После того, как клиент и сервер подключены, мы можем увидеть информацию о проверке, упомянутую выше, из журнала и запроса. Клиент:

客户端连接信息

Сервер:

服务端连接信息

После запуска сервера сначала поэкспериментируйте с нашим «100 миллионов ИИ», когда подключен только один пользователь, результат отправки информации показан на рисунке:

Несколько пользовательских подключений, здесь три подключенных пользователя используются для группового чата.

Пользователь один:

Пользователь два:

Пользователь третий:

До сих пор WebSocket помогал нам реализовать потребности в мгновенном общении, и я считаю, что все в основном начали использовать WebSocket.

Суммировать

Понимание этой статьи может помочь вам приступить к работе с WebSocket и решить некоторые текущие проблемы со связью через Интернет, которые могут существовать. Я также видел в двух проектах, что этот тип решения выполняется через регулярный опрос, что также приводит к определенной трате ресурсов сервера в большей или меньшей степени. Поскольку WebSocket сам по себе является более сложным и предоставляет больше API-интерфейсов, его не очень просто реально использовать или оптимизировать в процессе использования, и это нужно делать в соответствии с реальной сценой.


личный блог: ytao.top

Мой официальный аккаунт ytao

我的公众号