Netty实现WebSocket及分布式解决方案

2024-09-02 08:52

本文主要是介绍Netty实现WebSocket及分布式解决方案,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

在项目中面临了两个关键需求:一是实时数据获取,二是轻量级的即时通讯功能。传统的轮询机制虽然也能够从服务器获取数据,但它存在明显的不足:首先,它无法实现真正的实时性,;其次,频繁的请求会占用宝贵的客户端资源,影响用户体验,并增加服务器的负载。

针对这些挑战,我们选择了WebSocket协议作为解决方案。WebSocket提供了一种持久的连接方式,允许服务器主动、即时地向客户端推送最新数据,从而确保了信息的实时更新。这种全双工通信机制不仅提高了数据传输的效率,还显著降低了客户端的资源消耗。同时通过WebSocket,我们能够构建一个轻量级、响应迅速的通讯平台,简单的实现了用户间的互动。

WebSocket与Http对比

  • 协议类型
    • HTTP:是一个无状态、请求/响应式的协议,通常用于客户端(如浏览器)和服务器之间的一次性事务处理。它是基于TCP的文本协议,用于分布式、协作式、超媒体信息系统。
    • WebSocket:是一个持久化的协议,提供了全双工通信机制。它允许服务器主动向客户端发送消息,而不需要客户端再次发起请求。WebSocket 也是基于TCP的,但它在建立连接后可以保持长连接状态。
  • 头部开销:
    • HTTP:每次请求和响应都需要携带完整的头部信息,这在频繁通信时会增加数据传输的开销。
    • WebSocket:一旦连接建立,后续的消息传输不需要携带HTTP头部,只有较小的数据包头,这减少了通信的开销。
  • 安全性:
    • HTTP:可以通过HTTPS来提供加密的传输。
    • WebSocket:也有对应的安全版本称为WebSocket Secure(WSS),它在WebSocket的基础上增加了TLS/SSL加密。
  • 握手过程:
    • HTTP:客户端发送请求,服务器响应请求,然后连接关闭或保持(取决于是否keep-alive)。
    • WebSocket:客户端发送一个特殊的Upgrade请求来初始化WebSocket连接,服务器响应并升级连接为WebSocket连接。
WebSocket握手过程示例
ws://ip:port/ws
请求方法: GET
状态代码: 101 Switching Protocols

请求报文

Request:
GET ws://ip:port/ws HTTP/1.1
Host: host
Connection: Upgrade
Pragma: no-cache
Cache-Control: no-cache
User-Agent: 
Upgrade: websocket
Origin: http://ip:port
Sec-WebSocket-Version: 13
Accept-Encoding: gzip, deflate
Accept-Language: zh-CN,zh;q=0.9
Sec-WebSocket-Key: o88HTN24liAvfOOFHhVyKQ==
Sec-WebSocket-Extensions: permessage-deflate; client_max_window_bits

响应报文

HTTP/1.1 101 Switching Protocols
Server: nginx/1.20.2
Date: Mon, 01 Aug 2024 09:43:58 GMT
Connection: upgrade
upgrade: websocket
sec-websocket-accept: FMID9mDELWawU35fGSFJ/1H2790=

netty实现websocket

为什么选择netty实现websocket服务端?

  • 高性能和高并发:Netty是一个基于NIO的异步事件驱动的网络应用框架,它提供了高性能的网络通信能力。
  • 简化的编程模型:Netty提供了简化的编程模型,通过ChannelHandler和ChannelPipeline等组件,开发者可以方便地构建复杂的网络通信逻辑。这种模型不仅简化了代码,还提高了代码的可维护性和可读性。
  • 协议支持:Netty支持多种协议,包括HTTP、HTTPS、TCP、UDP等,这意味着可以在Netty中轻松实现WebSocket协议的支持。
