RabbitMQ入门教程 For Java【3】 - Publish/Subscribe

2024-02-25 19:38

本文主要是介绍RabbitMQ入门教程 For Java【3】 - Publish/Subscribe,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

我的开发环境: 
操作系统: Windows7 64bit 
开发环境: JDK 1.7 - 1.7.0_55 
开发工具: Eclipse Kepler SR2 
RabbitMQ版本: 3.6.0 
Elang版本: erl7.2.1 
关于Windows7下安装RabbitMQ的教程请先在网上找一下,有空我再补安装教程。 
源码地址 
https://github.com/chwshuang/rabbitmq.git



 在上一章中,我们学习创建了一个消息队列,她的每个任务消息只发送给一个工人。这一章,我们会将同一个任务消息发送给多个工人。这种模式就是“发布/订阅”。
  • 1
  • 2

为了说明这种模式,我们将以一个日志系统进行讲解:一个日志发送者,两个日志接收者,接收者1可以把这条日志写入到磁盘上,另外一个接收者2可以将这条日志打印到控制台中。

“发布/订阅”模式的基础是将消息广播到所有的接收器上。


交换器


在之前的教程中,我们都是直接在消息队列中进行发送和接收消息,现在开始要介绍RabbitMQ完整的消息模型了。 
首先,我们先来回顾一下之前学到关于RabbitMQ的内容:

  • 生产者是发送消息的应用程序
  • 队列是存储消息的缓冲区
  • 消费者是接收消息的应用程序

实际上,RabbitMQ中消息传递模型的核心思想是:生产者不直接发送消息到队列。实际的运行环境中,生产者是不知道消息会发送到那个队列上,她只会将消息发送到一个交换器,交换器也像一个生产线,她一边接收生产者发来的消息,另外一边则根据交换规则,将消息放到队列中。交换器必须知道她所接收的消息是什么?它应该被放到那个队列中?它应该被添加到多个队列吗?还是应该丢弃?这些规则都是按照交换器的规则来确定的。 
这里写图片描述

交换器的规则有:
  • direct (直连)
  • topic (主题)
  • headers (标题)
  • fanout (分发)也有翻译为扇出的。

我们将使用【fanout】类型创建一个名称为 logs的交换器,

channel.exchangeDeclare("logs", "fanout");
  • 1

分发交换器很简单,你通过名称也能想到,她是广播所有的消息,

交换器列表 
通过rabbitmqctl list_exchanges指令列出服务器上所有可用的交换器

$ sudo rabbitmqctl list_exchanges
Listing exchanges ...direct
amq.direct      direct
amq.fanout      fanout
amq.headers     headers
amq.match       headers
amq.rabbitmq.log        topic
amq.rabbitmq.trace      topic
amq.topic       topic
logs    fanout
...done.
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10
  • 11
  • 12

这个列表里面所有以【amq.*】开头的交换器都是RabbitMQ默认创建的。在生产环境中,可以自己定义。

匿名交换器 
在之前的教程中,我们知道,发送消息到队列时根本没有使用交换器,但是消息也能发送到队列。这是因为RabbitMQ选择了一个空“”字符串的默认交换器。 
来看看我们之前的代码:

channel.basicPublish("", "hello", null, message.getBytes());
  • 1

第一个参数就是交换器的名称。如果输入“”空字符串,表示使用默认的匿名交换器。 
第二个参数是【routingKey】路由线索 
匿名交换器规则: 
发送到routingKey名称对应的队列。

现在,我们可以发送消息到交换器中:

channel.basicPublish( "logs", "", null, message.getBytes());
  • 1

临时队列


记得前两章中使用的队列指定的名称吗?(Hello World和task_queue). 
如果要在生产者和消费者之间创建一个新的队列,又不想使用原来的队列,临时队列就是为这个场景而生的:

  1. 首先,每当我们连接到RabbitMQ,我们需要一个新的空队列,我们可以用一个随机名称来创建,或者说让服务器选择一个随机队列名称给我们。
  2. 一旦我们断开消费者,队列应该立即被删除。

    在Java客户端,提供queuedeclare()为我们创建一个非持久化、独立、自动删除的队列名称。

String queueName = channel.queueDeclare().getQueue();
  • 1

通过上面的代码就能获取到一个随机队列名称。 
例如:它可能是:amq.gen-jzty20brgko-hjmujj0wlg。


绑定


这里写图片描述 
如果我们已经创建了一个分发交换器和队列,现在我们就可以就将我们的队列跟交换器进行绑定。

channel.queueBind(queueName, "logs", "");
  • 1

