• Spring Cloud Stream实践


    概述

    不同中间件,有各自的使用方法,代码也不一样。

    可以使用Spring Cloud Stream解耦,切换中间件时,不需要修改代码。实现方式为使用绑定层,绑定层对生产者和消费者提供统一的编码方式,需要连接不同的中间件时,绑定层使用不同的绑定器即可,也就是把切换中间件需要做相应的修改工作交给绑定层来做。

    本文的操作是在 微服务调用链路追踪 的基础上进行。

    环境说明

    jdk1.8

    maven3.6.3

    mysql8

    spring cloud2021.0.8

    spring boot2.7.12

    idea2022

    rabbitmq3.12.4

    步骤

    消息生产者

    创建子模块stream_producer

    添加依赖

    1. <dependencies>
    2. <dependency>
    3. <groupId>org.springframework.cloudgroupId>
    4. <artifactId>spring-cloud-starter-stream-rabbitartifactId>
    5. dependency>
    6. dependencies>

    刷新依赖

    配置application.yml

    1. server:
    2. port: 7001
    3. spring:
    4. application:
    5. name: stream_producer
    6. rabbitmq:
    7. addresses: 127.0.0.1
    8. username: guest
    9. password: guest
    10. cloud:
    11. stream:
    12. bindings:
    13. output:
    14. destination: my-default #指定消息发送的目的地,值为rabbit的exchange的名称
    15. binders:
    16. defaultRabbit:
    17. type: rabbit #配置默认的绑定器为rabbit

    查看Source.class源码

    编写生产者代码,发送一条消息("hello world")到rabbitmq的my-default exchange中

    1. package org.example.stream;
    2. import org.springframework.beans.factory.annotation.Autowired;
    3. import org.springframework.boot.CommandLineRunner;
    4. import org.springframework.boot.SpringApplication;
    5. import org.springframework.boot.autoconfigure.SpringBootApplication;
    6. import org.springframework.cloud.stream.annotation.EnableBinding;
    7. import org.springframework.cloud.stream.messaging.Source;
    8. import org.springframework.messaging.MessageChannel;
    9. import org.springframework.messaging.support.MessageBuilder;
    10. @EnableBinding(Source.class)
    11. @SpringBootApplication
    12. public class StreamProductApplication implements CommandLineRunner {
    13. @Autowired
    14. private MessageChannel output;
    15. @Override
    16. public void run(String... args) throws Exception {
    17. //发送消息
    18. // messageBuilder 工具类,创建消息
    19. output.send(MessageBuilder.withPayload("hello world").build());
    20. }
    21. public static void main(String[] args) {
    22. SpringApplication.run(StreamProductApplication.class, args);
    23. }
    24. }

    查看rabbitmq web UI

    http://localhost:15672/

    看到Exchanges中还没有my-default

    运行StreamProductApplication

    刷新rabbitmq Web UI,看到了my-dafault的exchange

    消息消费者

    创建子模块stream_consumer

    添加依赖

    1. <dependencies>
    2. <dependency>
    3. <groupId>org.springframework.cloudgroupId>
    4. <artifactId>spring-cloud-starter-stream-rabbitartifactId>
    5. dependency>
    6. dependencies>

    配置application.yml

    1. server:
    2. port: 7002
    3. spring:
    4. application:
    5. name: stream_consumer
    6. rabbitmq:
    7. addresses: 127.0.0.1
    8. username: guest
    9. password: guest
    10. cloud:
    11. stream:
    12. bindings:
    13. input: #内置获取消息的通道,从destination配置值的exchange中获取信息
    14. destination: my-default #指定消息发送的目的地,值为rabbit的exchange的名称
    15. binders:
    16. defaultRabbit:
    17. type: rabbit #配置默认的绑定器为rabbit

    查看内置通道名称为input

    编写消息消费者启动类,在启动类监听接收消息

    1. package org.example.stream;
    2. import org.springframework.boot.SpringApplication;
    3. import org.springframework.boot.autoconfigure.SpringBootApplication;
    4. import org.springframework.cloud.stream.annotation.EnableBinding;
    5. import org.springframework.cloud.stream.annotation.StreamListener;
    6. import org.springframework.cloud.stream.messaging.Sink;
    7. import org.springframework.messaging.Message;
    8. @SpringBootApplication
    9. @EnableBinding(Sink.class)
    10. public class StreamConsumerApplication {
    11. @StreamListener(Sink.INPUT)
    12. public void input(Message message){
    13. System.out.println("监听收到:" + message.getPayload());
    14. }
    15. public static void main(String[] args) {
    16. SpringApplication.run(StreamConsumerApplication.class, args);
    17. }
    18. }

    运行stream_consumer消费者服务,监听消息

    运行stream_producer生产者服务,发送消息

    查看消费者服务控制台日志,接收到了消息

    优化代码

    之前把生产和消费的消息都写在启动类中了,代码耦合高。

    优化思路是把不同功能的代码分开放。

    消息生产者

    stream_producer 代码结构如下

    1. package org.example.stream.producer;
    2. import org.springframework.beans.factory.annotation.Autowired;
    3. import org.springframework.cloud.stream.annotation.EnableBinding;
    4. import org.springframework.cloud.stream.messaging.Source;
    5. import org.springframework.messaging.MessageChannel;
    6. import org.springframework.messaging.support.MessageBuilder;
    7. import org.springframework.stereotype.Component;
    8. /**
    9. * 向中间件发送数据
    10. */
    11. @Component
    12. @EnableBinding(Source.class)
    13. public class MessageSender {
    14. @Autowired
    15. private MessageChannel output;//通道
    16. //发送消息
    17. public void send(Object obj){
    18. output.send(MessageBuilder.withPayload(obj).build());
    19. }
    20. }

    修改启动类

    1. package org.example.stream;
    2. import org.springframework.beans.factory.annotation.Autowired;
    3. import org.springframework.boot.CommandLineRunner;
    4. import org.springframework.boot.SpringApplication;
    5. import org.springframework.boot.autoconfigure.SpringBootApplication;
    6. import org.springframework.cloud.stream.annotation.EnableBinding;
    7. import org.springframework.cloud.stream.messaging.Source;
    8. import org.springframework.messaging.MessageChannel;
    9. import org.springframework.messaging.support.MessageBuilder;
    10. @SpringBootApplication
    11. public class StreamProductApplication {
    12. public static void main(String[] args) {
    13. SpringApplication.run(StreamProductApplication.class, args);
    14. }
    15. }

    pom.xml添加junit依赖

    1. <dependency>
    2. <groupId>junitgroupId>
    3. <artifactId>junitartifactId>
    4. <scope>testscope>
    5. dependency>

    刷新依赖

    编写测试类

    在stream_producerm模块的src/test目录下,新建org.example.stream包,再建出ProducerTest类,代码如下

    1. package org.example.stream;
    2. import org.example.stream.producer.MessageSender;
    3. import org.junit.Test;
    4. import org.junit.runner.RunWith;
    5. import org.springframework.beans.factory.annotation.Autowired;
    6. import org.springframework.boot.test.context.SpringBootTest;
    7. import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
    8. @RunWith(SpringJUnit4ClassRunner.class)
    9. @SpringBootTest
    10. public class ProducerTest {
    11. @Autowired
    12. private MessageSender messageSender;//注入发送消息工具类
    13. @Test
    14. public void testSend(){
    15. messageSender.send("hello world");
    16. }
    17. }

     

    消息消费者

    stream_consumer代码结构如下

    添加MessageListener类获取消息

    1. package org.example.stream.consumer;
    2. import org.springframework.cloud.stream.annotation.EnableBinding;
    3. import org.springframework.cloud.stream.annotation.StreamListener;
    4. import org.springframework.cloud.stream.messaging.Sink;
    5. import org.springframework.stereotype.Component;
    6. @Component
    7. @EnableBinding(Sink.class)
    8. public class MessageListener {
    9. // 监听binding中的信息
    10. @StreamListener(Sink.INPUT)
    11. public void input(String message){
    12. System.out.println("获取信息:" + message);
    13. }
    14. }

    修改启动类

    1. package org.example.stream;
    2. import org.springframework.boot.SpringApplication;
    3. import org.springframework.boot.autoconfigure.SpringBootApplication;
    4. import org.springframework.cloud.stream.annotation.EnableBinding;
    5. import org.springframework.cloud.stream.annotation.StreamListener;
    6. import org.springframework.cloud.stream.messaging.Sink;
    7. import org.springframework.messaging.Message;
    8. @SpringBootApplication
    9. public class StreamConsumerApplication {
    10. public static void main(String[] args) {
    11. SpringApplication.run(StreamConsumerApplication.class, args);
    12. }
    13. }

    启动consumer接收消息

    执行producer单元测试类ProducerTest的testSend()方法,发送消息

    查看consumer控制台输出,接收到信息了

    代码解耦后,同样能成功生产消息和消费消息。

    自定义消息通道

    此前使用默认的消息通道outputinput。

    也可以自己定义消息通道,例如:myoutputmyinput

    消息生产者

    org.example.stream包下新建channel包,在channel包下新建MyProcessor接口类

    1. package org.example.stream.channel;
    2. import org.springframework.cloud.stream.annotation.Input;
    3. import org.springframework.cloud.stream.annotation.Output;
    4. import org.springframework.messaging.MessageChannel;
    5. import org.springframework.messaging.SubscribableChannel;
    6. /**
    7. * 自定义的消息通道
    8. */
    9. public interface MyProcessor {
    10. /**
    11. * 消息生产这的配置
    12. */
    13. String MYOUTPUT = "myoutput";
    14. @Output("myoutput")
    15. MessageChannel myoutput();
    16. /**
    17. * 消息消费者的配置
    18. */
    19. String MYINPUT = "myinput";
    20. @Input("myinput")
    21. SubscribableChannel myinput();
    22. }

    修改MessageSender

    1. package org.example.stream.producer;
    2. import org.example.stream.channel.MyProcessor;
    3. import org.springframework.beans.factory.annotation.Autowired;
    4. import org.springframework.cloud.stream.annotation.EnableBinding;
    5. import org.springframework.cloud.stream.messaging.Source;
    6. import org.springframework.messaging.MessageChannel;
    7. import org.springframework.messaging.support.MessageBuilder;
    8. import org.springframework.stereotype.Component;
    9. /**
    10. * 向中间件发送数据
    11. */
    12. @Component
    13. @EnableBinding(MyProcessor.class)
    14. public class MessageSender {
    15. @Autowired
    16. private MessageChannel myoutput;//通道
    17. //发送消息
    18. public void send(Object obj){
    19. myoutput.send(MessageBuilder.withPayload(obj).build());
    20. }
    21. }

    修改application.yml

    1. cloud:
    2. stream:
    3. bindings:
    4. output:
    5. destination: my-default #指定消息发送的目的地
    6. myoutput:
    7. destination: custom-output

    消息消费者

    在stream_consumer服务的org.example.stream包下新建channel包,在channel包下新建MyProcessor接口类

    1. package org.example.stream.channel;
    2. import org.springframework.cloud.stream.annotation.Input;
    3. import org.springframework.cloud.stream.annotation.Output;
    4. import org.springframework.messaging.MessageChannel;
    5. import org.springframework.messaging.SubscribableChannel;
    6. /**
    7. * 自定义的消息通道
    8. */
    9. public interface MyProcessor {
    10. /**
    11. * 消息生产者的配置
    12. */
    13. String MYOUTPUT = "myoutput";
    14. @Output("myoutput")
    15. MessageChannel myoutput();
    16. /**
    17. * 消息消费者的配置
    18. */
    19. String MYINPUT = "myinput";
    20. @Input("myinput")
    21. SubscribableChannel myinput();
    22. }

    修改MessageListener

    1. package org.example.stream.stream;
    2. import org.example.stream.channel.MyProcessor;
    3. import org.springframework.cloud.stream.annotation.EnableBinding;
    4. import org.springframework.cloud.stream.annotation.StreamListener;
    5. import org.springframework.cloud.stream.messaging.Sink;
    6. import org.springframework.stereotype.Component;
    7. @Component
    8. @EnableBinding(MyProcessor.class)
    9. public class MessageListener {
    10. // 监听binding中的信息
    11. @StreamListener(MyProcessor.MYINPUT)
    12. public void input(String message){
    13. System.out.println("获取信息:" + message);
    14. }
    15. }

    修改application.yml配置

    1. cloud:
    2. stream:
    3. bindings:
    4. input: #内置获取消息的通道,从destination配置值的exchange中获取信息
    5. destination: my-default #指定消息发送的目的地
    6. myinput:
    7. destination: custom-output

    测试

    启动stream_consumer

    运行单元测试的testSend()方法生产消息

    查看stream_consumer控制台,能看到生产的消息,如下

    获取信息:hello world

    消息分组

    采用复制配置方式运行两个消费者

    启动第一个消费者(端口为7002)

    修改端口为7003,copy configuration,再启动另一个消费者

    执行生产者单元测试生产消息,看到两个消费者都接收到了信息

    说明:如果有两个消费者,生产一条消息后,两个消费者均能收到信息。

    但当我们发送一条消息只需要其中一个消费者消费消息时,这时候就需要用到消息分组,发送一条消息消费者组内只有一个消费者消费到。

    我们只需要在服务消费者端设置spring.cloud.stream.bindings.input.group 属性即可

    重启两个消费者

    修改端口号为7002,重新启动第一个消费者

    修改端口号为7003,重新启动第二个消费者

    生产者生产一条消息

    查看消费者接收消息情况,只有一个消费者接收到信息。

    消息分区

    消息分区就是实现特定消息只往特定机器发送。

    修改生产者配置

    1. cloud:
    2. stream:
    3. bindings:
    4. output:
    5. destination: my-default #指定消息发送的目的地,值为rabbit的exchange的名称
    6. myoutput:
    7. destination: custom-output
    8. producer:
    9. partition-key-expression: payload #分区关键字 可以是对象中的id,或对象
    10. partition-count: 2 #分区数量

    修改消费者1的application.yml配置

    1. server:
    2. port: 7002
    3. spring:
    4. application:
    5. name: stream_consumer
    6. rabbitmq:
    7. addresses: 127.0.0.1
    8. username: guest
    9. password: guest
    10. cloud:
    11. stream:
    12. bindings:
    13. input: #内置获取消息的通道,从destination配置值的exchange中获取信息
    14. destination: my-default #指定消息发送的目的地,值为rabbit的exchange的名称
    15. myinput:
    16. destination: custom-output
    17. group: group1 #消息分组,有多个消费者时,只有一个消费者接收到信息
    18. consumer:
    19. partitioned: true #开启分区支持
    20. binders:
    21. defaultRabbit:
    22. type: rabbit #配置默认的绑定器为rabbit
    23. instance-count: 2 #消费者总数
    24. instance-index: 0 #当前消费者的索引

    启动消费者1

    修改消费者2的配置

    1. server:
    2. port: 7003
    3. spring:
    4. application:
    5. name: stream_consumer
    6. rabbitmq:
    7. addresses: 127.0.0.1
    8. username: guest
    9. password: guest
    10. cloud:
    11. stream:
    12. bindings:
    13. input: #内置获取消息的通道,从destination配置值的exchange中获取信息
    14. destination: my-default #指定消息发送的目的地,值为rabbit的exchange的名称
    15. myinput:
    16. destination: custom-output
    17. group: group1 #消息分组,有多个消费者时,只有一个消费者接收到信息
    18. consumer:
    19. partitioned: true #开启分区支持
    20. binders:
    21. defaultRabbit:
    22. type: rabbit #配置默认的绑定器为rabbit
    23. instance-count: 2 #消费者总数
    24. instance-index: 1 #当前消费者的索引

    修改端口号为7003,当前消费者的索引instance-index的值修改为1

    启动消费者2

    生产者发送消息,看到只有Application(2)接收到消息

    再用生产者发送一次消息,也是Application(2)接收到消息

    说明实现了消息分区

    也可以更改发送的数据,看是否能发送到不同消费者

    修改生产者,发送数据由hello world变为hello world1,同时发送5次

    1. public void testSend(){
    2. for (int i = 0; i < 5; i++) {
    3. messageSender.send("hello world1");
    4. }
    5. }

    看到hello world1全部被Application消费

    所以消息分区是根据发送的消息不同,发送到不同消费者中。

    完成!enjoy it!

  • 相关阅读:
    AQS之Condition分析 (六)
    分布式存储技术解读系列之三:Swift | 架构进阶
    jsp基站管理系统servlet开发sqlserver数据库MVC结构java编程计算机网页项目
    Bug排查思路
    Nginx内置变量详解
    iCloud照片无法上传或同步怎么办?
    伴随对象的初始化
    Python哪个版本最稳定好用2023.10.19
    派生属性-架构案例2020(三十七)
    springboot毕设项目大学生请假系统 fq91k(java+VUE+Mybatis+Maven+Mysql)
  • 原文地址:https://blog.csdn.net/qq_42881421/article/details/134431700