最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Spring Cloud Stream如何實現(xiàn)服務(wù)之間的通訊

 更新時間:2019年10月15日 11:37:59   作者:維晟  
這篇文章主要介紹了Spring Cloud Stream如何實現(xiàn)服務(wù)之間的通訊,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下

Spring Cloud Stream

Srping cloud Bus的底層實現(xiàn)就是Spring Cloud Stream,Spring Cloud Stream的目的是用于構(gòu)建基于消息驅(qū)動(或事件驅(qū)動)的微服務(wù)架構(gòu)。Spring Cloud Stream本身對Spring Messaging、Spring Integration、Spring Boot Actuator、Spring Boot Externalized Configuration等模塊進行封裝(整合)和擴展,下面我們實現(xiàn)兩個服務(wù)之間的通訊來演示Spring Cloud Stream的使用方法。

整體概述

服務(wù)要想與其他服務(wù)通訊要定義通道,一般會定義輸出通道和輸入通道,輸出通道用于發(fā)送消息,輸入通道用于接收消息,每個通道都會有個名字(輸入和輸出只是通道類型,可以用不同的名字定義很多很多通道),不同通道的名字不能相同否則會報錯(輸入通道和輸出通道不同類型的通道名稱也不能相同),綁定器是操作RabbitMQ或Kafka的抽象層,為了屏蔽操作這些消息中間件的復(fù)雜性和不一致性,綁定器會用通道的名字在消息中間件中定義主題,一個主題內(nèi)的消息生產(chǎn)者來自多個服務(wù),一個主題內(nèi)消息的消費者也是多個服務(wù),也就是說消息的發(fā)布和消費是通過主題進行定義和組織的,通道的名字就是主題的名字,在RabbitMQ中主題使用Exchanges實現(xiàn),在Kafka中主題使用Topic實現(xiàn)。

準(zhǔn)備環(huán)境

創(chuàng)建兩個項目spring-cloud-stream-a和spring-cloud-stream-b,spring-cloud-stream-a我們用Spring Cloud Stream實現(xiàn)通訊,spring-cloud-stream-b我們用Spring Cloud Stream的底層模塊Spring Integration實現(xiàn)通訊。

兩個項目的POM文件依賴都是:

<dependencies>
    <dependency>
      <groupId>org.springframework.cloud</groupId>
      <artifactId>spring-cloud-stream</artifactId>
    </dependency>

    <dependency>
      <groupId>org.springframework.cloud</groupId>
      <artifactId>spring-cloud-stream-binder-rabbit</artifactId>
    </dependency>
    <dependency>
      <groupId>org.springframework.boot</groupId>
      <artifactId>spring-boot-starter-test</artifactId>
      <scope>test</scope>
    </dependency>
    <dependency>
      <groupId>org.springframework.cloud</groupId>
      <artifactId>spring-cloud-stream-test-support</artifactId>
      <scope>test</scope>
    </dependency>
  </dependencies>

spring-cloud-stream-binder-rabbit是指綁定器的實現(xiàn)使用RabbitMQ。

項目配置內(nèi)容application.properties:

spring.application.name=spring-cloud-stream-a
server.port=9010

#設(shè)置默認(rèn)綁定器
spring.cloud.stream.defaultBinder = rabbit

spring.rabbitmq.host=127.0.0.1
spring.rabbitmq.port=5672
spring.rabbitmq.username=guest
spring.rabbitmq.password=guest
spring.application.name=spring-cloud-stream-b
server.port=9011

#設(shè)置默認(rèn)綁定器
spring.cloud.stream.defaultBinder = rabbit

spring.rabbitmq.host=127.0.0.1
spring.rabbitmq.port=5672
spring.rabbitmq.username=guest
spring.rabbitmq.password=guest

啟動一個rabbitmq:

docker pull rabbitmq:3-management
docker run -d --hostname my-rabbit --name rabbit -p 5672:5672 -p 15672:15672 rabbitmq:3-management

編寫A項目代碼

在A項目中定義一個輸入通道一個輸出通道,定義通道在接口中使用@Input和@Output注解定義,程序啟動的時候Spring Cloud Stream會根據(jù)接口定義將實現(xiàn)類自動注入(Spring Cloud Stream自動實現(xiàn)該接口不需要寫代碼)。

A服務(wù)輸入通道,通道名稱ChatExchanges.A.Input,接口定義輸入通道必須返回SubscribableChannel:

