spring messaging是一套统一,抽象的消息编程模型,,为了屏蔽不同消息中间件(rabbitMq,activeMq)。不同通讯协议比如(stomp,jms,AMQP)的底层差异,,用一致的方式编写异步消息通讯代码,,无论是做websocket,还是消息队列消费,还是跨服务异步通知,都能复用核心的API

spring messaging是spring对消息通讯的标准化封装

  • channel: 通道,只负责传递消息,,不处理数据,,channel就是去找到所有的订阅者,匹配的就转发
  • handler: 不负责传消息,只负责处理消息,,处理完了之后可以把他丢入下一个通道中,,最后通过outboundChannel转发到客户端

核心类:

  • Message

  • MessageChannel : 消息通道

    • ExecutorSubscribableChannel : inboundChannel,outboundChannel,brokerChannel都是这个类型,, 可以配置通道的拦截器,决定通道的消息转发到哪个MessageHandler进行处理
  • MessageHandler : 消息处理器

    • SimpAnnotationMethodMessageHandler : 处理@MessageMapping的处理器
    • SimpleBrokerMessageHandler : 内存消息代理broker,,维护订阅,群发消息
    • UserDestinationMessageHandler : 处理私聊,需要将user解析成具体的session,定向发送
    • StompBrokerRelayMessageHandler : 外部broker (RabbitMQ / ActiveMQ),spring不自己处理,转发给mq处理
    • SubProtocolWebSocketHandler : 负责协议的转换
  • ChannelInterceptor : 拦截器 ,加在通道上,用来鉴权,日志,限流


  @MessageMapping("/greeting")
    @SendTo("/topic/greeting")  // 转发到订阅这个频道的人
    public Message greeting(Message message){
        System.out.println("message = " + message);
        message.setTo("public");
        return message;
    }

比如说 SimpAnnotationMethodMessageHandler处理之后,,如果有@SendTo就会转发到 broker的channel中,,然后SimpleBrokerMessageHandler再去处理 brokerChannel中的数据,,最后转发到 outboundChannel中


代码:

@Data
public class Message {
    private String destination;
    String payload;
    Map<String,Object> headers = new HashMap<>();

    public Message() {
    }

    public Message(String destination, String payload) {
        this.destination = destination;
        this.payload = payload;
    }
}

拦截器:

public interface ChannelInterceptor {

    /**
     * 消息正式发送到通道之前,,,进行拦截  ===》  校验消息合法性,,修改消息内容,记录发送前日志,权限校验
     * @param message
     *         返回原message  :  消息正常发送
     *         返回修改后的message : 消息被篡改
     *         返回null :  中断消息发送
     * @return
     */
    Message preSend(Message message);


    /**
     * 消息成功发送到通道之后,,, ===》 无论消费者是否处理,,只要消息进入通道之后就会执行
     *
     *   记录消息发送成功的日志,,, 统计消息发送量,记录发送耗时
     *   发送后置通知:  比如消息发送成功后,通知监控系统
     *   清理临时资源,,比如清理发送过程中创建的临时缓存
     * @param message
     */
    void postSend(Message message);
}
public class AuthInterceptor implements ChannelInterceptor{
    @Override
    public Message preSend(Message message) {
        System.out.println("[auth] 校验token");
        message.getHeaders().put("user","zs");
        return message;
    }

    @Override
    public void postSend(Message message) {

        System.out.println("[auth校验完成]");
    }
}

管道:

public interface MessageChannel {

    void send(Message message);
}

abstract class AbstractSubscribableChannel implements MessageChannel{

    protected List<MessageHandler> handlers = new ArrayList<>();
    protected List<ChannelInterceptor> interceptors = new CopyOnWriteArrayList<>();


    public void subscribe(MessageHandler handler){
        handlers.add(handler);
    }

    public void addInterceptor(ChannelInterceptor interceptor){
        interceptors.add(interceptor);
    }


    /**
     * 消息进入通道之前
     * @param message
     * @return
     */
    protected Message applyPreSend(Message message){
        for (ChannelInterceptor interceptor : interceptors) {
            message = interceptor.preSend(message);
            if (message == null){
                return null;
            }
        }

        return message;
    }


    /**
     * 消息进入通道之后
     * @param message
     */
    protected void applyPostSend(Message message){
        for (ChannelInterceptor interceptor : interceptors) {
            interceptor.postSend(message);
        }
    }

}

public class ExecutorSubscribableChannel extends AbstractSubscribableChannel{

    private String name;

    public ExecutorSubscribableChannel(String name) {
        this.name = name;
    }

    @Override
    public void send(Message message) {
        // 进入通道之前执行
       message = applyPreSend(message);
       if (message == null){
           return;
       }

        System.out.println("[channel : ]"+name+"===>"+message.getDestination());

        for (MessageHandler handler : handlers) {
            if (handler.supports(message)) {
                handler.handleMessage(message);
                // channel处理之后的拦截
                applyPostSend(message);
                return;
            }
        }


        System.out.println("[channel:]"+name+"无处理器");

    }
}

public class ClientOutboundChannel implements MessageChannel{
    @Override
    public void send(Message message) {
        System.out.println("[client outbound channel] 推送==》"+message.getPayload());
    }
}

处理器:


public interface MessageHandler {

    boolean supports(Message message);

    void handleMessage(Message message);
}
abstract class AbstractBrokerMessageHandler implements MessageHandler{

