Spring Cloud Stream之rocketmq
文档网址:参考:Spring Cloud Stream 是一个用来为微服务应用构建消息驱动能力的框架。它可以基于 Spring Boot 来创建独立的、可用于生产的 Spring 应用程序。Spring Cloud Stream 为一些供应商的消息中间件产品提供了个性化的自动化配置实现,并引入了发布-订阅、消费组、分区这三个核心概念。简单的说,Spring Cloud Stream本质上就是整合了
文档网址: 进阶指南 | Spring Cloud Alibaba
部署文档:部署方式 | RocketMQ
参考: SpringCloud之Stream消息驱动RocketMQ讲解_spring cloud stream rocketmq-CSDN博客
Spring Cloud Stream 是一个用来为微服务应用构建消息驱动能力的框架。它可以基于 Spring Boot 来创建独立的、可用于生产的 Spring 应用程序。Spring Cloud Stream 为一些供应商的消息中间件产品提供了个性化的自动化配置实现,并引入了发布-订阅、消费组、分区这三个核心概念。简单的说,Spring Cloud Stream本质上就是整合了Spring Boot和Spring Integration,实现了一套轻量级的消息驱动的微服务框架。

RocketMQ整合binder架构实现:

spring-cloud-statrer-stream-rocketmq 去除了对 RocketMQ-Spring 框架的依赖 。 Spring Cloud Stream Binder 核心类 RocketMQMessageChannelBinder 实现了 Spring Cloud Stream 规范,内部会构建 RocketMQInboundChannelAdapter 和 RocketMQProducerMessageHandler。
RocketMQProducerMessageHandler 会基于 Binding 配置通过 RocketMQProduceFactory 构造 RocketMQ Producer,其内部会把 spring-messaging 模块内 org.springframework.messaging.Message 消息类转换成 RocketMQ 的消息类 org.apache.rocketmq.common.message.Message,然后发送出去。
RocketMQInboundChannelAdapter 也会基于 Binding 配置通过 RocketMQConsumerFactory 构造 DefaultMQPushConsumer,其内部会启动 RocketMQ Consumer 接收消息。
目前 Binder 支持在 Header 中设置相关的 key 来进行 RocketMQ Message 消息的特性设置:
MessageBuilder builder = MessageBuilder.withPayload(msg)
.setHeader(RocketMQHeaders.TAGS, "binder")
.setHeader(RocketMQHeaders.KEYS, "my-key");
Message message = builder.build();
output().send(message);
或者使用 StreamBridge:
MessageBuilder builder = MessageBuilder.withPayload(msg)
.setHeader(RocketMQHeaders.TAGS, "binder")
.setHeader(RocketMQHeaders.KEYS, "my-key");
Message message = builder.build();
streamBridge.send("producer-out-0", message);
Springboot整合RocketMQ:
pom.xml依赖:
<!-- RocketMQ -->
<dependency>
<groupId>com.alibaba.cloud</groupId>
<artifactId>spring-cloud-starter-stream-rocketmq</artifactId>
</dependency>
application.yml配置:
spring:
cloud:
function:
definition: consumer1;consumer2
stream:
rocketmq:
binder:
name-server: localhost:9876
group: producer_group
bindings:
consumer1-in-0:
destination: topic1
binder: rocketmq
group: consumer_group1
content-type: application/json
consumer2-in-0:
destination: topic2
binder: rocketmq
group: consumer_group2
content-type: application/json
producer1-out-0:
destination: topic1
binder: rocketmq
发送消息:
@Autowired
private StreamBridge streamBridge;
public void sendMessage(TaskMsg taskMsg) {
Message<DefaultTaskEvent> message = MessageBuilder
.withPayload(taskMsg)
.setPriority(taskMsg.getTaskPriority())
.setHeader(MessageConst.PROPERTY_DELAY_TIME_LEVEL, taskMsg.getDelayLevel())
.build();
boolean send = streamBridge.send("producer1-out-0", message);
}
消费消息:
@Bean
public Consumer<TaskMsg> consumer1() {
return taskMsg-> {
System.out.println(taskMsg);
}
更多推荐



所有评论(0)