执行完这段代码后,日志交换器会将消息添加到我们的队列中。

绑定列表 
如果要查看绑定列表,可以执行【rabbitmqctl list_bindings】命令


全部代码


这里写图片描述

目录

这里写图片描述

生产者程序,他负责发送日志消息,与之前不同的是它不是将消息发送到匿名交换器中,而是发送到一个名为【logs】的交换器中。我们提供一个空字符串的routingkey,它的功能被交换器的分发类型代替了。下面是EmitLog.java的代码:

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;public class EmitLog {private static final String EXCHANGE_NAME = "logs";public static void main(String[] argv) throws Exception {ConnectionFactory factory = new ConnectionFactory();factory.setHost("localhost");Connection connection = factory.newConnection();Channel channel = connection.createChannel();channel.exchangeDeclare(EXCHANGE_NAME, "fanout");//      分发消息for(int i = 0 ; i < 5; i++){String message = "Hello World! " + i;channel.basicPublish(EXCHANGE_NAME, "", null, message.getBytes());System.out.println(" [x] Sent '" + message + "'");}channel.close();connection.close();}
}
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10
  • 11
  • 12
  • 13
  • 14
  • 15
  • 16
  • 17
  • 18
  • 19
  • 20
  • 21
  • 22
  • 23
  • 24
  • 25
  • 26
  • 27

上面的代码中,在建立连接后,我们声明了一个交互。如果当前没有队列被绑定到交换器,消息将被丢弃,因为没有消费者监听,这条消息将被丢弃。

下面的代码是接收日志ReceiveLogs1.java 和ReceiveLogs2.java:

import com.rabbitmq.client.*;import java.io.IOException;public class ReceiveLogs1 {private static final String EXCHANGE_NAME = "logs";public static void main(String[] argv) throws Exception {ConnectionFactory factory = new ConnectionFactory();factory.setHost("localhost");Connection connection = factory.newConnection();Channel channel = connection.createChannel();channel.exchangeDeclare(EXCHANGE_NAME, "fanout");String queueName = channel.queueDeclare().getQueue();channel.queueBind(queueName, EXCHANGE_NAME, "");System.out.println(" [*] Waiting for messages. To exit press CTRL+C");Consumer consumer = new DefaultConsumer(channel) {@Overridepublic void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {String message = new String(body, "UTF-8");System.out.println(" [x] Received '" + message + "'");}};channel.basicConsume(queueName, true, consumer);}
}
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10
  • 11
  • 12
  • 13
  • 14
  • 15
  • 16
  • 17
  • 18
  • 19
  • 20
  • 21
  • 22
  • 23
  • 24
  • 25
  • 26
  • 27
  • 28
  • 29
  • 30
import com.rabbitmq.client.*;import java.io.IOException;public class ReceiveLogs1 {private static final String EXCHANGE_NAME = "logs";public static void main(String[] argv) throws Exception {ConnectionFactory factory = new ConnectionFactory();factory.setHost("localhost");Connection connection = factory.newConnection();Channel channel = connection.createChannel();channel.exchangeDeclare(EXCHANGE_NAME, "fanout");String queueName = channel.queueDeclare().getQueue();channel.queueBind(queueName, EXCHANGE_NAME, "");System.out.println(" [*] Waiting for messages. To exit press CTRL+C");Consumer consumer = new DefaultConsumer(channel) {@Overridepublic void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {String message = new String(body, "UTF-8");System.out.println(" [x] Received '" + message + "'");}};channel.basicConsume(queueName, true, consumer);}
}
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10
  • 11
  • 12
  • 13
  • 14
  • 15
  • 16
  • 17
  • 18
  • 19
  • 20
  • 21
  • 22
  • 23
  • 24
  • 25
  • 26
  • 27
  • 28
  • 29
  • 30

运行


先运行ReceiveLogs1和ReceiveLogs2可以看到日志:

 [*] Waiting for messages. To exit press CTRL+C
  • 1

然后运行EmitLog:

EmitLog日志:[x] Sent 'Hello World! 0'[x] Sent 'Hello World! 1'[x] Sent 'Hello World! 2'[x] Sent 'Hello World! 3'[x] Sent 'Hello World! 4'ReceiveLogs1和ReceiveLogs2日志[*] Waiting for messages. To exit press CTRL+C[x] Received 'Hello World! 0'[x] Received 'Hello World! 1'[x] Received 'Hello World! 2'[x] Received 'Hello World! 3'[x] Received 'Hello World! 4'
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10
  • 11
  • 12
  • 13
  • 14

看到这里,说明我们的程序运行正常,消费者通过声明【logs】交换器和【fanout】类型,接收到了来自【logs】交换器的所有消息。

使用【rabbitmqctl list_bindings】命令可以看到两个临时队列的名称

$ sudo rabbitmqctl list_bindings
Listing bindings ...
logs    exchange        amq.gen-JzTY20BRgKO-HjmUJj0wLg  queue           []
logs    exchange        amq.gen-vso0PVvyiRIL2WoV3i48Yg  queue           []
...done.
  • 1
  • 2
  • 3
  • 4
  • 5

以上就是这一章讲的发布/订阅模式,下一章将介绍消息路由(Routing)

这篇关于RabbitMQ入门教程 For Java【3】 - Publish/Subscribe的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/746489

相关文章

Java学习手册之Filter和Listener使用方法

《Java学习手册之Filter和Listener使用方法》:本文主要介绍Java学习手册之Filter和Listener使用方法的相关资料,Filter是一种拦截器,可以在请求到达Servl... 目录一、Filter(过滤器)1. Filter 的工作原理2. Filter 的配置与使用二、Listen

Spring Boot中JSON数值溢出问题从报错到优雅解决办法

《SpringBoot中JSON数值溢出问题从报错到优雅解决办法》:本文主要介绍SpringBoot中JSON数值溢出问题从报错到优雅的解决办法,通过修改字段类型为Long、添加全局异常处理和... 目录一、问题背景:为什么我的接口突然报错了?二、为什么会发生这个错误?1. Java 数据类型的“容量”限制

Java对象转换的实现方式汇总

《Java对象转换的实现方式汇总》:本文主要介绍Java对象转换的多种实现方式,本文通过实例代码给大家介绍的非常详细,对大家的学习或工作具有一定的参考借鉴价值,需要的朋友参考下吧... 目录Java对象转换的多种实现方式1. 手动映射(Manual Mapping)2. Builder模式3. 工具类辅助映

SpringBoot请求参数接收控制指南分享

《SpringBoot请求参数接收控制指南分享》:本文主要介绍SpringBoot请求参数接收控制指南,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录Spring Boot 请求参数接收控制指南1. 概述2. 有注解时参数接收方式对比3. 无注解时接收参数默认位置

SpringBoot基于配置实现短信服务策略的动态切换

《SpringBoot基于配置实现短信服务策略的动态切换》这篇文章主要为大家详细介绍了SpringBoot在接入多个短信服务商(如阿里云、腾讯云、华为云)后,如何根据配置或环境切换使用不同的服务商,需... 目录目标功能示例配置(application.yml)配置类绑定短信发送策略接口示例:阿里云 & 腾

SpringBoot项目中报错The field screenShot exceeds its maximum permitted size of 1048576 bytes.的问题及解决

《SpringBoot项目中报错ThefieldscreenShotexceedsitsmaximumpermittedsizeof1048576bytes.的问题及解决》这篇文章... 目录项目场景问题描述原因分析解决方案总结项目场景javascript提示:项目相关背景:项目场景:基于Spring

Spring Boot 整合 SSE的高级实践(Server-Sent Events)

《SpringBoot整合SSE的高级实践(Server-SentEvents)》SSE(Server-SentEvents)是一种基于HTTP协议的单向通信机制,允许服务器向浏览器持续发送实... 目录1、简述2、Spring Boot 中的SSE实现2.1 添加依赖2.2 实现后端接口2.3 配置超时时

Spring Boot读取配置文件的五种方式小结

《SpringBoot读取配置文件的五种方式小结》SpringBoot提供了灵活多样的方式来读取配置文件,这篇文章为大家介绍了5种常见的读取方式,文中的示例代码简洁易懂,大家可以根据自己的需要进... 目录1. 配置文件位置与加载顺序2. 读取配置文件的方式汇总方式一:使用 @Value 注解读取配置方式二

一文详解Java异常处理你都了解哪些知识

《一文详解Java异常处理你都了解哪些知识》:本文主要介绍Java异常处理的相关资料,包括异常的分类、捕获和处理异常的语法、常见的异常类型以及自定义异常的实现,文中通过代码介绍的非常详细,需要的朋... 目录前言一、什么是异常二、异常的分类2.1 受检异常2.2 非受检异常三、异常处理的语法3.1 try-

Java中的@SneakyThrows注解用法详解

《Java中的@SneakyThrows注解用法详解》:本文主要介绍Java中的@SneakyThrows注解用法的相关资料,Lombok的@SneakyThrows注解简化了Java方法中的异常... 目录前言一、@SneakyThrows 简介1.1 什么是 Lombok?二、@SneakyThrows