    protected MessageChannel clientOutboundChannel;

    public AbstractBrokerMessageHandler(MessageChannel clientOutboundChannel) {
        this.clientOutboundChannel = clientOutboundChannel;
    }

}
public class SimpAnnotationMethodMessageHandler implements MessageHandler {

    private MessageChannel brokerChannel;


    public SimpAnnotationMethodMessageHandler(MessageChannel brokerChannel) {
        this.brokerChannel = brokerChannel;
    }

    @Override
    public boolean supports(Message message) {
        return message.getDestination().startsWith("/app");
    }

    @Override
    public void handleMessage(Message message) {
        System.out.println("[@MessageMapping] 处理:"+message.getPayload());

        // 模拟 @SendTo
        brokerChannel.send(new Message("/topic/chat","结果:"+message.getPayload()));


        brokerChannel.send(new Message("/mq/test","mq:"+message.getPayload()));


    }
}

public class UserDestinationMessageHandler implements MessageHandler{

    private MessageChannel brokerChannel;

    public UserDestinationMessageHandler(MessageChannel brokerChannel) {
        this.brokerChannel = brokerChannel;
    }

    @Override
    public boolean supports(Message message) {
        return message.getDestination().startsWith("/user");
    }

    @Override
    public void handleMessage(Message message) {


        String target = message.getDestination().replace("/user", "/queue");


        System.out.println("[user destination handler]转发到"+target);

        brokerChannel.send(new Message(target,message.getPayload()));

    }
}

public class SimpleBrokerMessageHandler extends AbstractBrokerMessageHandler{
    private Map<String, List<String>> subscribers = new HashMap<>();

    public SimpleBrokerMessageHandler(MessageChannel clientOutboundChannel) {
        super(clientOutboundChannel);
    }

    public void subscribe(String destination,String user){
        subscribers.computeIfAbsent(destination,k->new ArrayList<>()).add(user);
    }

    @Override
    public boolean supports(Message message) {
        return message.getDestination().startsWith("/topic") || message.getDestination().startsWith("/queue");
    }

    @Override
    public void handleMessage(Message message) {
        System.out.println("[simple broker]广播"+message.getDestination());

        List<String> users = subscribers.getOrDefault(message.getDestination(), List.of());

        for (String user : users) {
            clientOutboundChannel.send(new Message(message.getDestination(),user+"收到消息"+message.getPayload()));
        }
    }
}
public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler{


    public StompBrokerRelayMessageHandler(MessageChannel clientOutboundChannel) {
        super(clientOutboundChannel);
    }

    @Override
    public boolean supports(Message message) {
        return message.getDestination().startsWith("/mq");
    }

    @Override
    public void handleMessage(Message message) {
        System.out.println("[发送到mq处理:]"+message.getPayload());

        clientOutboundChannel.send(new Message(message.getDestination(),"mq返回:"+message.getPayload()));
    }
}

template:

public class SimpMessagingTemplate {


    private MessageChannel brokerChannel;

    public SimpMessagingTemplate(MessageChannel brokerChannel) {
        this.brokerChannel = brokerChannel;
    }


    public void convertAndSend(String dest,String payload){
        brokerChannel.send(new Message(dest,payload));
    }
}

测试:

    public static void main(String[] args) {
        ExecutorSubscribableChannel inbound = new ExecutorSubscribableChannel("inbound");
        ClientOutboundChannel outbound = new ClientOutboundChannel();
        ExecutorSubscribableChannel broker = new ExecutorSubscribableChannel("broker");

        // interceptor ,,, 只有inbound的 channel有拦截器,,其他的没有拦截器
        inbound.addInterceptor(new AuthInterceptor());


        // 处理器
        SimpAnnotationMethodMessageHandler controller = new SimpAnnotationMethodMessageHandler(broker);

        UserDestinationMessageHandler userHandler = new UserDestinationMessageHandler(broker);

        SimpleBrokerMessageHandler simpleBroker = new SimpleBrokerMessageHandler(outbound);

        StompBrokerRelayMessageHandler mq = new StompBrokerRelayMessageHandler(outbound);



        // 订阅通道,,将通道消息转发到  handler中
        inbound.subscribe(controller);
        inbound.subscribe(userHandler);

        //  handle订阅管道中的消息
        broker.subscribe(simpleBroker);
        broker.subscribe(mq);


        // 用户去订阅  broker中的消息
        simpleBroker.subscribe("/topic/chat","用户A");
        simpleBroker.subscribe("/topic/chat","用户C");
        simpleBroker.subscribe("/queue/msg","用户B");


        SimpMessagingTemplate template = new SimpMessagingTemplate(broker);




//        System.out.println("发送消息");
//        inbound.send(new Message("/app/chat","hello"));



        inbound.send(new Message("/user/msg","私聊消息"));


        template.convertAndSend("/topic/chat","系统消息");



    }
}

每个通道都会有自己的 拦截器 和 订阅者(消息处理者),,, springboot中websocket,,有三个核心的channel,,inboundChannel ,brokerChannel ,outboundChannel,
消息通过inboundChannel进去,之后要么走@MessageMapping ,要么直接进入brokerChannel转发,,
最终需要推送给客户端的消息都会从 outboundChannel发出去

@MessageMapping标记的方法没有返回值或者没有发送,,就不会进入brokerchannel,,也就结束了

公众号:代码源记

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