public interface ChatInput {
  String INPUT = "ChatExchanges.A.Input";
  @Input(ChatInput.INPUT)
  SubscribableChannel input();
}

A服務(wù)輸出通道,通道名稱ChatExchanges.A.Output,輸出通道必須返回MessageChannel:

public interface ChatOutput {

  String OUTPUT = "ChatExchanges.A.Output";

  @Output(ChatOutput.OUTPUT)
  MessageChannel output();
}

定義消息實體類:

public class ChatMessage implements Serializable {

  private String name;
  private String message;
  private Date chatDate;

  //沒有無參數(shù)的構(gòu)造函數(shù)并行化會出錯
  private ChatMessage(){}

  public ChatMessage(String name,String message,Date chatDate){
    this.name = name;
    this.message = message;
    this.chatDate = chatDate;
  }

  public String getName(){
    return this.name;
  }

  public String getMessage(){
    return this.message;
  }

  public Date getChatDate() { return this.chatDate; }

  public String ShowMessage(){
    return String.format("聊天消息:%s的時候,%s說%s。",this.chatDate,this.name,this.message);
  }
}

在業(yè)務(wù)處理類上用@EnableBinding注解綁定輸入通道和輸出通道,這個綁定動作其實就是創(chuàng)建并注冊輸入和輸出通道的實現(xiàn)類到Bean中,所以可以直接是使用@Autowired進行注入使用,另外消息的串行化默認(rèn)使用application/json格式(com.fastexml.jackson),最后用@StreamListener注解進行指定通道消息的監(jiān)聽:

//ChatInput.class的輸入通道不在這里綁定,監(jiān)聽到數(shù)據(jù)會找不到AClient類的引用。
//Input和Output通道定義的名字不能一樣,否則程序啟動會拋異常。
@EnableBinding({ChatOutput.class,ChatInput.class})
public class AClient {

  private static Logger logger = LoggerFactory.getLogger(AClient.class);

  @Autowired
  private ChatOutput chatOutput;

  //StreamListener自帶了Json轉(zhuǎn)對象的能力,收到B的消息打印并回復(fù)B一個新的消息。
  @StreamListener(ChatInput.INPUT)
  public void PrintInput(ChatMessage message) {

    logger.info(message.ShowMessage());

    ChatMessage replyMessage = new ChatMessage("ClientA","A To B Message.", new Date());

    chatOutput.output().send(MessageBuilder.withPayload(replyMessage).build());
  }
}

到此A項目代碼編寫完成。

編寫B(tài)項目代碼

B項目使用Spring Integration實現(xiàn)消息的發(fā)布和消費,定義通道時我們要交換輸入通道和輸出通道的名稱:

public interface ChatProcessor {

  String OUTPUT = "ChatExchanges.A.Input";
  String INPUT = "ChatExchanges.A.Output";

  @Input(ChatProcessor.INPUT)
  SubscribableChannel input();

  @Output(ChatProcessor.OUTPUT)
  MessageChannel output();
}

消息實體類:

public class ChatMessage {
  private String name;
  private String message;
  private Date chatDate;

  //沒有無參數(shù)的構(gòu)造函數(shù)并行化會出錯
  private ChatMessage(){}

  public ChatMessage(String name,String message,Date chatDate){
    this.name = name;
    this.message = message;
    this.chatDate = chatDate;
  }

  public String getName(){
    return this.name;
  }

  public String getMessage(){
    return this.message;
  }

  public Date getChatDate() { return this.chatDate; }

  public String ShowMessage(){
    return String.format("聊天消息:%s的時候,%s說%s。",this.chatDate,this.name,this.message);
  }
}

業(yè)務(wù)處理類用@ServiceActivator注解代替@StreamListener,用@InboundChannelAdapter注解發(fā)布消息:

@EnableBinding(ChatProcessor.class)
public class BClient {

  private static Logger logger = LoggerFactory.getLogger(BClient.class);

  //@ServiceActivator沒有Json轉(zhuǎn)對象的能力需要借助@Transformer注解
  @ServiceActivator(inputChannel=ChatProcessor.INPUT)
  public void PrintInput(ChatMessage message) {

    logger.info(message.ShowMessage());
  }

  @Transformer(inputChannel = ChatProcessor.INPUT,outputChannel = ChatProcessor.INPUT)
  public ChatMessage transform(String message) throws Exception{
    ObjectMapper objectMapper = new ObjectMapper();
    return objectMapper.readValue(message,ChatMessage.class);
  }

