智能图库模块-协同编辑
终于来到了最后一块功能的开发学习,对于这种多人实时协作的业务,需要考虑到怎么防止用户之间的协作冲突、协作的交互以及如何提高协作的效率等等,而本项目只是对于团队空间的图片有协同编辑的需求,没有其他太多复杂的东西,因此采用Websocket来实现协同编辑。
Websocket是另一种通信协议,与Http不同的是,它可以持续保持连接,且是全双工通信。对于图片编辑,我们规定有进入编辑状态、开始编辑状态和退出编辑状态,执行的动作与之前开发图片编辑的功能一样是改变大小方向和裁剪。由于websocket不能像http那样从request中获取到用户信息和图片信息等,需要通过拦截器来在即将建立websocket连接的时候给会话指定信息:
▼java复制代码@Component @Slf4j public class WsHandshakeInterceptor implements HandshakeInterceptor { @Resource private UserService userService; @Resource private PictureService pictureService; @Resource private SpaceService spaceService; @Resource private SpaceUserAuthManager spaceUserAuthManager; @Override public boolean beforeHandshake(@NotNull ServerHttpRequest request, @NotNull ServerHttpResponse response, @NotNull WebSocketHandler wsHandler, @NotNull Map<String, Object> attributes) { if (request instanceof ServletServerHttpRequest) { HttpServletRequest servletRequest = ((ServletServerHttpRequest) request).getServletRequest(); // 获取请求参数 String pictureId = servletRequest.getParameter("pictureId"); if (StrUtil.isBlank(pictureId)) { log.error("缺少图片参数,拒绝握手"); return false; } User loginUser = userService.getLoginUser(servletRequest); if (ObjUtil.isEmpty(loginUser)) { log.error("用户未登录,拒绝握手"); return false; } // 校验用户是否有该图片的权限 Picture picture = pictureService.getById(pictureId); if (picture == null) { log.error("图片不存在,拒绝握手"); return false; } Long spaceId = picture.getSpaceId(); Space space = null; if (spaceId != null) { space = spaceService.getById(spaceId); if (space == null) { log.error("空间不存在,拒绝握手"); return false; } if (space.getSpaceType() != SpaceTypeEnum.TEAM.getValue()) { log.info("不是团队空间,拒绝握手"); return false; } } List<String> permissionList = spaceUserAuthManager.getPermissionList(space, loginUser); if (!permissionList.contains(SpaceUserPermissionConstant.PICTURE_EDIT)) { log.error("没有图片编辑权限,拒绝握手"); return false; } // 设置 attributes attributes.put("user", loginUser); attributes.put("userId", loginUser.getId()); attributes.put("pictureId", Long.valueOf(pictureId)); // 记得转换为 Long 类型 } return true; } @Override public void afterHandshake(@NotNull ServerHttpRequest request, @NotNull ServerHttpResponse response, @NotNull WebSocketHandler wsHandler, Exception exception) { } }
编写完拦截器后,需要在配置中引入该拦截器,因此编写websocket配置类WebSocketConfig:
▼java复制代码@Configuration @EnableWebSocket public class WebSocketConfig implements WebSocketConfigurer { @Resource private PictureEditHandler pictureEditHandler; @Resource private WsHandshakeInterceptor wsHandshakeInterceptor; @Override public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) { registry.addHandler(pictureEditHandler, "/ws/picture/edit") .addInterceptors(wsHandshakeInterceptor) .setAllowedOrigins("*"); } }
除了拦截器,也需要websocket处理器来在接收到消息时进行相应的处理,定义一个PictureEditHandler处理器,它继承TextWebSocketHandler以便可以通过字符串形式接收和发送消息,我们主要实现其中的afterConnectionEstablished、handleTextMessage和afterConnectionClosed方法,它们分别是连接建立后的操作,处理连接中的消息和关闭连接后的处理,为了避免冲突,定义两个集合,一个集合用于表示当前图片谁在编辑,一个用于表示这个图片进入编辑状态的有哪些会话:
▼java复制代码// 每张图片的编辑状态,key: pictureId, value: 当前正在编辑的用户 ID private final Map<Long, Long> pictureEditingUsers = new ConcurrentHashMap<>(); // 保存所有连接的会话,key: pictureId, value: 用户会话集合 private final Map<Long, Set<WebSocketSession>> pictureSessions = new ConcurrentHashMap<>();
由于可能会有多个websocket客户端建立连接,为了在修改map数据时保证线程安全,使用ConcurrentHashMap来保证。
在处理消息时,需要通知会话中的其他用户,即广播:
▼java复制代码public void broadcastToPicture(Long pictureId, PictureEditResponseMessage pictureEditResponseMessage,WebSocketSession excludeSession) throws IOException { Set<WebSocketSession> sessionSet = pictureSessions.get(pictureId); if (CollUtil.isNotEmpty(sessionSet)){ //创建ObjectMapper ObjectMapper objectMapper = new ObjectMapper(); //配置序列化 将Long类型转为String 解决精度丢失 SimpleModule module = new SimpleModule(); module.addSerializer(Long.class, ToStringSerializer.instance); module.addSerializer(Long.TYPE, ToStringSerializer.instance);//支持long类型 objectMapper.registerModule(module); //序列化为json字符串 String message = objectMapper.writeValueAsString(pictureEditResponseMessage); TextMessage textMessage = new TextMessage(message); for (WebSocketSession session : sessionSet) { //排除掉的session不发送 if (excludeSession != null && session.equals(excludeSession)){ continue; } if (session.isOpen()){ session.sendMessage(textMessage); } } } }
该方法的excludeSession参数是排除的某个session,在用户自己编辑图片的时候他并不用被通知他自己正在干什么,当想全部人都通知的时候,这个参数为null即可。在建立好连接后会执行afterConnectionEstablished方法,在该方法中保存相应的会话并进行广播通知:
▼java复制代码public void afterConnectionEstablished(WebSocketSession session) throws Exception { super.afterConnectionEstablished(session); //保存会话到集合中 User user = (User) session.getAttributes().get("user"); Long pictureId = (Long) session.getAttributes().get("pictureId"); pictureSessions.putIfAbsent(pictureId, ConcurrentHashMap.newKeySet()); pictureSessions.get(pictureId).add(session); //构造响应,发送加入编辑的消息通知 PictureEditResponseMessage pictureEditResponseMessage = new PictureEditResponseMessage(); pictureEditResponseMessage.setType(PictureEditMessageTypeEnum.INFO.getValue()); String message = String.format("用户 %s 加入编辑", user.getUserName()); pictureEditResponseMessage.setMessage(message); pictureEditResponseMessage.setUser(userService.getUserVo(user)); //广播给该图片的所有用户 broadcastToPicture(pictureId,pictureEditResponseMessage); }
考虑到websocket通常是一个长连接,是需要占用资源的,且对于单个客户端而言,它连续发送多个消息的时候服务器是按照接收顺序来进行处理的,并且接收信息和处理消息是同一个线程,当处理耗时变长且客户端多起来的时候就会导致一个客户端在执行的时候其他客户端需要等待很长一段时间,甚至耗尽服务器资源,对于这种多个线程等待的情况,就可以考虑使用异步的方式来处理,即我们把原本由一个人负责接收信息和处理消息变成两个人分别负责接收和处理,异步的同时也需要引入一个队列来存放任务以保证任务能够按顺序执行,因此这里引入了Disruptor无锁队列来实现。
Disruptor相较于传统的BlockingQueue,采用的是环形队列作为缓冲区,这样在规定完缓冲区大小后不会再进行扩容,且disruptor会在生产者赶上消费者且消费者未处理完数据时进行等待保证不被覆盖,此外对于并发控制它采用的是乐观锁CAS和内存屏障,降低了线程切换的开销,还利用了CPU的缓存局部性,提高了访问速度。
Disruptor是一个封装了RingBuffer、Producer和Consumer的类,有着生产者和消费者的概念,按照消息队列来看,我们需要定义事件以及事件的生产者和消费者,对于我们的需求来说,这个事件的内容要包含当前的websocket session、该session现在是哪个用户,哪个图片以及它的消息是什么,因此定义事件如下:
▼java复制代码public class PictureEditEvent { /** * 消息 */ private PictureEditRequestMessage pictureEditRequestMessage; /** * 当前用户的 session */ private WebSocketSession session; /** * 当前用户 */ private User user; /** * 图片 id */ private Long pictureId; }
接下来定义消费者PictureEditEventWorkHandler,消费者的职责是根据消息类型来分发到对应的处理器:
▼java复制代码public class PictureEditEventWorkHandler implements WorkHandler<PictureEditEvent> { @Resource private PictureEditHandler pictureEditHandler; @Resource private UserService userService; @Override public void onEvent(PictureEditEvent pictureEditEvent) throws Exception { PictureEditRequestMessage pictureEditRequestMessage = pictureEditEvent.getPictureEditRequestMessage(); WebSocketSession session = pictureEditEvent.getSession(); User user = pictureEditEvent.getUser(); Long pictureId = pictureEditEvent.getPictureId(); String type = pictureEditRequestMessage.getType(); PictureEditMessageTypeEnum pictureEditMessageTypeEnum = PictureEditMessageTypeEnum.getEnumByValue(type); //根据消息类型处理消息 switch (pictureEditMessageTypeEnum){ case ENTER_EDIT: pictureEditHandler.handleEnterEditMessage(pictureEditRequestMessage,session,user,pictureId); break; case EDIT_ACTION: pictureEditHandler.handleEditActionMessage(pictureEditRequestMessage,session,user,pictureId); break; case EXIT_EDIT: pictureEditHandler.handleExitEditMessage(pictureEditRequestMessage,session,user,pictureId); break; default: //其他消息类型 返回错误提示 PictureEditResponseMessage pictureEditResponseMessage = new PictureEditResponseMessage(); pictureEditResponseMessage.setType(PictureEditMessageTypeEnum.ERROR.getValue()); pictureEditResponseMessage.setMessage("消息类型错误"); pictureEditResponseMessage.setUser(userService.getUserVo(user)); session.sendMessage(new TextMessage(JSONUtil.toJsonStr(pictureEditResponseMessage))); break; } } }
其中处理消息的方法写在了之前的PictureEditHandler中,它们的逻辑都很简单,对于进入编辑状态,在没有用户正在编辑图片时才将用户加入会话并发送通知,而编辑动作和退出编辑还需额外验证当前用户是否为当前编辑者。
编写Disruptor配置类来将事件和消费者关联到disruptor实例中:
▼java复制代码@Configuration public class PictureEditEventDisruptorConfig { @Resource private PictureEditEventWorkHandler pictureEditEventWorkHandler; @Bean("pictureEditEventDisruptor") public Disruptor<PictureEditEvent> messageModelRingBuffer() { // ringBuffer 的大小 int bufferSize = 1024 * 256; Disruptor<PictureEditEvent> disruptor = new Disruptor<>( PictureEditEvent::new, bufferSize, ThreadFactoryBuilder.create().setNamePrefix("pictureEditEventDisruptor").build() ); // 设置消费者 disruptor.handleEventsWithWorkerPool(pictureEditEventWorkHandler); // 开启 disruptor disruptor.start(); return disruptor; } }
接着编写生产者,生产者要将事件放入到disruptor的环形缓冲区中,并且通过shutdown完成优雅关机保证服务停止时所有事件都能够被处理:
▼java复制代码public class PictureEditEventProducer { @Resource private Disruptor<PictureEditEvent> pictureEditEventDisruptor; /** * 发布事件 * @param pictureEditRequestMessage * @param session * @param user * @param pictureId */ public void publishEvent(PictureEditRequestMessage pictureEditRequestMessage, WebSocketSession session, User user, Long pictureId){ RingBuffer<PictureEditEvent> ringBuffer = pictureEditEventDisruptor.getRingBuffer(); //获取到可以放置事件的位置 long next = ringBuffer.next(); PictureEditEvent pictureEditEvent = ringBuffer.get(next); pictureEditEvent.setPictureEditRequestMessage(pictureEditRequestMessage); pictureEditEvent.setSession(session); pictureEditEvent.setUser(user); pictureEditEvent.setPictureId(pictureId); //发布事件 ringBuffer.publish(next); } /** * 优雅停机 */ @PreDestroy public void destroy(){ pictureEditEventDisruptor.shutdown(); } }
定义完事件、生产者和消费者后,回到前面的websocket处理器PictureEditHandler,实现handleTextMessage方法:
▼java复制代码@Override protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception { super.handleTextMessage(session, message); //获取消息内容,将消息字符串转为为PictureEditMessage PictureEditRequestMessage pictureEditRequestMessage = JSONUtil.toBean(message.getPayload(), PictureEditRequestMessage.class); //从Session属性获取到公共参数 User user = (User) session.getAttributes().get("user"); Long pictureId = (Long) session.getAttributes().get("pictureId"); //根据消息类型处理消息 生成消息到disruptor队列中 pictureEditEventProducer.publishEvent(pictureEditRequestMessage,session,user,pictureId); }
最后还要记得在websocket连接关闭后释放掉资源,实现afterConnectionClosed方法:
▼java复制代码@Override public void afterConnectionClosed(WebSocketSession session, CloseStatus status) throws Exception { super.afterConnectionClosed(session, status); User user = (User) session.getAttributes().get("user"); Long pictureId = (Long) session.getAttributes().get("pictureId"); //移除当前用户的编辑状态 handleExitEditMessage(null,session,user,pictureId); //删除会话 Set<WebSocketSession> sessionSet = pictureSessions.get(pictureId); if (sessionSet!= null){ sessionSet.remove(session); if (sessionSet.isEmpty()){ pictureSessions.remove(pictureId); } } //通知其他用户当前用户离开了编辑 PictureEditResponseMessage pictureEditResponseMessage = new PictureEditResponseMessage(); pictureEditResponseMessage.setType(PictureEditMessageTypeEnum.INFO.getValue()); String message = String.format("用户 %s 离开了编辑", user.getUserName()); pictureEditResponseMessage.setMessage(message); pictureEditResponseMessage.setUser(userService.getUserVo(user)); broadcastToPicture(pictureId,pictureEditResponseMessage); }
这样就实现了基于disruptor的异步消息处理。前端部分这里就不作记录了,至此,本项目只差部署上线了,明天将项目部署后就要开始针对着写简历和准备八股算法了。
