IT数码 购物 网址 头条 软件 日历 阅读 图书馆
TxT小说阅读器
↓语音阅读,小说下载,古典文学↓
图片批量下载器
↓批量下载图片,美女图库↓
图片自动播放器
↓图片自动播放器↓
一键清除垃圾
↓轻轻一点,清除系统垃圾↓
开发: C++知识库 Java知识库 JavaScript Python PHP知识库 人工智能 区块链 大数据 移动开发 嵌入式 开发工具 数据结构与算法 开发测试 游戏开发 网络协议 系统运维
教程: HTML教程 CSS教程 JavaScript教程 Go语言教程 JQuery教程 VUE教程 VUE3教程 Bootstrap教程 SQL数据库教程 C语言教程 C++教程 Java教程 Python教程 Python3教程 C#教程
数码: 电脑 笔记本 显卡 显示器 固态硬盘 硬盘 耳机 手机 iphone vivo oppo 小米 华为 单反 装机 图拉丁
 
   -> 大数据 -> RabbitMQ 学习笔记3 - Java 使用 RabbitMQ 示例 -> 正文阅读

[大数据]RabbitMQ 学习笔记3 - Java 使用 RabbitMQ 示例

1. 背景

本节讲述 Java 使用 RabbitMQ 的示例,和 发送者确认回调,消费者回执的内容。

2.知识

高级消息队列协议 (AMQP) 是面向消息的中间件的平台中立的协议。Spring AMQP 项目将 Spring 的概念应用于 AMQP,形成解决方案的开发。

AMQP 的一些基本概念:
开始之前, 要使用 RabbitMQ 首先要了解 AMQP 协议的基本概念,更多可阅读我的另一篇文章

  • 生产者:一个发送消息的程序,它产生消息并发送到队列。这里是用Go写的发送端示程序例。
  • 消息队列:即 RabbitMQ 内部的队列,它安装在一个服务器中。做为消息中间件,它与具体开发语言无关,支持 Go,Java等接入连接。
  • 消费者:消费者是一个等待消息,接收消息的接收端程序示例
  • 交换机(Exchange)可以理解成邮局,交换机将收到的消息根据路由规则分发给绑定的队列(Queue)
image.png

安装 RabbitMQ
参考我的另一篇文章:https://www.jianshu.com/p/53ba4fbd0d03

我们使用 Spring AMQP 框架来 操作 RabbitMQ 收发消息。

Spring AMQP 框架

Spring AMQP 项目将核心 Spring 概念应用于基于 AMQP 的消息传递解决方案的开发。它提供了一个“模板”作为发送和接收消息的高级抽象。

该项目由两部分组成;spring-amqp 是基础抽象,spring-rabbit 是 RabbitMQ 实现。

Spring AMQP 的特征

  • 用于异步处理入站消息的侦听器容器
  • RabbitTemplate 用于发送和接收消息
  • RabbitAdmin 用于自动声明队列、交换和绑定

3. 示例

下面通过一个示例看下基本的收发消息的操作。

3.1 编写程序“生产者”

第一步:配置Rabbit的数据连接

编辑 application.yml, 指定 rabbitmq 的服务器地址,端口号,账户名密码等。

spring:
  application:
    name: producer
  rabbitmq:
    host: localhost
    virtual-host: /
    port: 5672
    username: admin
    password: admin

第二步:配置好 队列,交换机,和绑定(queue,exchange,binding)

队列里存储了消息,交换机类似邮局,而“绑定”是个“ 队列+交换机”关联关系。通过“绑定” binding 将 交换机和 队列连线在一起。

@Configuration
public class RabbitConfig {
    // 路由的 key
    public static final String ROUTING_KEY = "hello_routing_key";
    public static final String EXCHANGE_NAME = "zyf_direct_exchange";
    public static final String QUEUE_NAME = "first_queue";

    //  一个队列
    @Bean
    public Queue getFirstQueue() {
        return new Queue(QUEUE_NAME);
    }

    // 一个 直接交换机
    @Bean
    public DirectExchange getDirectExchange() {
        return new DirectExchange(EXCHANGE_NAME);
    }

    // 进行绑定
    @Bean
    public Binding getBinding() {
        return BindingBuilder.bind(getFirstQueue()).to(getDirectExchange()).with(ROUTING_KEY);
    }
}