自定义handler
@ChannelHandler.Sharable
public class WebSocketBusinessHandler extends SimpleChannelInboundHandler<TextWebSocketFrame> {@Overrideprotected void channelRead0(ChannelHandlerContext ctx, TextWebSocketFrame msg) throws Exception {// 处理客户端传输过来的消息String content = msg.text();log.debug("接收到消息:{}", content);Response resp = execute(content);ctx.writeAndFlush(new TextWebSocketFrame(JsonUtils.toJson(resp)));}@Overridepublic void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {log.error("处理命令错误", cause);Response response = Response.fail(((BusinessException) cause).getErrCode());ctx.channel().writeAndFlush(new TextWebSocketFrame(JsonUtils.toJson(response)));}
}
实现ChannelInitializer.initChannel方法
ChannelPipeline pipeline = ch.pipeline();// websocket 基于http协议,所以要有http编解码器
pipeline.addLast(new HttpServerCodec());
// 对写大数据流的支持
pipeline.addLast(new ChunkedWriteHandler());
// 对httpMessage进行聚合,聚合成FullHttpRequest或FullHttpResponse
pipeline.addLast(new HttpObjectAggregator(1024 * 64));
// ====================== 以上是用于支持http协议    ======================
//可以增加一些内部接口// ====================== 以下是支持httpWebsocket ======================
// WebSocket 数据压缩扩展
pipeline.addLast(new WebSocketServerCompressionHandler());
/*** websocket 服务器处理的协议,用于指定给客户端连接访问的路由 : /ws* 本handler会帮你处理一些繁重的复杂的事* 会帮你处理握手动作: handshaking(close, ping, pong) ping + pong = 心跳* 对于websocket来讲,都是以frames进行传输的,不同的数据类型对应的frames也不同*/
pipeline.addLast(new WebSocketServerProtocolHandler("/ws"));
// 自定义的handler
pipeline.addLast(new WebSocketBusinessHandler());
服务端启动
EventLoopGroup mainGroup = new NioEventLoopGroup();
EventLoopGroup subGroup = new NioEventLoopGroup();try {ServerBootstrap serverBootstrap = new ServerBootstrap().group(mainGroup, subGroup).channel(NioServerSocketChannel.class).childHandler(initializer);ChannelFuture channelFuture = serverBootstrap.bind(port).sync();channelFuture.channel().closeFuture().sync();} catch (InterruptedException e) {e.printStackTrace();
} finally {mainGroup.shutdownGracefully();subGroup.shutdownGracefully();
}

在nginx、ingress暴露出来

# ingress暴露ws
- backend:serviceName: websocketservicePort: websocketpath: /wspathType: ImplementationSpecific
  http{# websocket# 根据客户端请求中 $http_upgrade 的值,来构造改变 $connection_upgrade 的值#  $connection_upgrade 的值会一直是 upgrade。然后如果 $http_upgrade 为空字符串的话,那值会是 close。map $http_upgrade $connection_upgrade {default upgrade;'' close;}}location  ^~  /ws{proxy_pass http://test.com;proxy_redirect    off;proxy_set_header X-Real-IP $remote_addr;proxy_set_header Host $host;proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;proxy_http_version 1.1;proxy_set_header Upgrade $http_upgrade;proxy_set_header Connection $connection_upgrade;}

nginx -s reload

分布式解决方案

WebSocket是有状态协议的,客户端连接服务器时只和集群中一个节点连接,数据传输过程中也只与这一节点通信。因此,WebSocket集群需要解决会话共享的问题。如果只采用单节点部署,虽然可以避免这一问题,但无法水平扩展支撑更高负载,有单点的风险。
WebSocket是有状态的,无法像直接HTTP以集群方式实现负载均衡,长连接建立后即与服务端某个节点保持着会话,因此集群下想要得知会话属于哪个节点,有两种方案,一种是使用类似微服务的注册中心来维护全局的会话映射关系,一种是使用事件广播由各节点自行判断是否持有会话。

以方案二为例:

连接管理
/*** WebSocket的连接对象。*/
public class WebSocketConn {/*** 通道*/private Channel channel;/*** 应用标签 application*/private String ;/*** 人员id*/private Long personId;/*** 人员名称*/private String personName;public WebSocketConn(Channel channel) {this.channel = channel;}public void bind(String appId, Long personId, String personName) {this.app = appId;this.personId = personId;this.personName = personName;}/*** 查询符合的数据*/public boolean match(@NonNull String appId, @NonNull Long targetPersonId) {if(app.equals(appId) && personId.equals(targetPersonId)){return true;}return false;}public Channel getChannel() {return channel;}public String getApp() {return app;}
}public class WebSocketConnRegister {private Map<Channel, WebSocketConn> connMap = new ConcurrentHashMap<>();public  void register(WebSocketConn conn) {connMap.put(conn.getChannel(), conn);}public void unregister(Channel channel) {connMap.remove(channel);}public WebSocketConn getByChannel(Channel channel) {return connMap.get(channel);}public List<WebSocketConn> findByTarget(String application, Long personId) {return connMap.values().stream().filter(w -> w.match(application, personId)).collect(Collectors.toList());}}
连接接入或断开

WebSocketBusinessHandler 实现ChannelHandlerAdapter的handlerAdded、handlerRemoved方法

@Override
public void handlerAdded(ChannelHandlerContext ctx) throws Exception {log.debug("有连接进入,连接channel_id:{}", ctx.channel().id().asLongText());WebSocketConn conn = new WebSocketConn(ctx.channel());register.register(conn);
}@Override
public void handlerRemoved(ChannelHandlerContext ctx) throws Exception {log.debug("连接断开,连接channel:{}", ctx.channel().id().asLongText());register.unregister(ctx.channel());
}

WebSocketBusinessHandler.channelRead0 增加处理逻辑
1、根据收到消息来;
2、发送channel 不在此jvm实例中,则发送到消息队列中

WebSocketConn conn = register.getByChannel(ctx.channel());
if (conn != null) {Response<ClientCommandRet> resp = 处理业务逻辑;ctx.writeAndFlush(new TextWebSocketFrame(JsonUtils.toJson(resp)));return;
}else{// 不在当前实例,发送消息到RocketMQapplicationContext.getBean(RocketMQTemplate.class).convertAndSend("topicName",event);ctx.channel().close();
}

其它实例上监听topic后,同样查找目标发送channel,如果不存在,直接结束。

@Component
@RocketMQMessageListener(topic = "topicName",consumerGroup = "groupName",messageModel = MessageModel.BROADCASTING)
public static class WebsocketConsumerListener implements RocketMQListener<String> {public void onMessage(String event) {log.info("consume event: {}", event);try {List<WebSocketConn> webSocketConnList = register.findByTarget(event.getApplication(), event.getTargetPersonId());String msg = assembleMsg(event);webSocketConns.stream().parallel().forEach(conn -> sendMsg2Channel(msg, conn));} catch (Exception e) {e.printStackTrace();}}
}private void sendMsg2Channel(String msg, WebSocketConn conn) {try {conn.getChannel().writeAndFlush(new TextWebSocketFrame(msg));} catch (Exception e) {log.error("推送消息失败", e);try {conn.getChannel().close();} catch (Exception ex) {}}
}

注意:
1、rocketmq使用广播模式,因为每条消息每个节点都要消费一次。
2、kafka则要保证所有的partition都要订阅到。

List<PartitionInfo> partitionInfos = kafkaConsumer.partitionsFor("topicName");

参考

http://www.ruanyifeng.com/blog/2017/05/websocket.html
http://coolaf.com/tool/chattest 测试工具

这篇关于Netty实现WebSocket及分布式解决方案的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/1129528

相关文章

Python实现精准提取 PDF中的文本,表格与图片

《Python实现精准提取PDF中的文本,表格与图片》在实际的系统开发中,处理PDF文件不仅限于读取整页文本,还有提取文档中的表格数据,图片或特定区域的内容,下面我们来看看如何使用Python实... 目录安装 python 库提取 PDF 文本内容:获取整页文本与指定区域内容获取页面上的所有文本内容获取

基于Python实现一个Windows Tree命令工具

《基于Python实现一个WindowsTree命令工具》今天想要在Windows平台的CMD命令终端窗口中使用像Linux下的tree命令,打印一下目录结构层级树,然而还真有tree命令,但是发现... 目录引言实现代码使用说明可用选项示例用法功能特点添加到环境变量方法一:创建批处理文件并添加到PATH1

Java使用HttpClient实现图片下载与本地保存功能

《Java使用HttpClient实现图片下载与本地保存功能》在当今数字化时代,网络资源的获取与处理已成为软件开发中的常见需求,其中,图片作为网络上最常见的资源之一,其下载与保存功能在许多应用场景中都... 目录引言一、Apache HttpClient简介二、技术栈与环境准备三、实现图片下载与保存功能1.

canal实现mysql数据同步的详细过程

《canal实现mysql数据同步的详细过程》:本文主要介绍canal实现mysql数据同步的详细过程,本文通过实例图文相结合给大家介绍的非常详细,对大家的学习或工作具有一定的参考借鉴价值,需要的... 目录1、canal下载2、mysql同步用户创建和授权3、canal admin安装和启动4、canal

Nexus安装和启动的实现教程

《Nexus安装和启动的实现教程》:本文主要介绍Nexus安装和启动的实现教程,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录一、Nexus下载二、Nexus安装和启动三、关闭Nexus总结一、Nexus下载官方下载链接:DownloadWindows系统根

SpringBoot集成LiteFlow实现轻量级工作流引擎的详细过程

《SpringBoot集成LiteFlow实现轻量级工作流引擎的详细过程》LiteFlow是一款专注于逻辑驱动流程编排的轻量级框架,它以组件化方式快速构建和执行业务流程,有效解耦复杂业务逻辑,下面给大... 目录一、基础概念1.1 组件(Component)1.2 规则(Rule)1.3 上下文(Conte

MySQL 横向衍生表(Lateral Derived Tables)的实现

《MySQL横向衍生表(LateralDerivedTables)的实现》横向衍生表适用于在需要通过子查询获取中间结果集的场景,相对于普通衍生表,横向衍生表可以引用在其之前出现过的表名,本文就来... 目录一、横向衍生表用法示例1.1 用法示例1.2 使用建议前面我们介绍过mysql中的衍生表(From子句

Mybatis的分页实现方式

《Mybatis的分页实现方式》MyBatis的分页实现方式主要有以下几种,每种方式适用于不同的场景,且在性能、灵活性和代码侵入性上有所差异,对Mybatis的分页实现方式感兴趣的朋友一起看看吧... 目录​1. 原生 SQL 分页(物理分页)​​2. RowBounds 分页(逻辑分页)​​3. Page

MyBatis Plus 中 update_time 字段自动填充失效的原因分析及解决方案(最新整理)

《MyBatisPlus中update_time字段自动填充失效的原因分析及解决方案(最新整理)》在使用MyBatisPlus时,通常我们会在数据库表中设置create_time和update... 目录前言一、问题现象二、原因分析三、总结:常见原因与解决方法对照表四、推荐写法前言在使用 MyBATis

Python基于微信OCR引擎实现高效图片文字识别

《Python基于微信OCR引擎实现高效图片文字识别》这篇文章主要为大家详细介绍了一款基于微信OCR引擎的图片文字识别桌面应用开发全过程,可以实现从图片拖拽识别到文字提取,感兴趣的小伙伴可以跟随小编一... 目录一、项目概述1.1 开发背景1.2 技术选型1.3 核心优势二、功能详解2.1 核心功能模块2.