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



推荐阅读
  • 个人学习使用:谨慎参考1Client类importcom.thoughtworks.gauge.Step;importcom.thoughtworks.gauge.T ... [详细]
  • Java容器中的compareto方法排序原理解析
    本文从源码解析Java容器中的compareto方法的排序原理,讲解了在使用数组存储数据时的限制以及存储效率的问题。同时提到了Redis的五大数据结构和list、set等知识点,回忆了作者大学时代的Java学习经历。文章以作者做的思维导图作为目录,展示了整个讲解过程。 ... [详细]
  • 本文介绍了如何在给定的有序字符序列中插入新字符,并保持序列的有序性。通过示例代码演示了插入过程,以及插入后的字符序列。 ... [详细]
  • Spring特性实现接口多类的动态调用详解
    本文详细介绍了如何使用Spring特性实现接口多类的动态调用。通过对Spring IoC容器的基础类BeanFactory和ApplicationContext的介绍,以及getBeansOfType方法的应用,解决了在实际工作中遇到的接口及多个实现类的问题。同时,文章还提到了SPI使用的不便之处,并介绍了借助ApplicationContext实现需求的方法。阅读本文,你将了解到Spring特性的实现原理和实际应用方式。 ... [详细]
  • 本文讨论了一个关于cuowu类的问题,作者在使用cuowu类时遇到了错误提示和使用AdjustmentListener的问题。文章提供了16个解决方案,并给出了两个可能导致错误的原因。 ... [详细]
  • [大整数乘法] java代码实现
    本文介绍了使用java代码实现大整数乘法的过程,同时也涉及到大整数加法和大整数减法的计算方法。通过分治算法来提高计算效率,并对算法的时间复杂度进行了研究。详细代码实现请参考文章链接。 ... [详细]
  • Spring源码解密之默认标签的解析方式分析
    本文分析了Spring源码解密中默认标签的解析方式。通过对命名空间的判断,区分默认命名空间和自定义命名空间,并采用不同的解析方式。其中,bean标签的解析最为复杂和重要。 ... [详细]
  • SpringBoot uri统一权限管理的实现方法及步骤详解
    本文详细介绍了SpringBoot中实现uri统一权限管理的方法,包括表结构定义、自动统计URI并自动删除脏数据、程序启动加载等步骤。通过该方法可以提高系统的安全性,实现对系统任意接口的权限拦截验证。 ... [详细]
  • 本文分享了一个关于在C#中使用异步代码的问题,作者在控制台中运行时代码正常工作,但在Windows窗体中却无法正常工作。作者尝试搜索局域网上的主机,但在窗体中计数器没有减少。文章提供了相关的代码和解决思路。 ... [详细]
  • javascript  – 概述在Firefox上无法正常工作
    我试图提出一些自定义大纲,以达到一些Web可访问性建议.但我不能用Firefox制作.这就是它在Chrome上的外观:而那个图标实际上是一个锚点.在Firefox上,它只概述了整个 ... [详细]
  • JavaSE笔试题-接口、抽象类、多态等问题解答
    本文解答了JavaSE笔试题中关于接口、抽象类、多态等问题。包括Math类的取整数方法、接口是否可继承、抽象类是否可实现接口、抽象类是否可继承具体类、抽象类中是否可以有静态main方法等问题。同时介绍了面向对象的特征,以及Java中实现多态的机制。 ... [详细]
  • 本文介绍了计算机网络的定义和通信流程,包括客户端编译文件、二进制转换、三层路由设备等。同时,还介绍了计算机网络中常用的关键词,如MAC地址和IP地址。 ... [详细]
  • importjava.util.ArrayList;publicclassPageIndex{privateintpageSize;每页要显示的行privateintpageNum ... [详细]
  • 本文探讨了C语言中指针的应用与价值,指针在C语言中具有灵活性和可变性,通过指针可以操作系统内存和控制外部I/O端口。文章介绍了指针变量和指针的指向变量的含义和用法,以及判断变量数据类型和指向变量或成员变量的类型的方法。还讨论了指针访问数组元素和下标法数组元素的等价关系,以及指针作为函数参数可以改变主调函数变量的值的特点。此外,文章还提到了指针在动态存储分配、链表创建和相关操作中的应用,以及类成员指针与外部变量的区分方法。通过本文的阐述,读者可以更好地理解和应用C语言中的指针。 ... [详细]
  • 猜字母游戏
    猜字母游戏猜字母游戏——设计数据结构猜字母游戏——设计程序结构猜字母游戏——实现字母生成方法猜字母游戏——实现字母检测方法猜字母游戏——实现主方法1猜字母游戏——设计数据结构1.1 ... [详细]
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社区 版权所有