精品专区-精品自拍9-精品自拍三级乱伦-精品自拍视频-精品自拍视频曝光-精品自拍小视频

網站建設資訊

NEWS

網站建設資訊

SpringCloudStream怎么實現服務之間的通訊

小編給大家分享一下Spring Cloud Stream怎么實現服務之間的通訊,希望大家閱讀完這篇文章之后都有所收獲,下面讓我們一起去探討吧!

創新互聯建站是專業的博山網站建設公司,博山接單;提供網站制作、成都網站制作,網頁設計,網站設計,建網站,PHP網站建設等專業做網站服務;采用PHP框架,可快速的進行博山網站開發網頁制作和功能擴展;專業做搜索引擎喜愛的網站,專業的做網站團隊,希望更多企業前來合作!

Spring Cloud Stream

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

整體概述

Spring Cloud Stream怎么實現服務之間的通訊

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

準備環境

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

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


    
      org.springframework.cloud
      spring-cloud-stream
    

    
      org.springframework.cloud
      spring-cloud-stream-binder-rabbit
    
    
      org.springframework.boot
      spring-boot-starter-test
      test
    
    
      org.springframework.cloud
      spring-cloud-stream-test-support
      test
    
  

spring-cloud-stream-binder-rabbit是指綁定器的實現使用RabbitMQ。

項目配置內容application.properties:

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

#設置默認綁定器
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

#設置默認綁定器
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會根據接口定義將實現類自動注入(Spring Cloud Stream自動實現該接口不需要寫代碼)。

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

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

A服務輸出通道,通道名稱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;

  //沒有無參數的構造函數并行化會出錯
  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);
  }
}

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

//ChatInput.class的輸入通道不在這里綁定,監聽到數據會找不到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轉對象的能力,收到B的消息打印并回復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項目代碼

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

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;

  //沒有無參數的構造函數并行化會出錯
  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);
  }
}

業務處理類用@ServiceActivator注解代替@StreamListener,用@InboundChannelAdapter注解發布消息:

@EnableBinding(ChatProcessor.class)
public class BClient {

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

  //@ServiceActivator沒有Json轉對象的能力需要借助@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);
  }

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

運行程序

啟動A項目和B項目:

Spring Cloud Stream怎么實現服務之間的通訊

Spring Cloud Stream怎么實現服務之間的通訊

看完了這篇文章,相信你對“Spring Cloud Stream怎么實現服務之間的通訊”有了一定的了解,如果想了解更多相關知識,歡迎關注創新互聯行業資訊頻道,感謝各位的閱讀!


文章標題:SpringCloudStream怎么實現服務之間的通訊
文章分享:http://m.jcarcd.cn/article/pdccgg.html
主站蜘蛛池模板: 韩国乱伦天堂网 | 国产高清| 日韩一二区| 国产日本精品视频 | 蜜桃视频 | 欧美日韩国产亚洲一 | 国产制服亚洲 | 99精品热 | 九一福利在线 | 三级欧美日本国产 | 日本韩国三级 | 国产另类日韩制 | 无码超乳爆乳中文字幕在线看伦 | 精品国精品国 | 成人激情视 | 日韩欧美亚洲一区 | 97国产 | 欧美午夜在线观看 | 成人论坛网 | 国产美女精品在线 | 国产高清午夜自 | 欧美一级在线全免费 | 天美麻花 | 日韩中文有码高清 | 国产盗摄不卡 | 国产精品视频免费的 | 国产后入在线观 | 国产激情精品一 | 国产高清在线视频色 | 精品免费在线视频 | 青草青在线 | 中文字幕精 | 国语我和子的乱视频 | 日韩v在线观看亚洲 | 日本淫秽视频在线 | 另类图片欧美小 | 日本在线观 | 国产95在 | 国产亚洲精在线看 | 韩国免费一级a一片 | 成人淫色午夜福利 |