Spring Cloud Stream is a framework for building highly scalable event-driven microservices connected with shared messaging systems.html
The framework provides a flexible programming model built on already established and familiar Spring idioms and best practices, including support for persistent pub/sub semantics, consumer groups, and stateful partitions.java
Spring Cloud Stream是一个用于构建与共享消息传递系统链接的高度可扩展的事件驱动型微服务的框架。git
应用程序经过inputs或outputs来与Spring Cloud Stream中binder对象交互,binder对象负责与消息中间件交互。也就是说:Spring Cloud Stream可以屏蔽底层消息中间件【RabbitMQ,kafka等】的差别,下降切换成本,统一消息的编程模型。所以,若是咱们想要使用消息驱动,咱们不须要了解各类消息中间件,咱们只须要了解Spring Cloud Stream就行了。web
SpringCloud Stream经过Spring Integration来链接消息代理中间件以实现消息事件驱动。spring
SpringCloud Stream还为一些供应商的消息中间件产品提供了个性化的自动化配置实现,引用了发布-订阅,消费组,分区三个核心概念。编程
目前支持的消息中间件官网能够查询:像RabbitMQ,kafka都是支持的,本篇文章基于RabbitMQ消息中间件。json
标准消息队列的特色:api
如何统一底层差别?经过定义binder绑定器做为中间层,完美地实现了应用程序与消息中间件细节之间的隔离。经过向应用程序暴露统一的Channel通道,使得应用程序再也不须要考虑各类不一样的消息中间件的实现。app
Stream中的消息通讯方式遵循发布-订阅模式,使用Topic主题进行广播。框架
组成 | 说明 |
---|---|
Middleware | 中间件,目前只支持RabbitMQ和Kafka |
Binder | Binder是应用与消息中间件之间的封装,目前实行了Kafka和RabbitMQ的Binder, 经过Binder能够很方便链接中间件,能够动态地改变消息类型 |
@Input | 注解标识输入通道,经过该输出通道接收到地消息进入应用程序 |
@Output | 注解标识输出通道,发布的消息将经过该通道离开应用程序 |
@StreamListener | 监听队列,用于消费者的队列的消息接收 |
@EnableBinding | 指信道channel和exchange绑定在一块儿 |
cloud-stream-rabbitmq-provider8801
,做为生产者进行发消息模块。cloud-stream-rabbitmq-consumer8802
,做为消息接收模块。cloud-stream-rabbitmq-consumer8803
,做为消息接收模块。<dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-starter-stream-rabbit</artifactId> </dependency>
其实没有涉及到服务注册发现,但为了完整性,仍是将该服务注册进服务注册中心的配置加上。
server: port: 8801 spring: application: name: cloud-stream-provider cloud: stream: binders: # 在此处配置要绑定的rabbitmq的服务信息; defaultRabbit: # 表示定义的名称,用于于binding整合 type: rabbit # 消息组件类型 environment: # 设置rabbitmq的相关的环境配置 spring: rabbitmq: host: [hostname] port: 5672 username: guest password: guest bindings: # 服务的整合处理 output: # 这个名字是一个通道的名称 destination: studyExchange # 表示要使用的Exchange名称定义 content-type: application/json # 设置消息类型,本次为json,文本则设置“text/plain” binder: defaultRabbit # 设置要绑定的消息服务的具体设置 eureka: client: # 客户端进行Eureka注册的配置 service-url: defaultZone: http://localhost:7001/eureka instance: lease-renewal-interval-in-seconds: 2 # 设置心跳的时间间隔(默认是30秒) lease-expiration-duration-in-seconds: 5 # 若是如今超过了5秒的间隔(默认是90秒) instance-id: send-8801.com # 在信息列表时显示主机名称 prefer-ip-address: true # 访问的路径变为IP地址
@SpringBootApplication public class StreamMQMain8801 { public static void main(String[] args) { SpringApplication.run(StreamMQMain8801.class, args); } }
@EnableBinding(Source.class) //定义消息的推送管道 public class MessageProviderImpl implements IMessageProvider { /* * 消息发送管道 */ @Resource private MessageChannel output; @Override public String send() { String serial = UUID.randomUUID().toString(); output.send(MessageBuilder.withPayload(serial).build()); System.out.println("==> serial + " + serial); return null; } }
@RestController public class SendMessageController { @Resource private IMessageProvider messageProvider; @GetMapping(value = "/sendMessage") public String sendMessage() { return messageProvider.send(); } }
启动7001注册中心,启动8801生产者模块,接着登陆rabbitMQ的web图形化界面,咱们将会看到一个新建的exchange:studyExchange,这是咱们在yml中配置的。
接着访问接口,localhost:8801/sendMessage
,控制台输出serial不断,rabbitMQ中也成功发送消息,测试成功。
一样引入spring-cloud-starter-stream-rabbit
依赖。
<dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-starter-stream-rabbit</artifactId> </dependency>
server: port: 8802 spring: application: name: cloud-stream-consumer cloud: stream: binders: # 在此处配置要绑定的rabbitmq的服务信息; defaultRabbit: # 表示定义的名称,用于于binding整合 type: rabbit # 消息组件类型 environment: # 设置rabbitmq的相关的环境配置 spring: rabbitmq: host: [hostname] port: 5672 username: guest password: guest bindings: # 服务的整合处理 input: # 这个名字是一个通道的名称 destination: studyExchange # 表示要使用的Exchange名称定义 content-type: application/json # 设置消息类型,本次为对象json,若是是文本则设置“text/plain” binder: defaultRabbit # 设置要绑定的消息服务的具体设置 eureka: client: # 客户端进行Eureka注册的配置 service-url: defaultZone: http://localhost:7001/eureka instance: lease-renewal-interval-in-seconds: 2 # 设置心跳的时间间隔(默认是30秒) lease-expiration-duration-in-seconds: 5 # 若是如今超过了5秒的间隔(默认是90秒) instance-id: receive-8802.com # 在信息列表时显示主机名称 prefer-ip-address: true # 访问的路径变为IP地址
@SpringBootApplication public class StreamMQMain8802 { public static void main(String[] args) { SpringApplication.run(StreamMQMain8802.class, args); } }
@Component @EnableBinding(Sink.class) public class ReceiveMessageListenerController { @Value("${server.port}") private String serverPort; @StreamListener(Sink.INPUT) public void input(Message<String> message) { System.out.println("消费者1号,----->接受到的消息: " + message.getPayload() + "\t port: " + serverPort); } }
在刚刚两个服务启动以后,再启动刚刚建立的8802模块。
访问:localhost:8801/sendMessage
,8802模块控制台成功打印消息。
按照上面的步骤运行下来,貌似没什么问题,其实并非?
为了演示这个问题,咱们仿照8802模块,克隆一份8803模块,并依次启动7001注册中心,8801消息生产者,880二、8803消息消费者。
咱们只需访问:localhost:8801/sendMessage
就可以显示消息重复消费的问题,对此,咱们能够经过Stream中的分组消费来解决,同一个组的消费者存在竞争关系,只能有一个能够进行消费。
咱们经过rabbitMQ的管理页面就能看到,这两个消费者默认的组流水号是不一样的,解决的办法也很简单,指定他们的流水号相同便可。
咱们在8802和8803中的yml中配置:spring.cloud.stream.bindings.input.group=A
便可实现轮询消费。
消息持久化是很关键的步骤,若是不具有消息持久化的功能,假设某一消费者忽然宕机,生产者持续发送消息,消费者没法消费,会致使消息丢失。
加上group属性以后,就已经具有了消息持久化,演示也很简单,关闭消费者服务,生产者不断发送信息,重启消费者服务,发现启动以后,将宕机时的错过的消息消费。
本系列文章为《尚硅谷SpringCloud教程》的学习笔记【版本稍微有些不一样,后续遇到bug再作相关说明】,主要作一个长期的记录,为之后学习的同窗提供示例,代码同步更新到Gitee:https://gitee.com/tqbx/spring-cloud-learning,而且以标签的形式详细区分每一个步骤,这个系列文章也会同步更新。