spring Messaging
spring messaging是一套统一,抽象的消息编程模型,,为了屏蔽不同消息中间件(rabbitMq,activeMq)。不同通讯协议比如(stomp,jms,AMQP)的底层差异,,用一致的方式编写异步消息通讯代码,,无论是做websocket,还是消息队列消费,还是跨服务异步通知,都能复用核心的API
spring messaging是spring对消息通讯的标准化封装
- channel: 通道,只负责传递消息,,不处理数据,,channel就是去找到所有的订阅者,匹配的就转发
- handler: 不负责传消息,只负责处理消息,,处理完了之后可以把他丢入下一个通道中,,最后通过outboundChannel转发到客户端
核心类:
-
Message
-
MessageChannel : 消息通道
- ExecutorSubscribableChannel :
inboundChannel,outboundChannel,brokerChannel都是这个类型,, 可以配置通道的拦截器,决定通道的消息转发到哪个MessageHandler进行处理
- ExecutorSubscribableChannel :
-
MessageHandler : 消息处理器
- SimpAnnotationMethodMessageHandler : 处理
@MessageMapping的处理器 - SimpleBrokerMessageHandler : 内存消息代理broker,,维护订阅,群发消息
- UserDestinationMessageHandler : 处理私聊,需要将user解析成具体的session,定向发送
- StompBrokerRelayMessageHandler : 外部broker (RabbitMQ / ActiveMQ),spring不自己处理,转发给mq处理
- SubProtocolWebSocketHandler : 负责协议的转换
- SimpAnnotationMethodMessageHandler : 处理
-
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,,也就结束了
公众号:代码源记
更多推荐



所有评论(0)