为网站制定一个推广计划,网站LOGO透明底色PNG格式怎么做的,不需要充值的传奇手游,企业所得税怎么合理节税Stream消息驱动
gitee:springcloud_study: springcloud#xff1a;服务集群、注册中心、配置中心#xff08;热更新#xff09;、服务网关#xff08;校验、路由、负载均衡#xff09;、分布式缓存、分布式搜索、消息队列#xff08;异步通信#xff09;、数据库集群、…Stream消息驱动
gitee:springcloud_study: springcloud服务集群、注册中心、配置中心热更新、服务网关校验、路由、负载均衡、分布式缓存、分布式搜索、消息队列异步通信、数据库集群、分布式日志、系统监控链路追踪。
1. 消息驱动概述
作用屏蔽底层消息中间件的差异,降低切换成本统—消息的编程模型。底层不管是什么中间件如kafka、rabbitmqStream可以解决不同中间件的通信。 官网Spring Cloud Stream
官方定义Spring Cloud Stream是一个构建消息驱动微服务的框架。
应用程序通过 inputs 或者 outputsj来与Spring Cloud Stream中binder对象交互。通过我们配置来binding(绑定)而Spring Cloud Stream的 binder对象负责与消息中间件交互。所以我们只需要搞清楚如何与Spring Cloud Stream交互就可以方便使用消息驱动的方式。
通过使用Spring Integration来连接消息代理中间件以实现消息事件驱动。
Spring Cloud Stream为一些供应商的消息中间件产品提供了个性化的自动化配置实现引用了发布-订阅、消费组、分区的三个核心概念。
但是Stream只支持kafka、rabbitmq。 设计思想 标准的MQ
1.生产者/消费者之间靠消息媒介传递信息内容Message
2.消息必须走特定的通道消息通道MessageChannel
3.消息通道里的消息如何被消费呢谁负责收发处理消息通道MessageChannel的子接口SubscribableChannel由MessageHandler消息处理器所订阅
Cloud Stream
Stream利用Binder来绑定中间件的输入流和输出流。如果系统使用到了两个中间件kafka、rabbitmq这些中间件的差异性导致我们实际项目开发给我们造成了一定的困扰我们如果用了两个消息队列的其中一种后面的业务需求我想往另外一种消息队列进行迁移这时候无疑就是一个灾难性的人—大堆东西都要重新推倒重新做因为它跟我们的系统耦合了这时候springcloud Stream给我们提供了一种解耦合的方式。
Stream中的消息通信方式遵循了发布-订阅模式
Topic在Rabbitmq中是Exchange、在kafka中是Topic。
Spring Cloud Stream标准流程套路 Middleware中间件目前只支持RabbitMQ和Kafka
Binder是应用与消息中间件之间的封装目前实行了Kafka和RabbitMQ的Binder通过Binder可以很方便的连接中间件可以动态的改变消息类型(对应于Kafka的topic,RabbitMQ的exchange)这些都可以通过配置文件来实现。
Input注解标识输入通道通过该输入通道接收到的消息进入应用程序
Output注解标识输出通道发布的消息将通过该通道离开应用程序
StreamListener监听队列。用于消费者的队列的消息接收
EnableBinding指信道channel和exchange绑定在一起 Binder很方便的连接中间件屏蔽差异 Channel通道是队列Queue的一种抽象在消息通讯系统中就是实现存储和转发的媒介通过Channel对队列进行配置。 Source和Sink简单的可理解为参照对象是Spring Cloud Stream自身从Stream发布消息就是输出接受消息就是输入。
2. 消息驱动之生产者
创建cloud-stream-rabbitmq-provider8801作为生产者进行发消息模块 pom文件
dependencygroupIdorg.springframework.cloud/groupIdartifactIdspring-cloud-starter-stream-rabbit/artifactId
/dependency application.yaml
server:port: 8801
spring:application:name: cloud-stream-providercloud:stream:binders: #在此处配置要绑定的rabbitmq的服务信息;defaultRabbit:type: rabbitenvironment:spring:rabbitmq:host: 192.168.25.153port: 5672username: adminpassword: aaaaaabindings: #服务的整合处理output: #destination: studyExchange #表示要使用的Exchange名称定义content-type: application/json #设置消息类型本次为json,文本则设置“text/plain”binder: defaultRabbit #设置要绑定的消.息服务的具体设置
eureka:client:service-url:defaultZone: http://eureka7001.com:7001/eureka,http://eureka7002.com:7002/eurekaregister-with-eureka: truefetch-registry: trueinstance:lease-renewal-interval-in-seconds: 2lease-expiration-duration-in-seconds: 5instance-id: send-8801.comprefer-ip-address: true 主启动类
SpringBootApplication
EnableEurekaClient
public class StreamMQMain8801 {public static void main(String[] args) {SpringApplication.run(StreamMQMain8801.class,args);}
} service层
public interface IMessageProvider {String send();
}
EnableBinding(Source.class) //定义消息的推送管道
public class MessageProviderImpl implements IMessageProvider {
Resourceprivate MessageChannel output; //消息发送管道
Overridepublic String send() {String serial UUID.randomUUID().toString();output.send(MessageBuilder.withPayload(serial).build());System.out.println(********serialserial);return null;}
} controller层
RestController
public class SendMessageController {
Resourceprivate IMessageProvider messageProvider;
GetMapping(value /sendMessage)public String sendMessage(){return messageProvider.send();}
}
测试 3. 消息驱动之消费者
创建cloud-stream-rabbitmq-consumer8802作为消息接收模块 pom文件
dependencygroupIdorg.springframework.cloud/groupIdartifactIdspring-cloud-starter-stream-rabbit/artifactId
/dependency application.yml
server:port: 8802
spring:application:name: cloud-stream-consumercloud:stream:binders: #在此处配置要绑定的rabbitmq的服务信息;defaultRabbit:type: rabbitenvironment:spring:rabbitmq:host: 192.168.25.153port: 5672username: adminpassword: aaaaaabindings: #服务的整合处理input: #destination: studyExchange #表示要使用的Exchange名称定义content-type: application/json #设置消息类型本次为json,文本则设置“text/plain”binder: defaultRabbit #设置要绑定的消.息服务的具体设置
eureka:client:service-url:defaultZone: http://eureka7001.com:7001/eureka,http://eureka7002.com:7002/eurekaregister-with-eureka: truefetch-registry: trueinstance:lease-renewal-interval-in-seconds: 2lease-expiration-duration-in-seconds: 5instance-id: receive-8802.comprefer-ip-address: true controller层
Component
EnableBinding(Sink.class)
public class ReceiveMessageController {
Value(${server.port})private String serverPort;
StreamListener(Sink.INPUT)public void input(MessageString message){System.out.println(消费者1号------接收到的消息message.getPayload()\t portserverPort);}
} 主启动类
SpringBootApplication
EnableEurekaClient
public class ConsumerMQMain8802 {public static void main(String[] args) {SpringApplication.run(ConsumerMQMain8802.class,args);}
}
测试
启动loccalhost:8801/sendMessage就可以了消费者就是一个监听器有message就消费。 4. 分组消费与持久化
根据cloud-stream-rabbitmq-consumer8802创建8803项目运行暴露问题
消息重复消费和消息持久化问题需要进行分组操作。注意在Stream中处于同一个group中的多个消费者是竞争关系就能够保证消息只会被其中一个应用消费一次。不同组是可以全面消费的(重复消费)。
解决重复消费方法加入同一个组下图是不同分组的情况
cloud-stream-rabbitmq-consumer8802和8803设置不同分组yicaiA/B
server:port: 8803
spring:application:name: cloud-stream-consumercloud:stream:binders: #在此处配置要绑定的rabbitmq的服务信息;defaultRabbit:type: rabbitenvironment:spring:rabbitmq:host: 192.168.25.153port: 5672username: adminpassword: aaaaaabindings: #服务的整合处理input: #destination: studyExchange #表示要使用的Exchange名称定义content-type: application/json #设置消息类型本次为json,文本则设置“text/plain”binder: defaultRabbit #设置要绑定的消.息服务的具体设置group: yicaiB
server:port: 8802
spring:application:name: cloud-stream-consumercloud:stream:binders: #在此处配置要绑定的rabbitmq的服务信息;defaultRabbit:type: rabbitenvironment:spring:rabbitmq:host: 192.168.25.153port: 5672username: adminpassword: aaaaaabindings: #服务的整合处理input: #destination: studyExchange #表示要使用的Exchange名称定义content-type: application/json #设置消息类型本次为json,文本则设置“text/plain”binder: defaultRabbit #设置要绑定的消.息服务的具体设置group: yicaiA
cloud-stream-rabbitmq-consumer8802和8803设置同一个组yicaiA
server:port: 8802
spring:application:name: cloud-stream-consumercloud:stream:binders: #在此处配置要绑定的rabbitmq的服务信息;defaultRabbit:type: rabbitenvironment:spring:rabbitmq:host: 192.168.25.153port: 5672username: adminpassword: aaaaaabindings: #服务的整合处理input: #destination: studyExchange #表示要使用的Exchange名称定义content-type: application/json #设置消息类型本次为json,文本则设置“text/plain”binder: defaultRabbit #设置要绑定的消.息服务的具体设置group: yicaiA
测试
持久化 加上group就算实现类持久化。所谓的持久化就是如果没有分组一个服务发送消息其他服务由于没有分组如果其他哪些服务断开又继续重启这样就会导致以前那些消息丢失。