You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Spring Boot服务层获取用户信息并发送指定用户WebSocket消息

解决Spring Boot WebSocket中ChannelInterceptor获取的User传递到服务层的问题

下面提供几种实用方案,帮你把拦截器中获取的User对象传递到服务层,实现定向消息推送:

方案1:将User存入WebSocketSession属性中

WebSocketSession自带attributes存储能力,可在拦截器中把User存入session,后续在服务层或消息处理器中直接读取。

步骤1:在ChannelInterceptor中存储User

@Component
public class UserChannelInterceptor implements ChannelInterceptor {
    @Override
    public Message<?> preSend(Message<?> message, MessageChannel channel) {
        StompHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class);
        // 处理连接请求时提取并存储User
        if (StompCommand.CONNECT.equals(accessor.getCommand())) {
            // 复用你已有的逻辑获取username和userId
            String username = accessor.getFirstNativeHeader("username");
            String userId = accessor.getFirstNativeHeader("userId");
            User user = new User(username, userId);
            
            WebSocketSession session = accessor.getSession();
            if (session != null) {
                session.getAttributes().put("currentUser", user);
            }
        }
        return message;
    }
}

步骤2:在服务层/消息处理器中获取User

如果是处理客户端消息的@MessageMapping方法,可直接注入WebSocketSession取出User:

@MessageMapping("/client/message")
public void handleClientMessage(WebSocketSession session, String content) {
    User currentUser = (User) session.getAttributes().get("currentUser");
    // 基于currentUser做业务逻辑,比如权限校验、定向转发等
}

方案2:维护全局Session-User映射表

在拦截器中维护线程安全的映射关系,存储WebSocketSession与User的对应关系,服务层直接通过userId查找会话或推送消息。

步骤1:在拦截器中维护映射

@Component
public class UserChannelInterceptor implements ChannelInterceptor {
    // 线程安全映射:sessionId -> User
    public static final ConcurrentHashMap<String, User> SESSION_USER_MAP = new ConcurrentHashMap<>();
    // 线程安全映射:userId -> WebSocketSession
    public static final ConcurrentHashMap<String, WebSocketSession> USER_SESSION_MAP = new ConcurrentHashMap<>();

    @Override
    public Message<?> preSend(Message<?> message, MessageChannel channel) {
        StompHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class);
        
        if (StompCommand.CONNECT.equals(accessor.getCommand())) {
            String username = accessor.getFirstNativeHeader("username");
            String userId = accessor.getFirstNativeHeader("userId");
            User user = new User(username, userId);
            WebSocketSession session = accessor.getSession();
            
            if (session != null) {
                String sessionId = session.getId();
                SESSION_USER_MAP.put(sessionId, user);
                USER_SESSION_MAP.put(userId, session);
            }
        } else if (StompCommand.DISCONNECT.equals(accessor.getCommand())) {
            // 用户断开时清理映射,避免内存泄漏
            WebSocketSession session = accessor.getSession();
            if (session != null) {
                String sessionId = session.getId();
                User user = SESSION_USER_MAP.remove(sessionId);
                if (user != null) {
                    USER_SESSION_MAP.remove(user.getUserId());
                }
            }
        }
        return message;
    }
}

步骤2:服务层定向推送

@Service
public class WebSocketMessageService {
    @Autowired
    private SimpMessagingTemplate messagingTemplate;

    public void sendMessageToTargetUser(String targetUserId, Object message) {
        // 方式1:直接通过会话发送消息
        WebSocketSession session = UserChannelInterceptor.USER_SESSION_MAP.get(targetUserId);
        if (session != null && session.isOpen()) {
            try {
                session.sendMessage(new TextMessage(JSON.toJSONString(message)));
            } catch (IOException e) {
                // 处理发送异常,比如记录日志、清理无效会话
            }
        }

        // 方式2:使用SimpMessagingTemplate定向推送(推荐)
        // 客户端需订阅 /user/{targetUserId}/notify 这类目的地
        messagingTemplate.convertAndSendToUser(targetUserId, "/notify", message);
    }
}

方案3:用ThreadLocal传递User(适合REST接口触发的推送)

如果其他微服务通过REST接口调用WebSocket服务推送消息,可在接口层把User存入ThreadLocal,服务层直接获取,避免参数传递繁琐。

步骤1:定义ThreadLocal工具类

public class UserContextHolder {
    private static final ThreadLocal<User> USER_THREAD_LOCAL = new ThreadLocal<>();

    public static void setUser(User user) {
        USER_THREAD_LOCAL.set(user);
    }

    public static User getUser() {
        return USER_THREAD_LOCAL.get();
    }

    // 请求结束后必须清理,防止内存泄漏
    public static void clear() {
        USER_THREAD_LOCAL.remove();
    }
}

步骤2:在REST接口中设置User

@RestController
@RequestMapping("/ws")
public class WebSocketPushController {
    @Autowired
    private WebSocketMessageService messageService;

    @PostMapping("/push")
    public ResponseEntity<Void> pushToUser(@RequestParam String targetUserId, 
                                           @RequestHeader String username,
                                           @RequestHeader String userId,
                                           @RequestBody Object message) {
        User currentUser = new User(username, userId);
        UserContextHolder.setUser(currentUser);
        try {
            messageService.sendMessageToTargetUser(targetUserId, message);
            return ResponseEntity.ok().build();
        } finally {
            // 无论请求成败,都清理ThreadLocal
            UserContextHolder.clear();
        }
    }
}

步骤3:服务层获取User

@Service
public class WebSocketMessageService {
    @Autowired
    private SimpMessagingTemplate messagingTemplate;

    public void sendMessageToTargetUser(String targetUserId, Object message) {
        User operator = UserContextHolder.getUser();
        // 基于operator做权限校验、日志记录等
        messagingTemplate.convertAndSendToUser(targetUserId, "/notify", message);
    }
}

注意事项

  • 使用convertAndSendToUser时,客户端订阅的目的地需以/user/开头,比如/user/{userId}/notify。
  • 全局映射表必须处理断开连接的场景,避免积累无效对象导致内存泄漏。
  • ThreadLocal必须在请求生命周期结束后清理,异步调用场景更要注意,防止线程复用导致脏数据。

内容的提问来源于stack exchange,提问作者Devmina

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.24 03:52:06