热门标签 | HotTags
当前位置:  开发笔记 > 编程语言 > 正文

开发笔记:springcloudstream

篇首语:本文由编程笔记#小编为大家整理,主要介绍了springcloudstream相关的知识,希望对你有一定的参考价值。创建springbo

篇首语:本文由编程笔记#小编为大家整理,主要介绍了spring cloud stream相关的知识,希望对你有一定的参考价值。



创建spring boot工程,添加pom依赖




org.springframework.cloud
spring-cloud-starter-stream-rabbit


View Code

添加消息接收SinkReceiver



import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.annotation.StreamListener;
import org.springframework.cloud.stream.messaging.Sink;
@EnableBinding(Sink.
class)
public class SinkReceiver {
private static Logger logger= LoggerFactory.getLogger(SinkReceiver.class);
@StreamListener(Sink.INPUT)
public void receive(Object payload){
logger.info(
"Received: "+payload);
}
}


View Code

配置



spring.application.name=stream-hello
spring.rabbitmq.host
=10.202.203.29
spring.rabbitmq.port
=5672
spring.rabbitmq.username
=springcloud
spring.rabbitmq.password
=123456


View Code

运行程序,打开rabbitmq监控界面,可以看到

推送消息

在控制台查看结果

 

创建一个消息发送类SinkSender



import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.messaging.Sink;
import org.springframework.context.annotation.Bean;
import org.springframework.integration.annotation.InboundChannelAdapter;
import org.springframework.integration.annotation.Poller;
import org.springframework.integration.core.MessageSource;
import org.springframework.messaging.support.GenericMessage;
import java.util.Date;
@EnableBinding(value
= {Sink.class})
public class SinkSender {
private static Logger logger= LoggerFactory.getLogger(SinkSender.class);
@Bean
@InboundChannelAdapter(value
= Sink.INPUT,poller = @Poller(fixedDelay = "2000"))
public MessageSource timerMessageSource(){
return ()-> new GenericMessage<>(new Date());
}
}


View Code

启动工程,可以在控制台看到每隔2秒收到信息

 

在SinkSender中添加日期转换



@Transformer(inputChannel = Sink.INPUT,outputChannel = Sink.INPUT)
public Object transform(Date message){
return new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").format(message);
}


View Code

控制台查看消息

 

添加一个User类



public class User {
private Integer id;
private String name;
private Integer age;
public Integer getId() {
return id;
}
public void setId(Integer id) {
this.id = id;
}
public String getName() {
return name;
}
public void setName(String name) {
this.name = name;
}
public Integer getAge() {
return age;
}
public void setAge(Integer age) {
this.age = age;
}
}


View Code

修改SinkSender



@Bean
@InboundChannelAdapter(value
= Sink.INPUT,poller = @Poller(fixedDelay = "2000"))
public MessageSource timerMessageSource(){
return ()->new GenericMessage<>("{\\"id\\":1,\\"name\\":\\"tom\\",\\"age\\":20}");
}


View Code

修改SinkReceiver



@ServiceActivator(inputChannel = Sink.INPUT)
public void receive(User user){
logger.info(
"Received: "+user);
}
@Transformer(inputChannel
= Sink.INPUT,outputChannel = Sink.INPUT)
public User transform(String message) throws Exception {
ObjectMapper objectMapper
=new ObjectMapper();
User user
=objectMapper.readValue(message,User.class);
return user;
}


View Code

这里使用@ServiceActivator必须指定@Transformer来处理自定义对象

改成就无需自定义@Transformer



@StreamListener(Sink.INPUT)
public void receive(User user){
logger.info(
"Received: "+user);
}


View Code

 


消息反馈

按上面项目再新建两个项目:App1和App2

App1



import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.annotation.StreamListener;
import org.springframework.cloud.stream.messaging.Processor;
import org.springframework.messaging.handler.annotation.SendTo;
@EnableBinding(value
= {Processor.class})
public class App1 {
private static Logger logger= LoggerFactory.getLogger(App1.class);
@StreamListener(Processor.INPUT)
@SendTo(Processor.OUTPUT)
public Object receiveFromInput(Object payload){
logger.info(
"Received: "+payload);
return "From Input Channel Return - "+payload;
}
}


View Code

App2



import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.annotation.StreamListener;
import org.springframework.cloud.stream.messaging.Processor;
import org.springframework.context.annotation.Bean;
import org.springframework.integration.annotation.InboundChannelAdapter;
import org.springframework.integration.annotation.Poller;
import org.springframework.integration.core.MessageSource;
import org.springframework.messaging.support.GenericMessage;
import java.util.Date;
@EnableBinding(value
= {Processor.class})
public class App2 {
private static Logger logger= LoggerFactory.getLogger(App2.class);
@Bean
@InboundChannelAdapter(value
= Processor.OUTPUT,poller = @Poller(fixedDelay = "2000"))
public MessageSource timeMessageSource(){
return ()->new GenericMessage<>(new Date());
}
@StreamListener(Processor.INPUT)
public void receiveFromOutput(Object payload){
logger.info(
"Received: "+payload);
}
}


