热门标签 | 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里都会有



推荐阅读
  • Android工程师面试准备及设计模式使用场景
    本文介绍了Android工程师面试准备的经验,包括面试流程和重点准备内容。同时,还介绍了建造者模式的使用场景,以及在Android开发中的具体应用。 ... [详细]
  • 个人学习使用:谨慎参考1Client类importcom.thoughtworks.gauge.Step;importcom.thoughtworks.gauge.T ... [详细]
  • 本文分享了一个关于在C#中使用异步代码的问题,作者在控制台中运行时代码正常工作,但在Windows窗体中却无法正常工作。作者尝试搜索局域网上的主机,但在窗体中计数器没有减少。文章提供了相关的代码和解决思路。 ... [详细]
  • 前景:当UI一个查询条件为多项选择,或录入多个条件的时候,比如查询所有名称里面包含以下动态条件,需要模糊查询里面每一项时比如是这样一个数组条件:newstring[]{兴业银行, ... [详细]
  • 如何自行分析定位SAP BSP错误
    The“BSPtag”Imentionedintheblogtitlemeansforexamplethetagchtmlb:configCelleratorbelowwhichi ... [详细]
  • Nginx使用(server参数配置)
    本文介绍了Nginx的使用,重点讲解了server参数配置,包括端口号、主机名、根目录等内容。同时,还介绍了Nginx的反向代理功能。 ... [详细]
  • 向QTextEdit拖放文件的方法及实现步骤
    本文介绍了在使用QTextEdit时如何实现拖放文件的功能,包括相关的方法和实现步骤。通过重写dragEnterEvent和dropEvent函数,并结合QMimeData和QUrl等类,可以轻松实现向QTextEdit拖放文件的功能。详细的代码实现和说明可以参考本文提供的示例代码。 ... [详细]
  • android listview OnItemClickListener失效原因
    最近在做listview时发现OnItemClickListener失效的问题,经过查找发现是因为button的原因。不仅listitem中存在button会影响OnItemClickListener事件的失效,还会导致单击后listview每个item的背景改变,使得item中的所有有关焦点的事件都失效。本文给出了一个范例来说明这种情况,并提供了解决方法。 ... [详细]
  • Spring特性实现接口多类的动态调用详解
    本文详细介绍了如何使用Spring特性实现接口多类的动态调用。通过对Spring IoC容器的基础类BeanFactory和ApplicationContext的介绍,以及getBeansOfType方法的应用,解决了在实际工作中遇到的接口及多个实现类的问题。同时,文章还提到了SPI使用的不便之处,并介绍了借助ApplicationContext实现需求的方法。阅读本文,你将了解到Spring特性的实现原理和实际应用方式。 ... [详细]
  • 本文讨论了一个关于cuowu类的问题,作者在使用cuowu类时遇到了错误提示和使用AdjustmentListener的问题。文章提供了16个解决方案,并给出了两个可能导致错误的原因。 ... [详细]
  • 拥抱Android Design Support Library新变化(导航视图、悬浮ActionBar)
    转载请注明明桑AndroidAndroid5.0Loollipop作为Android最重要的版本之一,为我们带来了全新的界面风格和设计语言。看起来很受欢迎࿰ ... [详细]
  • Android开发实现的计时器功能示例
    本文分享了Android开发实现的计时器功能示例,包括效果图、布局和按钮的使用。通过使用Chronometer控件,可以实现计时器功能。该示例适用于Android平台,供开发者参考。 ... [详细]
  • Ihavethefollowingonhtml我在html上有以下内容<html><head><scriptsrc..3003_Tes ... [详细]
  • 基于Socket的多个客户端之间的聊天功能实现方法
    本文介绍了基于Socket的多个客户端之间实现聊天功能的方法,包括服务器端的实现和客户端的实现。服务器端通过每个用户的输出流向特定用户发送消息,而客户端通过输入流接收消息。同时,还介绍了相关的实体类和Socket的基本概念。 ... [详细]
  • 本文介绍了Android中的assets目录和raw目录的共同点和区别,包括获取资源的方法、目录结构的限制以及列出资源的能力。同时,还解释了raw目录中资源文件生成的ID,并说明了这些目录的使用方法。 ... [详细]
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社区 版权所有