  //每秒發(fā)出一個消息給A
  @Bean
  @InboundChannelAdapter(value = ChatProcessor.OUTPUT,poller = @Poller(fixedDelay="1000"))
  public GenericMessage<ChatMessage> SendChatMessage(){
    ChatMessage message = new ChatMessage("ClientB","B To A Message.", new Date());
    GenericMessage<ChatMessage> gm = new GenericMessage<>(message);
    return gm;
  }
}

運行程序

啟動A項目和B項目:


源碼

Github倉庫:https://github.com/sunweisheng/spring-cloud-example

以上就是本文的全部內(nèi)容,希望對大家的學(xué)習(xí)有所幫助,也希望大家多多支持腳本之家。

相關(guān)文章

  • java中hashmap容量的初始化實現(xiàn)

    java中hashmap容量的初始化實現(xiàn)

    這篇文章主要介紹了java中hashmap容量的初始化實現(xiàn),文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2019-11-11
  • java設(shè)計模式之觀察者模式

    java設(shè)計模式之觀察者模式

    這篇文章主要為大家詳細(xì)介紹了java設(shè)計模式之觀察者模式的相關(guān)資料,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2016-12-12
  • 關(guān)于mybatisPlus?yml配置方式

    關(guān)于mybatisPlus?yml配置方式

    這篇文章主要介紹了mybatisPlus?yml配置方式,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2022-06-06
  • Spring Cloud Feign統(tǒng)一設(shè)置驗證token實現(xiàn)方法解析

    Spring Cloud Feign統(tǒng)一設(shè)置驗證token實現(xiàn)方法解析

    這篇文章主要介紹了Spring Cloud Feign統(tǒng)一設(shè)置驗證token實現(xiàn)方法解析,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下
    2020-08-08
  • 解決Springboot項目打包后的頁面丟失問題(thymeleaf報錯)

    解決Springboot項目打包后的頁面丟失問題(thymeleaf報錯)

    這篇文章主要介紹了解決Springboot項目打包后的頁面丟失問題(thymeleaf報錯),具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-11-11
  • 如何把JAR發(fā)布到maven中央倉庫的幾種方法

    如何把JAR發(fā)布到maven中央倉庫的幾種方法

    這篇文章主要介紹了如何把JAR發(fā)布到maven中央倉庫的幾種方法,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-05-05
  • Java查看線程運行狀態(tài)的方法詳解

    Java查看線程運行狀態(tài)的方法詳解

    這篇文章主要為大家詳細(xì)介紹了Java語言如何查看線程運行狀態(tài)的方法,文中的示例代碼講解詳細(xì),感興趣的小伙伴可以跟隨小編一起學(xué)習(xí)一下
    2022-08-08
  • 如何把VS Code打造成Java開發(fā)IDE

    如何把VS Code打造成Java開發(fā)IDE

    這篇文章主要介紹了如何把VS Code打造成Java開發(fā)IDE,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下
    2020-10-10
  • Java實現(xiàn)讀取及生成Excel文件的方法

    Java實現(xiàn)讀取及生成Excel文件的方法

    這篇文章主要介紹了Java實現(xiàn)讀取及生成Excel文件的方法,結(jié)合實例形式分析了java通過引入第三方j(luò)ar包poi-3.0.1-FINAL-20070705.jar實現(xiàn)針對Excel文件的讀取及生成功能,需要的朋友可以參考下
    2017-12-12
  • Java中Map遍歷的九種方式匯總

    Java中Map遍歷的九種方式匯總

    這篇文章主要介紹了Java中九種?Map?的遍歷方式匯總的相關(guān)資料,需要的朋友可以參考下
    2022-11-11

最新評論

佳木斯市| 华亭县| 洛浦县| 洱源县| 新竹县| 大洼县| 集贤县| 杨浦区| 汤阴县| 义马市| 天水市| 宜兴市| 德兴市| 五大连池市| 西盟| 东乡县| 南投县| 肥城市| 区。| 增城市| 三明市| 五河县| 壶关县| 柞水县| 体育| 南澳县| 长春市| 巍山| 广州市| 奈曼旗| 沙坪坝区| 全州县| 双城市| 天气| 新泰市| 古浪县| 望江县| 蒲城县| 南汇区| 祁阳县| 乐昌市|