View Code

App2的配置做个变更



spring.rabbitmq.host=10.202.203.29
spring.rabbitmq.port
=5672
spring.rabbitmq.username
=springcloud
spring.rabbitmq.password
=123456
spring.cloud.stream.bindings.input.destination
=output
spring.cloud.stream.bindings.output.destination
=input
server.port
=8001


View Code

启动两个项目

 


消费组

启动多个消费端App1和一个生产端App2,可以看到App2发送的消息被多个App1接收并处理

通过指定group可以然消息只被相应的group接收

App1-1


spring.cloud.stream.bindings.input.group=Service-A

App1-2


spring.cloud.stream.bindings.input.group=Service-A

App2


spring.cloud.stream.bindings.input.group=Service-A

这样App2发送的消息将被两个App1轮询处理

如果此时添加一个App1-3


spring.cloud.stream.bindings.input.group=Service-B

 

 从rabbitmq管理界面查看

 

一个exchange绑定了两个queue,从exchange里推送一条消息,两个queue里都会有



推荐阅读
  • 深入解析SpringMVC核心组件:DispatcherServlet的工作原理
    本文详细探讨了SpringMVC的核心组件——DispatcherServlet的运作机制,旨在帮助有一定Java和Spring基础的开发人员理解HTTP请求是如何被映射到Controller并执行的。文章将解答以下问题:1. HTTP请求如何映射到Controller;2. Controller是如何被执行的。 ... [详细]
  • 本文详细介绍了如何在 Android 中使用值动画(ValueAnimator)来动态调整 ImageView 的高度,并探讨了相关的关键属性和方法,包括图片填充后的高度、原始图片高度、动画变化因子以及布局重置等。 ... [详细]
  • 深入解析Spring启动过程
    本文详细介绍了Spring框架的启动流程,帮助开发者理解其内部机制。通过具体示例和代码片段,解释了Bean定义、工厂类、读取器以及条件评估等关键概念,使读者能够更全面地掌握Spring的初始化过程。 ... [详细]
  • 探讨ChatGPT在法律和版权方面的潜在风险及影响,分析其作为内容创造工具的合法性和合规性。 ... [详细]
  • ListView简单使用
    先上效果:主要实现了Listview的绑定和点击事件。项目资源结构如下:先创建一个动物类,用来装载数据:Animal类如下:packagecom.example.simplelis ... [详细]
  • springMVC JRS303验证 ... [详细]
  • 烤鸭|本文_Spring之Bean的生命周期详解
    烤鸭|本文_Spring之Bean的生命周期详解 ... [详细]
  • 本文介绍了如何使用JavaScript的Fetch API与Express服务器进行交互,涵盖了GET、POST、PUT和DELETE请求的实现,并展示了如何处理JSON响应。 ... [详细]
  • Java项目分层架构设计与实践
    本文探讨了Java项目中应用分层的最佳实践,不仅介绍了常见的三层架构(Controller、Service、DAO),还深入分析了各层的职责划分及优化建议。通过合理的分层设计,可以提高代码的可维护性、扩展性和团队协作效率。 ... [详细]
  • 本文探讨了在Django项目中,如何在对象详情页面添加前后导航链接,以提升用户体验。文章详细描述了遇到的问题及解决方案。 ... [详细]
  • SpringMVC RestTemplate的几种请求调用(转)
    SpringMVCRestTemplate的几种请求调用(转),Go语言社区,Golang程序员人脉社 ... [详细]
  • 本文介绍如何在C#中将GridView控件的内容保存为图片文件。通过代码示例,详细说明了创建位图、绘制图形并保存图像的步骤。 ... [详细]
  • 探讨了在 Spring MVC 框架下,JSP 页面使用 标签时遇到的数据无法正确显示的问题,并提供了可能的原因和解决方案。 ... [详细]
  • 本文详细探讨了在微服务架构中,使用Feign进行远程调用时出现的请求头丢失问题,并提供了具体的解决方案。重点讨论了单线程和异步调用两种场景下的处理方法。 ... [详细]
  • Spring Cloud Config 使用 Vault 作为配置存储
    本文探讨了如何在Spring Cloud Config中集成HashiCorp Vault作为配置存储解决方案,基于Spring Cloud Hoxton.RELEASE及Spring Boot 2.2.1.RELEASE版本。文章还提供了详细的配置示例和实践建议。 ... [详细]
author-avatar
十一
这个家伙很懒,什么也没留下!
PHP1.CN | 中国最专业的PHP中文社区 | DevBox开发工具箱 | json解析格式化 |PHP资讯 | PHP教程 | 数据库技术 | 服务器技术 | 前端开发技术 | PHP框架 | 开发工具 | 在线工具
Copyright © 1998 - 2020 PHP1.CN. All Rights Reserved | 京公网安备 11010802041100号 | 京ICP备19059560号-4 | PHP1.CN 第一PHP社区 版权所有