如何在ChannelInterceptor的preSend()中通过DB校验拒绝请求并解决阻塞问题
问题描述
我在WebSocket配置类中定义了如下代码:
@Configuration @EnableWebSocketMessageBroker class WebsocketConfig(<...>)
同时定义了一个拦截器Bean:
@Bean(name = ["csrfChannelInterceptor"]) fun csrfChannelInterceptor(): ChannelInterceptor { return object : ChannelInterceptor { override fun preSend(message: Message<*>, channel: MessageChannel): Message<*>? { val accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor::class.java) ?: return null when (accessor.command) { StompCommand.SUBSCRIBE -> { val chatId = extractChatIdFromDestination(accessor.destination) // 若用户不是参与者则阻止其订阅频道 if (!chatsService.existsByChatIdAndParticipantId(chatId, getCurrentUserId())) { throw AccessDeniedException("No permissions found") } } StompCommand.SEND -> { // 与SUBSCRIBE相同的校验逻辑 } else -> { // 无操作 } } return super.preSend(message, channel) } } }
当用户尝试订阅或发送消息到指定频道时,会触发如下错误:
Failed to send message to MessageChannel in session 912<...>8:Failed to send message to ExecutorSubscribableChannel[clientInboundChannel]
这个问题仅在preSend()中调用阻塞式数据库方法时出现,我想确认:
- 这种通过自定义拦截逻辑拒绝SUBSCRIBE/SEND请求的方式是否正确?
- 该如何解决当前的报错问题?
解决方案
一、校验方式的合理性
用ChannelInterceptor拦截preSend方法做权限校验的思路是正确的,这也是Spring官方推荐的WebSocket消息前置校验方案。但你当前的实现有两个关键问题:
- 直接抛出
AccessDeniedException会导致消息通道内部处理异常,无法给客户端返回可识别的错误响应 - 在
preSend中调用阻塞式数据库操作,会占满消息通道的同步线程池,引发通道发送失败的问题
二、具体修复步骤
- 替换阻塞式数据库操作,改用异步调用
Spring WebSocket的clientInboundChannel默认使用同步线程池,阻塞IO操作会快速耗尽线程资源,导致后续消息无法处理。你需要把数据库查询改成异步模式:
// 先将chatsService的方法改造为异步返回CompletableFuture fun existsByChatIdAndParticipantIdAsync(chatId: String, userId: String): CompletableFuture<Boolean> // 拦截器中改用异步处理逻辑 override fun preSend(message: Message<*>, channel: MessageChannel): Message<*>? { val accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor::class.java) ?: return null when (accessor.command) { StompCommand.SUBSCRIBE, StompCommand.SEND -> { val chatId = extractChatIdFromDestination(accessor.destination) val userId = getCurrentUserId() // 异步执行权限校验 chatsService.existsByChatIdAndParticipantIdAsync(chatId, userId) .whenComplete { hasPermission, throwable -> if (!hasPermission || throwable != null) { // 校验不通过时,主动给客户端发送ERROR帧 val errorAccessor = StompHeaderAccessor.create(StompCommand.ERROR) errorAccessor.message = "No permissions found" errorAccessor.sessionId = accessor.sessionId val errorMessage = MessageBuilder.createMessage(ByteArray(0), errorAccessor.messageHeaders) channel.send(errorMessage) // 标记原消息为错误状态,阻止其继续流转 accessor.setLeaveMutable(true) accessor.setError(true) } } } else -> {} } return if (accessor.isError) null else super.preSend(message, channel) }
正确返回错误响应,避免直接抛异常
直接抛出异常会导致通道内部报错,无法正确通知客户端权限校验失败。正确的做法是构建StompCommand.ERROR消息主动发送给客户端,同时标记原消息为错误状态,阻止其继续在通道中流转。调整线程池配置(可选)
如果因业务限制必须使用阻塞式操作,可以通过调整clientInboundChannel的线程池大小缓解问题:
override fun configureClientInboundChannel(registration: ChannelRegistration) { registration.taskExecutor().corePoolSize(10).maxPoolSize(20) }
注意:这种方式只是临时缓解,异步处理才是解决阻塞问题的根本方案。
内容的提问来源于stack exchange,提问作者Captain Jacky
相关产品推荐
相关产品推荐