第三步:发消息

RabbitTemplate 是操作发送消息的 “模板方法”,springboot 已帮忙配置好注入关系,直接拿来用就可以了。调用: rabbitTemplate.convertAndSend(...) 来发消息,需要指明 交换机的名称,和路由key。

@Service
public class BusinessService {

    @Autowired
    private RabbitTemplate rabbitTemplate;


    public void sendMessage(String msg) {
        System.out.println("# sendMessage,msg=" + msg);


        String routingKey = RabbitConfig.ROUTING_KEY;
        String exchangeName = RabbitConfig.EXCHANGE_NAME;
        rabbitTemplate.convertAndSend(exchangeName, routingKey, msg);
    }

}

总结:

  • 先配置好从 “交换机”到“队列”的连线。
  • 发送者 就可以通过 交换机名称和路由 key 来发送消息。

3.2 编写程序“消费者”

然后就是准备接收消息了。

第一步:配置好 rabbitmq 的数据连接。

和上面的 发送者一样,编辑 application.yml, 指定 rabbitmq 的服务器地址,端口号,账户名密码等。

第二步:配置 异步消息的监听器

接收消息配置一个回调即可。使用 @RabbitMessageListener 注解标注。

@Component
public class RabbitMessageListener {
    public static final String QUEUE_NAME = "first_queue";

    @RabbitListener(queues = QUEUE_NAME)
    public void receive(String msg) {
        System.out.println(msg);
    }
}

至此就完成了收发消息。我的代码示例见:https://github.com/vir56k/java_demo/tree/master/rabbitmq_demo1

4. 更多扩展

4.1 生产者发送时的结果回调(确认模式)

发布是异步的——如何检测成功和失败?

发布消息是一种异步机制,默认情况下,"无法路由的消息" 会被 RabbitMQ 丢弃。为了成功发布,您可以收到异步确认,如相关发布者确认和返回 中所述。

考虑两种失败情况:

  • 发消息到不存在的交换机。
  • 发消息到交换机,但没有匹配的队列。

第一种情况的场景是 指定了 错误的交换机名称。
第二种情况的场景是 “发送者的退货” 。

(1)发送者发送消息后的 “消息确认” 回调事件
对于发布者确认 ,RabbitTemplate 需要 设置:

connectionFactory.setPublisherConfirmType(CachingConnectionFactory.ConfirmType.CORRELATED);
和
rabbitTemplate.setConfirmCallback(confirmCallback());

然后注册回调的实习,示例:

// 用来确认生产者 producer 将消息发送到 broker 的回调
 @Bean
 public RabbitTemplate.ConfirmCallback confirmCallback() {
     return new RabbitTemplate.ConfirmCallback() {
         @Override
         public void confirm(CorrelationData correlationData, boolean ack, String cause) {
             String id = correlationData == null ? "" : correlationData.getId();
             if (ack) {
                 // log.info(String.format("# [%s] 投递到 broker 成功!, cause=%s", id, cause));
             } else {
                 log.error(String.format("# [%s] 投递到 broker 失败!, cause=%s", id, cause));
             }
         }
     };
 }

上面的 ack 参数 指示了是否投递(到交换机)成功。

注意:一个 ConfirmCallback 仅支持 一个RabbitTemplate。

**(2)发送者的 “退货” 回调事件

对于返回的消息,模板的 mandatory 属性必须设置为true 。也需要将 CachingConnectionFactory 其 publisherReturns属性设置为true
即:

connectionFactory.setPublisherReturns(true);
...
rabbitTemplate.setReturnsCallback(returnsCallback());
// 强制标志,当 setReturnsCallback 被设置时,这里要设置为 true
rabbitTemplate.setMandatory(true);

示例:

// 发布到交换机,但没有匹配的目标队列 时,退货
@Bean
public RabbitTemplate.ReturnsCallback returnsCallback() {
    return new RabbitTemplate.ReturnsCallback() {
        @Override
        public void returnedMessage(ReturnedMessage returnedMessage) {
            int replyCode = returnedMessage.getReplyCode();
            String replyText = returnedMessage.getReplyText();
            String exchange = returnedMessage.getExchange();
            System.out.println(String.format("# 退货消息:原因=%s, replyCode=%s, exchange=%s", replyText, replyCode, exchange));
        }
    };
}

