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



推荐阅读
  • 技术分享:使用 Flask、AngularJS 和 Jinja2 构建高效前后端交互系统
    技术分享:使用 Flask、AngularJS 和 Jinja2 构建高效前后端交互系统 ... [详细]
  • 本文介绍了在 Java 编程中遇到的一个常见错误:对象无法转换为 long 类型,并提供了详细的解决方案。 ... [详细]
  • 本文详细介绍了Java反射机制的基本概念、获取Class对象的方法、反射的主要功能及其在实际开发中的应用。通过具体示例,帮助读者更好地理解和使用Java反射。 ... [详细]
  • Spring – Bean Life Cycle
    Spring – Bean Life Cycle ... [详细]
  • DAO(Data Access Object)模式是一种用于抽象和封装所有对数据库或其他持久化机制访问的方法,它通过提供一个统一的接口来隐藏底层数据访问的复杂性。 ... [详细]
  • 本文探讨了如何在 Java 中将多参数方法通过 Lambda 表达式传递给一个接受 List 的 Function。具体分析了 `OrderUtil` 类中的 `runInBatches` 方法及其使用场景。 ... [详细]
  • oracle c3p0 dword 60,web_day10 dbcp c3p0 dbutils
    createdatabasemydbcharactersetutf8;alertdatabasemydbcharactersetutf8;1.自定义连接池为了不去经常创建连接和释放 ... [详细]
  • 本教程详细介绍了如何使用 Spring Boot 创建一个简单的 Hello World 应用程序。适合初学者快速上手。 ... [详细]
  • 深入解析 Lifecycle 的实现原理
    本文将详细介绍 Android Jetpack 中 Lifecycle 组件的实现原理,帮助开发者更好地理解和使用 Lifecycle,避免常见的内存泄漏问题。 ... [详细]
  • 解决Bootstrap DataTable Ajax请求重复问题
    在最近的一个项目中,我们使用了JQuery DataTable进行数据展示,虽然使用起来非常方便,但在测试过程中发现了一个问题:当查询条件改变时,有时查询结果的数据不正确。通过FireBug调试发现,点击搜索按钮时,会发送两次Ajax请求,一次是原条件的请求,一次是新条件的请求。 ... [详细]
  • 如何在Java中使用DButils类
    这期内容当中小编将会给大家带来有关如何在Java中使用DButils类,文章内容丰富且以专业的角度为大家分析和叙述,阅读完这篇文章希望大家可以有所收获。D ... [详细]
  • 秒建一个后台管理系统?用这5个开源免费的Java项目就够了
    秒建一个后台管理系统?用这5个开源免费的Java项目就够了 ... [详细]
  • 在本阶段的Java编程实战中,我们将深入探讨位运算的应用。具体任务是实现逻辑位运算。用户需从键盘输入一个位运算符(如AND、OR、XOR或NOT)及相应的操作数,系统将根据输入的运算符执行相应的位运算并输出结果。此练习旨在加强学员对位运算的理解和实际操作能力。 ... [详细]
  • 属性类 `Properties` 是 `Hashtable` 类的子类,用于存储键值对形式的数据。该类在 Java 中广泛应用于配置文件的读取与写入,支持字符串类型的键和值。通过 `Properties` 类,开发者可以方便地进行配置信息的管理,确保应用程序的灵活性和可维护性。此外,`Properties` 类还提供了加载和保存属性文件的方法,使其在实际开发中具有较高的实用价值。 ... [详细]
  • 在Java编程中,初始化List集合有多种高效的方法。本文介绍了六种常见的技术,包括使用常规方式、Arrays.asList、Collections.addAll、Java 8的Stream API、双重大括号初始化以及使用List.of。每种方法都有其特定的应用场景和优缺点,开发者可以根据实际需求选择最合适的方式。例如,常规方式通过直接创建ArrayList对象并逐个添加元素,适用于需要动态修改列表的情况;而List.of则提供了一种简洁的不可变列表初始化方式,适合于固定数据集的场景。 ... [详细]
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社区 版权所有