AWS Amazon MQ STOMP中继场景下客户端收到非订阅消息异常
问题现象
基于Amazon MQ实现STOMP主题消息推送时出现消息过滤失效:使用内存实现的SimpleMessageBroker时,客户端仅能收到自身已订阅的匹配消息;切换为外置Amazon MQ消息代理后,所有客户端都会接收到通过SimpMessagingTemplate广播的全量消息,未按照订阅规则完成过滤。
相关实现代码
WebSocket基础配置
public class WebSocketConfig implements WebSocketMessageBrokerConfigurer { if ("true".equalsIgnoreCase(environment.getProperty("aws.activemq.enabled"))) config.setApplicationDestinationPrefixes("/app") .enableStompBrokerRelay("/topic") .setUserDestinationBroadcast("/user") .setRelayHost(this.activeMQCreds.getHost()) .setRelayPort(this.activeMQCreds.getPort()) .setClientLogin(this.activeMQCreds.getUsername()) .setClientPasscode(this.activeMQCreds.getPassword()) .setAutoStartup(true) .setSystemLogin(this.activeMQCreds.getUsername()) .setSystemPasscode(this.activeMQCreds.getPassword()) .setTcpClient(this.createClient()); else config.setApplicationDestinationPrefixes("/app") .enableSimpleBroker("/topic"); @Override public void configureClientInboundChannel(ChannelRegistration registration) { registration.interceptors(webSocketSessionChannelInterceptor); } private TcpOperations<byte[]> createClient() { return new ReactorNettyTcpClient<>( (client) -> client.remoteAddress(this::getAddress).secure(), new StompReactorNettyCodec()); } private SocketAddress getAddress() { try { InetAddress address = InetAddress.getByName(this.activeMQCreds.getHost().replace("stomp+ssl://", "")); return new InetSocketAddress(address, this.activeMQCreds.getPort()); } catch (UnknownHostException e) { log.error("Exception", e); return null; } } }
订阅权限拦截逻辑
public class WebSocketSessionChannelInterceptor implements ChannelInterceptor { @Override public Message<?> preSend(final Message<?> message, final MessageChannel channel) throws AuthenticationException { log.info("Message: {}", message); final StompHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class); try { final String destination = accessor.getDestination(); if (StompCommand.CONNECT == accessor.getCommand()) { final String token = accessor.getFirstNativeHeader("Authorization"); final UsernamePasswordAuthenticationToken user = webSocketAuthenticatorService.getAuthenticatedOrFail(token); if (user != null) SecurityContextHolder.getContext().setAuthentication(user); accessor.setUser(user); } else if (StompCommand.SUBSCRIBE == accessor.getCommand() && StringUtils.isNotEmpty(destination) && destination.contains("conversations?orgId=")) { String orgId = destination.substring(destination.lastIndexOf('=') + 1); if (WebSocketAuthenticatorService.hasPermissionToOrgId(orgId, (Authentication) accessor.getUser())) return message; else throw new BadCredentialsException( String.format("you do not have privilege to this organization: %s", orgId)); } return message; } catch (Exception e) { log.error("Exception: ", e); return null; } } @Override public void afterSendCompletion(Message<?> message, MessageChannel channel, boolean sent, @Nullable Exception ex) { final StompHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class); final String destination = accessor.getDestination(); if (StompCommand.SUBSCRIBE == accessor.getCommand() && StringUtils.isNotEmpty(destination) && destination.contains("conversations?orgId=")) { String orgId = destination.substring(destination.lastIndexOf('=') + 1); if (WebSocketAuthenticatorService.hasPermissionToOrgId(orgId, (Authentication) accessor.getUser())) webSocketUtils.publishConversationsForOrgId(orgId); } } }
消息发送逻辑
this.simpMessagingTemplate.convertAndSend( ApplicationConstants.ORG_WEBSOCKET_TOPIC_PATH + orgId, conversationBO);
问题根因
- 主题路径命名不符合STOMP代理通用规范
内置SimpleMessageBroker是Spring自研的轻量实现,会把包含查询串的完整路径作为主题名做精确匹配,因此conversations?orgId=A和conversations?orgId=B会被识别为两个完全独立的主题,不会出现消息串流。但Amazon MQ底层基于ActiveMQ实现,其STOMP协议解析逻辑不会将URL查询参数作为主题名称的一部分,所有带conversations?orgId=xxx格式的订阅都会被统一识别为对conversations根主题的订阅,因此该主题下的所有消息会被推送给所有订阅者。 - 拦截逻辑仅做订阅权限校验,未适配外置代理路由规则
现有拦截器仅在客户端发起订阅时校验orgId访问权限,没有参与消息发送环节的路由过滤。内置broker靠路径精确匹配歪打正着实现了消息隔离,但外置代理完全按照自身的主题解析规则做路由,只要客户端被识别为订阅了根主题,就会收到全量消息,拦截器的权限校验无法影响代理侧的路由行为。 - STOMP中继配置存在错误
配置中将/user设置为用户目的地广播地址,/user是Spring WebSocket默认的点对点用户消息前缀,占用该路径做广播路由会导致跨节点的用户消息、主题消息路由出现异常。
修复方案
- 调整主题路径命名规则,将orgId作为路径层级而非查询参数。将原有订阅路径
/topic/conversations?orgId={orgId}修改为/topic/conversations/{orgId}格式,例如/topic/conversations/1001,发送消息时也向对应层级路径发送。该写法完全兼容STOMP协议标准,所有合规STOMP代理都会将不同层级路径识别为独立主题,天然实现路径级别的消息隔离,无需额外过滤逻辑。 - 修正STOMP中继的用户路由配置,将用户目的地广播、用户注册表广播路径设置为独立的系统主题,避免占用业务前缀:
config.enableStompBrokerRelay("/topic") .setUserDestinationBroadcast("/topic/sys/unresolved-user") .setUserRegistryBroadcast("/topic/sys/user-registry") // 其余 relay 配置保持不变 - 若确实需要在单主题下按属性做细粒度过滤,使用STOMP标准的Selector机制实现:客户端订阅时在SUBSCRIBE帧添加
selector: orgId = 'xxx'请求头,发送消息时给消息头设置对应orgId属性,Amazon MQ会在服务端按照selector表达式完成消息过滤,仅推送匹配的消息给对应订阅者。该方案复杂度高于路径层级隔离,无特殊需求优先使用路径分层方案。
内容的提问来源于stack exchange,提问作者Jonas Schreiber
相关产品推荐
相关产品推荐

