
作者 | 洛夜
來源 |
阿裡巴巴雲原生公衆号Spring Cloud Stream在 Spring Cloud 體系内用于建構高度可擴充的基于事件驅動的微服務,其目的是為了簡化消息在 Spring Cloud 應用程式中的開發。
Spring Cloud Stream (後面以 SCS 代替 Spring Cloud Stream) 本身内容很多,而且它還有很多外部的依賴,想要熟悉 SCS,必須要先了解 Spring Messaging 和 Spring Integration 這兩個項目,接下來,文章将圍繞以下三點進行展開:
- 什麼是 Spring Messaging
- 什麼是 Spring Integration
- 什麼是 SCS 體系及其原理
本文配套可互動教程已登入阿裡雲知行動手實驗室,PC 端登入
start.aliyun.com_ _在浏覽器中立即體驗。
Spring Messaging
Spring Messaging 是 Spring Framework 中的一個子產品,其作用就是統一消息的程式設計模型。
- 比如消息 Messaging 對應的模型就包括一個消息體 Payload 和消息頭 Header:
package org.springframework.messaging;
public interface Message<T> {
T getPayload();
MessageHeaders getHeaders();
}
- 消息通道 MessageChannel 用于接收消息,調用send方法可以将消息發送至該消息通道中:
@FunctionalInterface
public interface MessageChannel {
long INDEFINITE_TIMEOUT = -1;
default boolean send(Message<?> message) {
return send(message, INDEFINITE_TIMEOUT);
}
boolean send(Message<?> message, long timeout);
}
消息通道裡的消息如何被消費呢?
- 由消息通道的子接口可訂閱的消息通道SubscribableChannel實作,被MessageHandler消息處理器所訂閱:
public interface SubscribableChannel extends MessageChannel {
boolean subscribe(MessageHandler handler);
boolean unsubscribe(MessageHandler handler);
}
- 由MessageHandler真正地消費/處理消息:
@FunctionalInterface
public interface MessageHandler {
void handleMessage(Message<?> message) throws MessagingException;
}
Spring Messaging 内部在消息模型的基礎上衍生出了其它的一些功能,如:
- 消息接收參數及傳回值處理:消息接收參數處理器HandlerMethodArgumentResolver配合@Header, @Payload等注解使用;消息接收後的傳回值處理器HandlerMethodReturnValueHandler配合@SendTo注解使用;
- 消息體内容轉換器MessageConverter;
- 統一抽象的消息發送模闆AbstractMessageSendingTemplate;
- 消息通道攔截器ChannelInterceptor;
Spring Integration
Spring Integration 提供了 Spring 程式設計模型的擴充用來支援企業內建模式(Enterprise Integration Patterns),是對 Spring Messaging 的擴充。
它提出了不少新的概念,包括消息路由MessageRoute、消息分發MessageDispatcher、消息過濾Filter、消息轉換Transformer、消息聚合Aggregator、消息分割Splitter等等。同時還提供了MessageChannel和MessageHandler的實作,分别包括 DirectChannel、ExecutorChannel、PublishSubscribeChannel和MessageFilter、ServiceActivatingHandler、MethodInvokingSplitter 等内容。
這裡為大家介紹幾種消息的處理方式:
- 消息的分割:
- 消息的聚合:
- 消息的過濾:
- 消息的分發:
接下來,我們以一個最簡單的例子來嘗試一下 Spring Integration。
這段代碼解釋為:
SubscribableChannel messageChannel =new DirectChannel(); // 1
messageChannel.subscribe(msg-> { // 2
System.out.println("receive: " +msg.getPayload());
});
messageChannel.send(MessageBuilder.withPayload("msgfrom alibaba").build()); // 3
- 構造一個可訂閱的消息通道messageChannel。
- 使用MessageHandler去消費這個消息通道裡的消息。
- 發送一條消息到這個消息通道,消息最終被消息通道裡的MessageHandler所消費。
- 最後控制台列印出:receive: msg from alibaba。
DirectChannel内部有個UnicastingDispatcher類型的消息分發器,會分發到對應的消息通道MessageChannel中,從名字也可以看出來,UnicastingDispatcher是個單點傳播的分發器,隻能選擇一個消息通道。那麼如何選擇呢? 内部提供了LoadBalancingStrategy負載均衡政策,預設隻有輪詢的實作,可以進行擴充。
我們對上段代碼做一點修改,使用多個 MessageHandler 去處理消息:
SubscribableChannel messageChannel = new DirectChannel();
messageChannel.subscribe(msg -> {
System.out.println("receive1: " + msg.getPayload());
});
messageChannel.subscribe(msg -> {
System.out.println("receive2: " + msg.getPayload());
});
messageChannel.send(MessageBuilder.withPayload("msg from alibaba").build());
messageChannel.send(MessageBuilder.withPayload("msg from alibaba").build());
由于DirectChannel内部的消息分發器是UnicastingDispatcher單點傳播的方式,并且采用輪詢的負載均衡政策,是以這裡兩次的消費分别對應這兩個MessageHandler。控制台列印出:
receive1: msg from alibaba
receive2: msg from alibaba
既然存在單點傳播的消息分發器UnicastingDispatcher,必然也會存在廣播的消息分發器,那就是BroadcastingDispatcher,它被 PublishSubscribeChannel 這個消息通道所使用。廣播消息分發器會把消息分發給所有的 MessageHandler:
SubscribableChannel messageChannel = new PublishSubscribeChannel();
messageChannel.subscribe(msg -> {
System.out.println("receive1: " + msg.getPayload());
});
messageChannel.subscribe(msg -> {
System.out.println("receive2: " + msg.getPayload());
});
messageChannel.send(MessageBuilder.withPayload("msg from alibaba").build());
messageChannel.send(MessageBuilder.withPayload("msg from alibaba").build());
Spring Cloud Stream
SCS 與各子產品之間的關系是:
- SCS 在 Spring Integration 的基礎上進行了封裝,提出了Binder, Binding, @EnableBinding, @StreamListener等概念。
- SCS 與 Spring Boot Actuator 整合,提供了/bindings, /channelsendpoint。
- SCS 與 Spring Boot Externalized Configuration 整合,提供了BindingProperties, BinderProperties等外部化配置類。
- SCS 增強了消息發送失敗的和消費失敗情況下的處理邏輯等功能。
- SCS 是 Spring Integration 的加強,同時與 Spring Boot 體系進行了融合,也是 Spring Cloud Bus 的基礎。它屏蔽了底層消息中間件的實作細節,希望以統一的一套 API 來進行消息的發送/消費,底層消息中間件的實作細節由各消息中間件的 Binder 完成。
Binder是提供與外部消息中間件內建的元件,為構造Binding提供了 2 個方法,分别是bindConsumer和bindProducer,它們分别用于構造生産者和消費者。目前官方的實作有 Rabbit Binder 和 Kafka Binder, Spring Cloud Alibaba 内部已經實作了 RocketMQ Binder。
從圖中可以看出,Binding是連接配接應用程式跟消息中間件的橋梁,用于消息的消費和生産。我們來看一個最簡單的使用 RocketMQ Binder 的例子,然後分析一下它的底層處理原理:
- 啟動類及消息的發送:
@SpringBootApplication
@EnableBinding({ Source.class, Sink.class }) // 1
public class SendAndReceiveApplication {
public static void main(String[] args) {
SpringApplication.run(SendAndReceiveApplication.class, args);
}
@Bean // 2
public CustomRunner customRunner() {
return new CustomRunner();
}
public static class CustomRunner implements CommandLineRunner {
@Autowired
private Source source;
@Override
public void run(String... args) throws Exception {
int count = 5;
for (int index = 1; index <= count; index++) {
source.output().send(MessageBuilder.withPayload("msg-" + index).build()); // 3
}
}
}
}
- 消息的接收:
@Service
public class StreamListenerReceiveService {
@StreamListener(Sink.INPUT) // 4
public void receiveByStreamListener1(String receiveMsg) {
System.out.println("receiveByStreamListener: " + receiveMsg);
}
}
這段代碼很簡單,沒有涉及到 RocketMQ 相關的代碼,消息的發送和接收都是基于 SCS 體系完成的。如果想切換成 RabbitMQ 或 Kafka,隻需修改配置檔案即可,代碼無需修改。
我們來分析下這段代碼的原理:
1.@EnableBinding對應的兩個接口屬性Source和Sink是 SCS 内部提供的。SCS 内部會基于Source和Sink構造BindableProxyFactory,且對應的 output 和 input 方法傳回的 MessageChannel 是DirectChannel。output 和 input 方法修飾的注解對應的 value 是配置檔案中 binding 的 name。
public interface Source {
String OUTPUT = "output";
@Output(Source.OUTPUT)
MessageChannel output();
}
public interface Sink {
String INPUT = "input";
@Input(Sink.INPUT)
SubscribableChannel input();
}
配置檔案裡 bindings 的 name 為 output 和 input,對應Source和Sink接口的方法上的注解裡的 value:
spring.cloud.stream.bindings.output.destination=test-topic
spring.cloud.stream.bindings.output.content-type=text/plain
spring.cloud.stream.rocketmq.bindings.output.producer.group=demo-group
spring.cloud.stream.bindings.input.destination=test-topic
spring.cloud.stream.bindings.input.content-type=text/plain
spring.cloud.stream.bindings.input.group=test-group1
- 構造CommandLineRunner,程式啟動的時候會執行CustomRunner的run方法。
- 調用Source接口裡的 output 方法擷取DirectChannel,并發送消息到這個消息通道中。這裡跟之前 Spring Integration 章節裡的代碼一緻。
- Source 裡的 output 發送消息到DirectChannel消息通道之後會被AbstractMessageChannelBinder#SendingHandler這個MessageHandler處理,然後它會委托給AbstractMessageChannelBinder#createProducerMessageHandler建立的 MessageHandler 處理(該方法由不同的消息中間件實作)。
- 不同的消息中間件對應的AbstractMessageChannelBinder#createProducerMessageHandler方法傳回的 MessageHandler 内部會把 Spring Message 轉換成對應中間件的 Message 模型并發送到對應中間件的 broker。
- 使用@StreamListener進行消息的訂閱。請注意,注解裡的Sink.input對應的值是 "input",會根據配置檔案裡 binding 對應的 name 為 input 的值進行配置:
- 不同的消息中間件對應的AbstractMessageChannelBinder#createConsumerEndpoint方法會使用 Consumer 訂閱消息,訂閱到消息後内部會把中間件對應的 Message 模型轉換成 Spring Message。
- 消息轉換之後會把 Spring Message 發送至 name 為 input 的消息通道中。
- @StreamListener對應的StreamListenerMessageHandler訂閱了 name 為 input 的消息通道,進行了消息的消費。
這個過程文字描述有點啰嗦,用一張圖總結一下(黃色部分涉及到各消息中間件的 Binder 實作以及 MQ 基本的訂閱釋出功能):
SCS 章節的最後,我們來看一段 SCS 關于消息的處理方式的一段代碼:
@StreamListener(value = Sink.INPUT, condition = "headers['index']=='1'")
public void receiveByHeader(Message msg) {
System.out.println("receive by headers['index']=='1': " + msg);
}
@StreamListener(value = Sink.INPUT, condition = "headers['index']=='9999'")
public void receivePerson(@Payload Person person) {
System.out.println("receive Person: " + person);
}
@StreamListener(value = Sink.INPUT)
public void receiveAllMsg(String msg) {
System.out.println("receive allMsg by StreamListener. content: " + msg);
}
@StreamListener(value = Sink.INPUT)
public void receiveHeaderAndMsg(@Header("index") String index, Message msg) {
System.out.println("receive by HeaderAndMsg by StreamListener. content: " + msg);
}
有沒有發現這段代碼跟 Spring MVC Controller 中接收請求的代碼很像? 實際上他們的架構都是類似的,Spring MVC 對于 Controller 中參數和傳回值的處理類分别是org.springframework.web.method.support.HandlerMethodArgumentResolver、org.springframework.web.method.support.HandlerMethodReturnValueHandler。
Spring Messaging 中對于參數和傳回值的處理類之前也提到過,分别是org.springframework.messaging.handler.invocation.HandlerMethodArgumentResolver、org.springframework.messaging.handler.invocation.HandlerMethodReturnValueHandler。
它們的類名一模一樣,甚至内部的方法名也一樣。
總結
上圖是 SCS 體系相關類說明的總結,關于 SCS 以及 RocketMQ Binder 更多相關的示例,可以參考
RocketMQ Binder Demos,包含了消息的聚合、分割、過濾;消息異常處理;消息标簽、SQL 過濾;同步、異步消費等等。
歡迎大家使用釘釘掃描二維碼加入 Spring Cloud Alibaba 開源讨論群: