整合spring cloud云架构 -消息驱动 Spring Cloud Stream

简介: Spring Cloud Stream

在使用spring cloud云架构的时候,我们不得不使用Spring cloud Stream,因为消息中间件的使用在项目中无处不在,我们公司后面做了娱乐方面的APP,在使用spring cloud做架构的时候,其中消息的异步通知,业务的异步处理都需要使用消息中间件机制。spring cloud的官方给出的集成建议(使用rabbit mq和kafka),我看了一下源码和配置,只要把rabbit mq集成,kafka只是换了一个pom配置jar包而已,闲话少说,我们就直接进入配置实施:

  1. 简介:

Spring cloud Stream 数据流操作开发包,封装了与Redis,Rabbit、Kafka等发送接收消息。

  1. 使用工具:

rabbit,具体的下载和安装细节我这里不做太多讲解,网上的实例太多了

  1. 创建commonservice-mq-producer消息的发送者项目,在pom里面配置stream-rabbit的依赖


<groupId>org.springframework.cloud</groupId>  
<artifactId>spring-cloud-starter-stream-rabbit</artifactId>  

  1. 在yml文件里面配置rabbit mq

Java代码 收藏代码
server:
port: 5666
spring:
application:

name: commonservice-mq-producer  

profiles:

active: dev  

cloud:

config:  
  discovery:   
    enabled: true  
    service-id: commonservice-config-server  

# rabbitmq和kafka都有相关配置的默认值,如果修改,可以再次进行配置

stream:  
  bindings:  
    mqScoreOutput:   
      destination: honghu_exchange  
      contentType: application/json  
        

rabbitmq:

 host: localhost  
 port: 5672  
 username: honghu  
 password: honghu</span>  

eureka:
client:

service-url:  
  defaultZone: http://honghu:123456@localhost:8761/eureka  

instance:

prefer-ip-address: true</span>  
  1. 定义接口ProducerService

Java代码 收藏代码
package com.honghu.cloud.producer;

import org.springframework.cloud.stream.annotation.Output;
import org.springframework.messaging.SubscribableChannel;

public interface ProducerService {

  
String SCORE_OUPUT = "mqScoreOutput";  
  
@Output(ProducerService.SCORE_OUPUT)  
SubscribableChannel sendMessage();  

}

  1. 定义绑定

Java代码 收藏代码
package com.honghu.cloud.producer;

import org.springframework.cloud.stream.annotation.EnableBinding;

@EnableBinding(ProducerService.class)
public class SendServerConfig {

}

  1. 定义发送消息业务ProducerController

Java代码 收藏代码
package com.honghu.cloud.controller;

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.messaging.Message;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestMethod;
import org.springframework.web.bind.annotation.RestController;

import com.honghu.cloud.common.code.ResponseCode;
import com.honghu.cloud.common.code.ResponseVO;
import com.honghu.cloud.entity.User;
import com.honghu.cloud.producer.ProducerService;

import net.sf.json.JSONObject;

@RestController
@RequestMapping(value = "producer")
public class ProducerController {

  
@Autowired  
private ProducerService producerService;  
  
  
/** 
 * 通过get方式发送</span>对象<span style="font-size: 16px;"> 
 * @param name 路径参数 
 * @return 成功|失败 
 */  
@RequestMapping(value = "/sendObj", method = RequestMethod.GET)  
public ResponseVO sendObj() {  
    User user = new User(1, "hello User");  
    <span style="color: #ff0000;">Message<User> msg = MessageBuilder.withPayload(user).build();</span>  
    boolean result = producerService.sendMessage().send(msg);  
    if(result){  
        return ResponseCode.buildEnumResponseVO(ResponseCode.RESPONSE_CODE_SUCCESS, false);  
    }  
    return ResponseCode.buildEnumResponseVO(ResponseCode.RESPONSE_CODE_FAILURE, false);  
}  
  
  
/** 
 * 通过get方式发送字符串消息 
 * @param name 路径参数 
 * @return 成功|失败 
 */  
@RequestMapping(value = "/send/{name}", method = RequestMethod.GET)  
public ResponseVO send(@PathVariable(value = "name", required = true) String name) {  
    Message msg = MessageBuilder.withPayload(name.getBytes()).build();  
    boolean result = producerService.sendMessage().send(msg);  
    if(result){  
        return ResponseCode.buildEnumResponseVO(ResponseCode.RESPONSE_CODE_SUCCESS, false);  
    }  
    return ResponseCode.buildEnumResponseVO(ResponseCode.RESPONSE_CODE_FAILURE, false);  
}  
  
/** 
 * 通过post方式发送</span>json对象<span style="font-size: 16px;"> 
 * @param name 路径参数 
 * @return 成功|失败 
 */  
@RequestMapping(value = "/sendJsonObj", method = RequestMethod.POST)  
public ResponseVO sendJsonObj(@RequestBody JSONObject jsonObj) {  
    Message<JSONObject> msg = MessageBuilder.withPayload(jsonObj).build();  
    boolean result = producerService.sendMessage().send(msg);  
    if(result){  
        return ResponseCode.buildEnumResponseVO(ResponseCode.RESPONSE_CODE_SUCCESS, false);  
    }  
    return ResponseCode.buildEnumResponseVO(ResponseCode.RESPONSE_CODE_FAILURE, false);  
}  

}

  1. 创建commonservice-mq-consumer1消息的消费者项目,在pom里面配置stream-rabbit的依赖

Java代码 收藏代码

<groupId>org.springframework.cloud</groupId>  
<artifactId>spring-cloud-starter-stream-rabbit</artifactId>  

  1. 在yml文件中配置:

Java代码 收藏代码
server:
port: 5111
spring:
application:

name: commonservice-mq-consumer1  

profiles:

active: dev  

cloud:

config:  
  discovery:   
    enabled: true  
    service-id: commonservice-config-server  
      
<span style="color: #ff0000;">stream:  
  bindings:  
    mqScoreInput:  
      group: honghu_queue  
      destination: honghu_exchange  
      contentType: application/json  
        

rabbitmq:

 host: localhost  
 port: 5672  
 username: honghu  
 password: honghu</span>  

eureka:
client:

service-url:  
  defaultZone: http://honghu:123456@localhost:8761/eureka  

instance:

prefer-ip-address: true  
  1. 定义接口ConsumerService

Java代码 收藏代码
package com.honghu.cloud.consumer;

import org.springframework.cloud.stream.annotation.Input;
import org.springframework.messaging.SubscribableChannel;

public interface ConsumerService {

  
<span style="color: #ff0000;">String SCORE_INPUT = "mqScoreInput";  

@Input(ConsumerService.SCORE_INPUT)  
SubscribableChannel sendMessage();</span>  

}

  1. 定义启动类和消息消费

Java代码 收藏代码
package com.honghu.cloud;

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.cloud.netflix.eureka.EnableEurekaClient;
import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.annotation.StreamListener;

import com.honghu.cloud.consumer.ConsumerService;
import com.honghu.cloud.entity.User;

@EnableEurekaClient
@SpringBootApplication
@EnableBinding(ConsumerService.class) //可以绑定多个接口
public class ConsumerApplication {

  
public static void main(String[] args) {  
    SpringApplication.run(ConsumerApplication.class, args);  
}  
  
<span style="color: #ff0000;">@StreamListener(ConsumerService.SCORE_INPUT)  
public void onMessage(Object obj) {  
    System.out.println("消费者1,接收到的消息:" + obj);  
}</span>  

}

  1. 分别启动commonservice-mq-producer、commonservice-mq-consumer1
  2. 通过postman来验证消息的发送和接收

可以看到接收到了消息,下一章我们介绍mq的集群方案。

到此,整个消息中心方案集成完毕(企业架构源码可以加求球:三五三六二四七二五九)

欢迎大家和我一起学习spring cloud构建微服务云架构,我这边会将近期研发的spring cloud微服务云架构的搭建过程和精髓记录下来,帮助更多有兴趣研发spring cloud框架的朋友,大家来一起探讨spring cloud架构的搭建过程及如何运用于企业项目。

目录
相关文章
|
2月前
|
运维 监控 负载均衡
动态服务管理平台:驱动微服务架构的高效引擎
动态服务管理平台:驱动微服务架构的高效引擎
31 0
|
1月前
|
负载均衡 Java 开发者
深入探索Spring Cloud与Spring Boot:构建微服务架构的实践经验
深入探索Spring Cloud与Spring Boot:构建微服务架构的实践经验
148 5
|
2月前
|
消息中间件 Java Kafka
Spring Boot 与 Apache Kafka 集成详解:构建高效消息驱动应用
Spring Boot 与 Apache Kafka 集成详解:构建高效消息驱动应用
61 1
|
4月前
|
人工智能 开发框架 Java
重磅发布!AI 驱动的 Java 开发框架:Spring AI Alibaba
随着生成式 AI 的快速发展,基于 AI 开发框架构建 AI 应用的诉求迅速增长,涌现出了包括 LangChain、LlamaIndex 等开发框架,但大部分框架只提供了 Python 语言的实现。但这些开发框架对于国内习惯了 Spring 开发范式的 Java 开发者而言,并非十分友好和丝滑。因此,我们基于 Spring AI 发布并快速演进 Spring AI Alibaba,通过提供一种方便的 API 抽象,帮助 Java 开发者简化 AI 应用的开发。同时,提供了完整的开源配套,包括可观测、网关、消息队列、配置中心等。
3298 15
|
3月前
|
消息中间件 监控 NoSQL
驱动系统架构
【10月更文挑战第29天】
38 2
|
3月前
|
存储 前端开发 API
DDD领域驱动设计实战-分层架构
DDD分层架构通过明确各层职责及交互规则,有效降低了层间依赖。其基本原则是每层仅与下方层耦合,分为严格和松散两种形式。架构演进包括传统四层架构与改良版四层架构,后者采用依赖反转设计原则优化基础设施层位置。各层职责分明:用户接口层处理显示与请求;应用层负责服务编排与组合;领域层实现业务逻辑;基础层提供技术基础服务。通过合理设计聚合与依赖关系,DDD支持微服务架构灵活演进,提升系统适应性和可维护性。
|
3月前
|
负载均衡 Java API
【Spring Cloud生态】Spring Cloud Gateway基本配置
【Spring Cloud生态】Spring Cloud Gateway基本配置
68 0
|
5月前
|
消息中间件 Java RocketMQ
微服务架构师的福音:深度解析Spring Cloud RocketMQ,打造高可靠消息驱动系统的不二之选!
【8月更文挑战第29天】Spring Cloud RocketMQ结合了Spring Cloud生态与RocketMQ消息中间件的优势,简化了RocketMQ在微服务中的集成,使开发者能更专注业务逻辑。通过配置依赖和连接信息,可轻松搭建消息生产和消费流程,支持消息过滤、转换及分布式事务等功能,确保微服务间解耦的同时,提升了系统的稳定性和效率。掌握其应用,有助于构建复杂分布式系统。
78 0
|
2天前
|
Java 测试技术 应用服务中间件
Spring Boot 如何测试打包部署
本文介绍了 Spring Boot 项目的开发、调试、打包及投产上线的全流程。主要内容包括: 1. **单元测试**:通过添加 `spring-boot-starter-test` 包,使用 `@RunWith(SpringRunner.class)` 和 `@SpringBootTest` 注解进行测试类开发。 2. **集成测试**:支持热部署,通过添加 `spring-boot-devtools` 实现代码修改后自动重启。 3. **投产上线**:提供两种部署方案,一是打包成 jar 包直接运行,二是打包成 war 包部署到 Tomcat 服务器。
24 10
下一篇
开通oss服务