核心业务处理包含会话防并发创建、消息事务入库、WebSocket推送、离线消息补发、批量已读标记。package com.para.system.chat.service.impl;import com.para.common.utils.SecurityUtils;import com.para.system.chat.domain.ChatMessage;import com.para.system.chat.domain.ChatSession;import com.para.system.chat.domain.ChatMsgDTO;import com.para.system.chat.mapper.ChatMessageMapper;import com.para.system.chat.mapper.ChatSessionMapper;import com.para.system.chat.service.ChatService;import com.para.system.chat.websocket.ChatWebSocketServer;import lombok.extern.slf4j.Slf4j;import org.springframework.stereotype.Service;import org.springframework.transaction.annotation.Transactional;import javax.annotation.Resource;import java.util.Date;import java.util.List;import java.util.Set;import java.util.stream.Collectors;/**聊天业务实现类功能会话管理、消息存储、未读计数、消息推送、离线消息补发*/Slf4jServicepublic class ChatServiceImpl implements ChatService {Resourceprivate ChatSessionMapper sessionMapper;Resourceprivate ChatMessageMapper msgMapper;/**获取或创建两个用户之间的聊天会话防并发重复创建param fromUid 当前登录用户IDparam toUid 对方用户IDreturn 会话对象*/Overridepublic ChatSession getOrCreateSession(Long fromUid, Long toUid) {// 禁止自己和自己聊天if (fromUid.equals(toUid)) {throw new RuntimeException(“不允许与自己创建聊天会话”);}// 统一大小ID存储保证双人会话唯一Long minId Math.min(fromUid, toUid);Long maxId Math.max(fromUid, toUid);Long loginUid SecurityUtils.getLoginUser().getUserId();// 查询已有会话ChatSession session sessionMapper.selectByUserPair(minId, maxId, loginUid);if (session null) {try {session new ChatSession();session.setFromUserId(minId);session.setToUserId(maxId);sessionMapper.insertSession(session);} catch (Exception e) {// 唯一索引冲突其他线程已创建日志告警不抛出异常log.warn(“会话创建冲突重新查询会话记录”);}// 冲突后二次查询返回最新会话return sessionMapper.selectByUserPair(minId, maxId, loginUid);}return session;}/**保存消息并异步推送给接收方数据库操作在事务内WebSocket推送在事务外避免长连接阻塞事务*/Overridepublic void saveAndPushMessage(Long senderId, ChatMsgDTO dto) {Long receiverId dto.getReceiverId();// 事务入库、更新会话、未读数1saveMessageInTransaction(senderId, dto);// 事务外推送消息推送失败不影响消息存储try {ChatWebSocketServer.sendToUser(receiverId, dto);} catch (Exception e) {log.error(“WebSocket推送消息失败发送人:{},接收人:{}”, senderId, receiverId, e);}}/**事务方法保存消息、更新会话最后消息、未读数自增*/Transactional(rollbackFor Exception.class)public void saveMessageInTransaction(Long senderId, ChatMsgDTO dto) {Long sessionId dto.getSessionId();Long receiverId dto.getReceiverId();Date now new Date();// 1. 插入聊天消息默认未读0ChatMessage msg new ChatMessage();msg.setSessionId(sessionId);msg.setSenderId(senderId);msg.setReceiverId(receiverId);msg.setContent(dto.getContent());msg.setMsgType(dto.getMsgType());msg.setReadStatus(“0”);msg.setCreateTime(now);msgMapper.insert(msg);// 2. 更新会话最后一条消息内容和时间sessionMapper.updateLastMsg(sessionId, dto.getContent(), now);// 3. 当前会话未读计数1sessionMapper.incrUnread(sessionId);}/**查询会话分页历史消息*/Overridepublic List history(Long sessionId) {return msgMapper.selectBySession(sessionId);}/**查询当前用户会话列表支持对方昵称模糊检索*/Overridepublic List sessionList(Long uid, String nickName) {return sessionMapper.selectUserSessions(uid, nickName);}/**批量标记会话全部消息已读清空会话未读计数*/OverrideTransactionalpublic void markAllRead(Long sessionId, Long uid) {msgMapper.batchRead(sessionId, uid);sessionMapper.clearUnread(sessionId);}/**用户WebSocket上线补发所有离线未读消息*/Overridepublic void pushOfflineMessage(Long uid) {List unreadList msgMapper.selectUnread(uid);if (unreadList null || unreadList.isEmpty()) {return;}// 逐条推送离线消息for (ChatMessage msg : unreadList) {ChatMsgDTO dto new ChatMsgDTO();dto.setSessionId(msg.getSessionId());dto.setReceiverId(msg.getReceiverId());dto.setContent(msg.getContent());dto.setMsgType(msg.getMsgType());ChatWebSocketServer.sendToUser(uid, dto);}// 推送完成批量标记已读避免重复推送try {Set sessionIds unreadList.stream().map(ChatMessage::getSessionId).collect(Collectors.toSet());for (Long sessionId : sessionIds) {msgMapper.batchRead(sessionId, uid);sessionMapper.clearUnread(sessionId);}} catch (Exception e) {log.error(“离线消息标记已读失败”, e);}}}5.3 Mapper 接口补充ChatSessionMapper.javapackage com.para.system.chat.mapper;import com.para.system.chat.domain.ChatSession;import org.apache.ibatis.annotations.Param;import java.util.List;public interface ChatSessionMapper {ChatSession selectByUserPair(Param(from) Long minId, Param(to) Long maxId, Param(uid) Long loginUid); void insertSession(ChatSession session); void updateLastMsg(Param(sid) Long sessionId, Param(content) String content, Param(time) java.util.Date time); void incrUnread(Param(sid) Long sessionId); void clearUnread(Param(sid) Long sessionId); ListChatSession selectUserSessions(Param(uid) Long uid, Param(nickName) String nickName);}ChatMessageMapper.javapackage com.para.system.chat.mapper;import com.para.system.chat.domain.ChatMessage;import org.apache.ibatis.annotations.Param;import java.util.List;public interface ChatMessageMapper {void insert(ChatMessage message); ListChatMessage selectBySession(Param(sessionId) Long sessionId); void batchRead(Param(sessionId) Long sessionId, Param(uid) Long userId); ListChatMessage selectUnread(Param(uid) Long userId);}5.4 MyBatis Mapper XML ChatSessionMapper.xml核心复杂SQL联表用户获取昵称头像、子查询实时统计未读、昵称模糊搜索、双人会话唯一查询。?xml version1.0 encodingUTF-8?!-- 会话结果映射关联用户昵称、头像、实时未读数量 -- resultMap typecom.para.system.chat.domain.ChatSession idChatSessionResult result propertyid columnid / result propertyfromUserId columnfrom_user_id / result propertytoUserId columnto_user_id / result propertylastContent columnlast_content / result propertylastMsgTime columnlast_msg_time / result propertyunreadCount columnunread_count / result propertyrealSelfUnread columnreal_self_unread/ result propertycreateTime columncreate_time / result propertyupdateTime columnupdate_time / result columnnick_name propertytargetNickName/ result columnavatar propertytargetAvatar/ /resultMap !-- 根据双方用户ID查询唯一会话关联用户信息实时未读统计 -- select idselectByUserPair resultMapChatSessionResult SELECT cs.*, su.nick_name, su.avatar, COALESCE(cm.unread_total,0) AS real_self_unread FROM chat_session cs LEFT JOIN sys_user su ON su.user_id ( CASE WHEN cs.from_user_id #{uid} THEN cs.to_user_id ELSE cs.from_user_id END ) LEFT JOIN ( SELECT session_id, COUNT(1) unread_total FROM chat_message WHERE receiver_id #{uid} AND read_status 0 GROUP BY session_id ) cm ON cs.id cm.session_id WHERE (from_user_id #{from} AND to_user_id #{to}) OR (from_user_id #{to} AND to_user_id #{from}) LIMIT 1 /select !-- 新增双人会话 -- insert idinsertSession parameterTypeChatSession INSERT INTO chat_session ( from_user_id, to_user_id, last_content, last_msg_time, unread_count, create_time, update_time ) VALUES ( #{fromUserId}, #{toUserId}, #{lastContent}, #{lastMsgTime}, #{unreadCount}, sysdate(), sysdate() ) /insert !-- 更新会话最后一条消息内容和时间 -- update idupdateLastMsg UPDATE chat_session SET last_content #{content}, last_msg_time #{time}, update_time sysdate() WHERE id #{sid} /update !-- 会话未读计数 1 -- update idincrUnread UPDATE chat_session SET unread_count unread_count 1, update_time sysdate() WHERE id #{sid} /update !-- 清空会话未读计数 -- update idclearUnread