上面方法里的 ReturnedMessage 具有以下属性:

  • message - 返回的消息本身
  • replyCode - 指示退货原因的代码
  • replyText - 退货的文字原因 - 例如 NO_ROUTE
  • exchange - 消息发送到的交换
  • routingKey - 使用的路由密钥

每个 ReturnsCallback 仅支持一个RabbitTemplate。

4.2消费者回执(确认模式)

消息接收回执是指 消息接收者 收到消息后 向 “broker” 消息代理 回复的“ 确认消息 ”

注意:这里的回执和 发送者 “没有任何关系” 。它通知到 rabbitmq ,rabbitmq 根据回执决定是 重复,或者放弃。

有三种回执模式:

  • NONE:不发送确认。RabbitMQ 将此称为“自动确认”,因为代理假定所有消息都已确认,而消费者没有采取任何行动。
  • MANUAL:侦听器必须通过调用来确认所有消息Channel.basicAck()。
  • AUTO:容器自动确认消息,除非MessageListener抛出异常。

实现手动回执时,注入 ackMode 就可以了。

@RabbitListener(queues = QUEUE_NAME, ackMode = "MANUAL")

示例:

int i = 0;

    // 异步 接收消息
    @RabbitListener(queues = QUEUE_NAME, ackMode = "MANUAL")
    public void receiveForManual(String msg, Message message, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException {
        System.out.println(msg);

        String returned_message_correlation = message.getMessageProperties().getHeader("spring_returned_message_correlation");
        log.info(String.format("# 触发 receiveForManual msg =%s ,correlationId = %s", msg, returned_message_correlation));

        i++;
        if (i % 3 == 0) {
            log.info("# 接收消息...");
            channel.basicAck(tag, false);
        } else if (i % 3 == 1) {
            log.error("# 消息已重复处理失败,拒绝再次接收...");
            channel.basicReject(tag, false);
        } else if (i % 3 == 2) {
            log.error("# 消息即将再次返回队列处理...");
            channel.basicNack(tag, false, true);
        }

    }

我的代码示例见:https://github.com/vir56k/java_demo/tree/master/rabbitmq_demo2

5.参考:

Spring AMQP 文档https://spring.io/projects/spring-amqp

https://docs.spring.io/spring-boot/docs/current/reference/html/features.html#features.messaging.amqp

https://github.com/spring-projects/spring-amqp-samples

https://docs.spring.io/spring-framework/docs/current/reference/html/integration.html#scheduling

  大数据 最新文章
实现Kafka至少消费一次
亚马逊云科技:还在苦于ETL?Zero ETL的时代
初探MapReduce
【SpringBoot框架篇】32.基于注解+redis实现
Elasticsearch:如何减少 Elasticsearch 集
Go redis操作
Redis面试题
专题五 Redis高并发场景
基于GBase8s和Calcite的多数据源查询
Redis——底层数据结构原理
上一篇文章      下一篇文章      查看所有文章
加:2021-07-22 23:00:39  更:2021-07-22 23:01:12 
 
开发: C++知识库 Java知识库 JavaScript Python PHP知识库 人工智能 区块链 大数据 移动开发 嵌入式 开发工具 数据结构与算法 开发测试 游戏开发 网络协议 系统运维
教程: HTML教程 CSS教程 JavaScript教程 Go语言教程 JQuery教程 VUE教程 VUE3教程 Bootstrap教程 SQL数据库教程 C语言教程 C++教程 Java教程 Python教程 Python3教程 C#教程
数码: 电脑 笔记本 显卡 显示器 固态硬盘 硬盘 耳机 手机 iphone vivo oppo 小米 华为 单反 装机 图拉丁

360图书馆 购物 三丰科技 阅读网 日历 万年历 2024年5日历 -2024/5/6 16:52:39-

图片自动播放器
↓图片自动播放器↓
TxT小说阅读器
↓语音阅读,小说下载,古典文学↓
一键清除垃圾
↓轻轻一点,清除系统垃圾↓
图片批量下载器
↓批量下载图片,美女图库↓
  网站联系: qq:121756557 email:121756557@qq.com  IT数码