Websocket+Disruptor并发框架实现协同编辑
Websocket+Disruptor处理协同编辑流程总结
▼text复制代码A.首先在Websocket配置类中声明websocket处理器,注册到某一个端点上,用于用户连接。 B.前端(用户)和该端点连接成功并发送消息之后,**会先到达websocket的拦截器**,拦截器对参数合法性、用户权限进行校验,校验成功将参数存储到attributes中(类似于向作用域对象中存储值) C.此时消息到达**handler处理器**,如果此时没有disruptor,消息在这里就会被处理掉,然后向用户广播信息了;但是有disruptor的话,就要封装成**事件**发送到**生产者** D.**disruptor生产者**将事件存储到**无锁环形缓冲区**中(取next()下标值,然后将该下标指向的event取到,然后就地赋值) E. **disruptor消费者**检查到队列中有事件,执行onEvent方法,switch一下消息类型,然后广播消息! F.**用户**在接收到消息之后,做相应处理。
Websocket:
初步了解
协议的种类:
a.全双工通信协议:允许数据同时在两个方向上传输,比如微信聊天,两个人同时可以向对方发送消息!(websocket就属于全双工通信协议)
b. 半双工通信,在同一时刻,只能由一方向另一方发送数据,比如客户端向服务器端发送请求,然后服务器端返回响应。
c.单工通信:数据只能单向传输,不能反向传输;
websocket连接成功之前会有一个基于http协议的握手请求,这个请求是websocket建立连接的前提,表明希望切换协议;
如果服务器是支持websocket协议的,那么就会返回一个101状态码,表示协议切换成功
websocket过滤器Interceptor
我们要在用户和websocket连接之前(也就是用户和websocket服务器建立连接之前)要判断用户的合法性,并且把用户的session添加到会话集合中。校验规则有:用户必须是同一个团队空间中的成员,并且有权限编辑图片,userId、pictureIdb必须是合法非空的;websocket拦截器要实现HandShakeIntegerceptor接口,并且重写两个方法:一个是握手前执行的,一个握手后执行的。
▼java复制代码@Override public boolean beforeHandshake(ServerHttpRequest serverHttpRequest, ServerHttpResponse response, WebSocketHandler wsHandler, Map<String, Object> attributes) throws Exception { //1.获取请求参数 并且判断 比如说pictureId\loginUser\该用户是否有该图片的操作权限 if(serverHttpRequest instanceof ServletServerHttpRequest){ //1.1 判断pictureId HttpServletRequest request = ((ServletServerHttpRequest) serverHttpRequest).getServletRequest(); String pictureId = request.getParameter("pictureId"); if(StrUtil.isBlank(pictureId)){ log.info("缺少参数,拒绝握手"); return false; } //1.2 判断登录用户 LoginUserVO loginUser = userService.getLoginUser(request); if(ObjectUtil.isNull(loginUser)){ log.info("用户未登录,拒绝握手"); } //1.3 判断当前图片是否存在 Picture picture = pictureService.getById(Long.valueOf(pictureId)); if(ObjectUtil.isNull(picture)){ log.info("图片不存在,拒绝握手"); return false; } //1.4 判断当前图片是公共图库/个人空间/团队空间的 Long spaceId = picture.getSpaceId(); Space space=null; if(spaceId!=null){ space=spaceService.getById(spaceId); if(space.getSpaceType().equals(SpaceTypeEnum.PRIVATE_SPACE.getValue())){ log.info("不是团队空间,拒绝握手"); return false; } } //1.5 如果是团队空间 判断当前该用户是否有edit的权限 List<String> permissionList = spaceUserManager.getPermissionList(space, loginUser); if(!permissionList.contains(SpaceUserPermissionConstant.PICTURE_EDIT)){ log.info("没有操作该空间的权限,拒绝握手"); return false; } //2.参数校验完毕 向attributes添加参数 attributes.put("pictureId",Long.valueOf(pictureId)); attributes.put("userId",loginUser.getId()); attributes.put("user",loginUser); } return true; }
websocket处理器handler
在这个处理器中,我们定义了一些方法来处理客户端发送的消息。然后使用广播 在编写handler方法之前,先声明两个变量
▼java复制代码//把它理解成正在编辑某个图片的人 private final Map<Long, Long> pictureEditingUsers = new ConcurrentHashMap<>(); //把它理解成某个图片的编辑空间(在该空间里面的人可以看到编辑人的操作) private final Map<Long, Set<WebSocketSession>> pictureSessions = new ConcurrentHashMap<>();
A:连接建立成功函数。此时是进入这个编辑的空间
▼text复制代码 在这个方法中,我们要实现的逻辑是:需要将发送该请求的用户的session添加到该pictureId对应的set<websocketSession>中,(把这个set集合理解成一个编辑的空间,里面存的是准备编辑该图片的用户)然后创建一条消息,交给广播去发送。
▼java复制代码public void afterConnectionEstablished(WebSocketSession session) throws Exception throws IOException{ Long pictureId=session.getAttributes().get("pictureId"); User user=(User)session.getAttributes().get("user"); pictureSessions.put(pictureId,ConcurrentHashMap.newKeySet()); pictureSessions.get(pictureId).add(session); //然后创建广播的message PictureEditResponseMessage pictureEditResponseMessage=new PictureEditResponseMessage(); pictureEditResponseMessage.setType(PictureEditMessageTypeEnum.INFO.getValue()); String message =String.format("%s加入了编辑空间",userVO.getUserName()); pictureEditResponseMessage.setMessage(message); pictureEditResponseMessage.setUserVO(userVO); broadcastToPicture(pictureId,pictureEditResponseMessage,session); }
B. 当某个用户要对图片进行编辑的时候(进入编辑状态)
▼text复制代码 在这个方法中,我们要实现的逻辑是:首先判断是否有人正在编辑,如果没人的话,就把当前userId放进去,然后还需要我们发送一条广播消息,告知用户有人正在编辑
▼java复制代码public void handleEnterEditMessage(PictureEditRequestMessage requestMessage,WebSocketSession session,User user,Long pictureId) throws IOException { //没有用户正在编辑该图片 才能进入编辑 if(!pictureEditingUsers.containsKey(pictureId)){ //设置当前用户为编辑用户 pictureEditingUsers.put(pictureId,user.getId()); //然后广播一条消息 PictureEditResponseMessage pictureEditResponseMessage=new PictureEditResponseMessage(); String message=String.format("%s正在开始编辑图片",user.getUserName()); pictureEditResponseMessage.setMessage(message); pictureEditResponseMessage.setType(PictureEditMessageTypeEnum.ENTER_EDIT.getValue()); pictureEditResponseMessage.setUserVO(userService.getUserVO(user)); broadcastToPicture(pictureId,pictureEditResponseMessage,session); } }
C.当某用户对图片进行编辑动作的时候
▼text复制代码在这个方法中,我们要实现的逻辑是:首先,我们从requestMessage中获取到该动作的类型,然后要判断该动作是否合法。如果合法的话,还要判断正在编辑该照片的用户和从request获取的登录用户是不是同一个人。最后创建一条信息,告诉其他用户,xxx正在执行xxx动作
▼java复制代码public void handleEditActionMessage(PictureEditResponseMessage requestMessage,WebSocketSession session,User user,Long pictureId) throws IOException { //1.首先要判断操作的类型是否合规 String editAction = requestMessage.getEditAction(); PictureEditMessageTypeEnum enumByValue = PictureEditMessageTypeEnum.getEnumByValue(editAction); if(ObjectUtil.isNull(enumByValue)){ Log.info("当前操作类型不合规"); return ; } //2.然后判断当前图片 Long editUserId = pictureEditingUsers.get(pictureId); if(editUserId!=null&&editUserId.equals(user.getId())){ //如果当前该照片的编辑者确定是editUser PictureEditResponseMessage responseMessage=new PictureEditResponseMessage(); String message=String.format("%s执行了%s操作",user.getUserName(),enumByValue.getValue()); responseMessage.setEditAction(editAction); responseMessage.setMessage(message); responseMessage.setUserVO(userService.getUserVO(user)); broadcastToPicture(pictureId,responseMessage,session); } }
D.某用户退出编辑状态
▼text复制代码 在这个方法中,我们要实现的逻辑是:我们需要判断:当前用户是不是编辑该图片的用户,如果是的话,那我们就要从pictureEditingUsers中将该用户删除。 然后创建一条信息,广播。这里是退出编辑状态,但是session还在sessions的集合中
▼java复制代码public void handleExitEditMessage(PictureEditRequestMessage requestMessage,WebSocketSession session,User user,Long pictureId) throws IOException { Long editUserId = pictureEditingUsers.get(pictureId); if(editUserId!=null&&editUserId.equals(user.getId())){ //如果当前编辑的用户是当前登录的用户 pictureEditingUsers.remove(editUserId); //构造退出编辑的消息 PictureEditResponseMessage pictureEditResponseMessage=new PictureEditResponseMessage(); String message=String.format("%s用户退出了编辑",user.getUserName()); pictureEditResponseMessage.setMessage(message); pictureEditResponseMessage.setType(PictureEditMessageTypeEnum.EXIT_EDIT.getValue()); pictureEditResponseMessage.setUserVO(userService.getUserVO(user)); broadcastToPicture(pictureId,pictureEditResponseMessage,session); } }
E.用户断开连接后的操作
▼text复制代码 类似于handleExitEditMessage
▼java复制代码@Override public void afterConnectionClosed(WebSocketSession session, CloseStatus status) throws Exception { Map<String, Object> attributes = session.getAttributes(); Long pictureId = (Long)attributes.get("pictureId"); User user=(User) attributes.get("user"); handleExitEditMessage(null,session,user,pictureId); Set<WebSocketSession> webSocketSessions = pictureSessions.get(pictureId); if(CollUtil.isNotEmpty(webSocketSessions)){ webSocketSessions.remove(session); //如果该图片对应的编辑空间里都没人了 那就直接remove掉吧 if(webSocketSessions.isEmpty()){ pictureSessions.remove(pictureId); } } //制造消息 PictureEditResponseMessage responseMessage=new PictureEditResponseMessage(); String message=String.format("%s用户离开编辑",user.getUserName()); responseMessage.setType(PictureEditMessageTypeEnum.EXIT_EDIT.getValue()); responseMessage.setMessage(message); responseMessage.setUserVO(userService.getUserVO(user)); broadcastToPicture(pictureId,responseMessage,session); }
F: 广播消息的动作
广播消息就是挨个向每一个接受者发送消息,比如说三个用户和服务器端建立了连接,那么就要依次向三个用户发送消息(发送什么消息呢?就是xx用户干了xx事,用户(前端)在知道了这些消息之后,就能实现一些逻辑)
为什么要引入Disruptor?
在当前这个场景中,我们为什么要引入Disruptor呢?每一个websocket连接对应一个独立的WebsocketSession,消息的处理在WebsocketSession的线程中处理。基于websocket服务器在处理一些和它连接用户发送的请求的时候,是根据接受的顺序同步处理的,并不是并发执行。
也就是说 接收消息和处理消息是同一个线程执行的 ,当处理消息任务是一个耗时的操作时,系统响应时间就会变长,请求源源不断到达服务器时,资源就会耗尽 。
因此我们引入一个队列,将需要执行的任务都放到这个队列中,交给专门的线程去处理这件事情。
Disruptor是什么?
Disruptor是一种高性能的并发框架,底层是一种无锁的环形队列数据结构组成。用于解决高吞吐量和低延迟场景中的并发问题。
a.环形缓冲区:这样可以避免扩容引起的性能开销,以及避免一系列不需要的垃圾回收
b.无锁设计:降低了线程切换的成本
c.序列号机制:生产者是存储数据的时候,会返回一个下一个数据下标的序列号。
事件在消息队列的存放下标位置就是序列号机制决定的,这样就可以实现并发处理不同序列号的事件
Disruptor实战
生产者、消费者、配置类、事件
A.配置类
在配置类中,我们new了一个disruptor对象,其中包括它的环形缓冲区的大小,并且给disruptor对象制定了消费者
▼java复制代码@Bean("pictureEditEventDisruptor") public Disruptor<PictureEditEvent> messageModelRingBuffer(){ //ringBuffer的大小 int bufferSize=1024*265; Disruptor<PictureEditEvent> disruptor= new Disruptor<>( PictureEditEvent::new, bufferSize, ThreadFactoryBuilder.create().setNamePrefix("pictureEditEventDisruptor").build() ); //设置消费者 disruptor.handleEventsWithWorkerPool(pictureEditEventWorkHandler); //开启start disruptor.start(); return disruptor; }
B.生产者
因为我们是一个多人协同编辑的场景,当用户向服务器(websocket)发送消息的时候,websocket就会调用我们的Disruptor,传递参数封装成一个事件,存储到队列中。
在我们这方法中,首先要获取到disruptor的一个无锁环形缓冲区!然后计算一下当前存储的下标位置(这个在每次存放完一个事件的时候 next就会指向下一个空闲区),然后将我们接收到的参数封装成一个事件对象就OK
▼java复制代码@Resource @Lazy Disruptor<PictureEditEvent> disruptor; public void publishPictureEditEvent(PictureEditRequestMessage requestMessage, WebSocketSession session, User user,Long pictureId){ RingBuffer<PictureEditEvent> ringBuffer = disruptor.getRingBuffer(); long next = ringBuffer.next(); //获取到一个该next下标存储的东西 PictureEditEvent pictureEditEvent = ringBuffer.get(next); //这里直接修改的是next下标下存储的那个事件 相当于是一个环形缓冲区 然后每一个里面都放了一个事件 但是还没有被赋值 pictureEditEvent.setPictureId(pictureId); pictureEditEvent.setUser(user); pictureEditEvent.setPictureEditRequestMessage(requestMessage); pictureEditEvent.setSession(session); //发布事件 ringBuffer.publish(next); }
C.消费者
在定义disruptor的时候,已经和消费者进行绑定了。因此在队列中有事件的时候,就会自动调用消费者的onEvent方法?
我们的消费者在接收到pictureEditEvent事件之后,就会对该事件进行处理。1.获取到它的requestMessage对象,然后获取它的参数 2.然后根据type类型计算出相应的枚举类 3.然后switch选择执行某一种方法
▼java复制代码@Resource private PictureEditHandler pictureEditHandler; @Resource private UserService userService; @Override public void onEvent(PictureEditEvent pictureEditEvent) throws Exception { PictureEditRequestMessage requestMessage = pictureEditEvent.getPictureEditRequestMessage(); //获取到消息类别 String type = requestMessage.getType(); User user = pictureEditEvent.getUser(); WebSocketSession session = pictureEditEvent.getSession(); Long pictureId = pictureEditEvent.getPictureId(); PictureEditMessageTypeEnum messageTypeEnum = PictureEditMessageTypeEnum.valueOf(type); switch(messageTypeEnum){ case ENTER_EDIT: pictureEditHandler.handleEnterEditMessage(requestMessage,session,user,pictureId); break; case EDIT_ACTION: pictureEditHandler.handleEditActionMessage(requestMessage,session,user,pictureId); break; case EXIT_EDIT: pictureEditHandler.handleExitEditMessage(requestMessage,session,user,pictureId); break; default: PictureEditResponseMessage responseMessage=new PictureEditResponseMessage(); responseMessage.setMessage("消息类型错误"); responseMessage.setUserVO(userService.getUserVO(user)); responseMessage.setType(PictureEditMessageTypeEnum.ERROR.getValue()); session.sendMessage(new TextMessage(JSONUtil.toJsonStr(responseMessage))); break; } }
