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

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);
问题根因
  1. 主题路径命名不符合STOMP代理通用规范
    内置SimpleMessageBroker是Spring自研的轻量实现,会把包含查询串的完整路径作为主题名做精确匹配,因此conversations?orgId=A和conversations?orgId=B会被识别为两个完全独立的主题,不会出现消息串流。但Amazon MQ底层基于ActiveMQ实现,其STOMP协议解析逻辑不会将URL查询参数作为主题名称的一部分,所有带conversations?orgId=xxx格式的订阅都会被统一识别为对conversations根主题的订阅,因此该主题下的所有消息会被推送给所有订阅者。
  2. 拦截逻辑仅做订阅权限校验,未适配外置代理路由规则
    现有拦截器仅在客户端发起订阅时校验orgId访问权限,没有参与消息发送环节的路由过滤。内置broker靠路径精确匹配歪打正着实现了消息隔离,但外置代理完全按照自身的主题解析规则做路由,只要客户端被识别为订阅了根主题,就会收到全量消息,拦截器的权限校验无法影响代理侧的路由行为。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 13:19:07