智能图库模块-协同编辑

终于来到了最后一块功能的开发学习,对于这种多人实时协作的业务,需要考虑到怎么防止用户之间的协作冲突、协作的交互以及如何提高协作的效率等等,而本项目只是对于团队空间的图片有协同编辑的需求,没有其他太多复杂的东西,因此采用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以便可以通过字符串形式接收和发送消息,我们主要实现其中的afterConnectionEstablishedhandleTextMessageafterConnectionClosed方法,它们分别是连接建立后的操作,处理连接中的消息和关闭连接后的处理,为了避免冲突,定义两个集合,一个集合用于表示当前图片谁在编辑,一个用于表示这个图片进入编辑状态的有哪些会话:

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是一个封装了RingBufferProducerConsumer的类,有着生产者和消费者的概念,按照消息队列来看,我们需要定义事件以及事件的生产者和消费者,对于我们的需求来说,这个事件的内容要包含当前的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的异步消息处理。前端部分这里就不作记录了,至此,本项目只差部署上线了,明天将项目部署后就要开始针对着写简历和准备八股算法了。

0个评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
shark
作者分享
做了一个月的小程序,终于准备到即将交付的时候了,这段时间里可真是被微信小程序的各种限制折磨的不要不要的,突出一个为难开发者,有几次改代码改到了两点多才改完!小程序的后端技术栈和鱼皮的智能Bi项目采用了差不多的技术栈,通过注解+aop来对接口进行令牌桶限流,用rabbitmq来缓解调用Ai服务的压力,异步生成报告,对于订单采用了自定义线程池结合async+scheduled注解实现定时任务的异步执行,每分钟遍历删除库中的订单,还有些功能等真正交付后再做下总结。唉,不知道这个项目再加上智能图库项目能不能让我这种0相关工作经验的双非往届生有面试机会和找到工作呢,不想再次尝试付出与回报不成正比的滋味了😭😭
8
今天收到了甲方的第一笔付款💰,没想到毕业一年半居然能用自己专业技能接个活赚了点外快,还是有点小高兴的。 这两天的学习状态不太行,我好像总会在经历了持续一段时间的学习后会有两三天学不进的状态🤧,不知道大伙会不会有这种状态,虽然我很想克服这种状态,但心里就好像有东西堵在那影响着学不进去(不知道是不是焦虑),希望明天能恢复状态吧。 转眼到了3月份,接的活的时间也剩下半个月了,下半部分就要开发调用Ai接口的代码和微信支付(感觉好像挺麻烦?)的代码了,同时也需要开始准备看面试题了,想问下各位是怎么使用面试鸭来高效刷题的呢,这两天(可能跟我的状态有关)看面试题总是有种似记住但又没全记住的感觉,上次背八股还是22年秋招的时候,对于一些知识可能还有记得一些,但由于我是0开发经验社招,不知道重点要看哪些好,希望各位大佬能给小弟一点建议🙏
4
这几天没有打卡,不是摆烂了,而是前几天有个朋友推了个活给我问我说做不做,是做一个微信小程序,并且需要调用Ai的接口,他因为工作太忙没有空接这个活,而我本来正在学习Bi项目,打算学完之后写简历刷面试题投简历,所以一开始还在犹豫,觉得自己23届毕业,虽说大学自学了几年Java,但是0相关工作经验,有点没底气,不过后来我拉到了另一个同学来帮忙,那个同学在大厂实习过,有过企业开发的经验,他也乐意加入一起做项目,于是最后便接下了这个活(主要是有钱可以拿)。由于自己是前端苦手,所以也是一边问Ai一边改页面的布局,改到今天觉得差不多了,甲方也觉得ok,后面就准备和同学一起写后端的接口了,在这几天也对着项目需求跟同学讨论了许多功能和结构上的设计,也算是一种积累经验的过程了,希望这一个月的工期内能把项目做好上线,最好是能写到简历上吧,孩子真的很想要一份工作!!这一个月干的活希望不要是白忙活,不然我真的要道心破碎了😭😭😭
7
智能Bi项目-限流优化
7
智能BI项目-调用硅基流动的DeepSeek
19
下载 APP