ActiveMQ Topic消息重发

MQ学习系列:

  1. 消息队列概念与认知
  2. ActiveMQ Topic消息重发

ActiveMQ Topic 消息重发

准备工作

windows下ActiveMQ的下载与启动

  • 百度的教程:链接 ←这里包含基本的下载安装启动以及简单的配置账号
  • 登录控制台主页:http://localhost:8161/admin/

启动错误以及解决方案

activeMQ启动错误 BeanFactory not initialized

JMS 消息确认机制

在session接口中定义的几个常量:

  • AUTO_ACKNOWLEDGE = 1 自动确认
  • CLIENT_ACKNOWLEDGE = 2 客户端手动确认
  • DUPS_OK_ACKNOWLEDGE = 3 自动批量确认
  • SESSION_TRANSACTED = 0 事务提交并确认

代码实现

消息消费端在创建Session对象时需要指定应答模式为客户端手动应答,当消费者获取到消息并成功处理后需要调用message.acknowledge()方法进行应答,通知Broker消费成功。如果处理过程中出现异常,需要调用session.recover()通知Broker重复消息,默认最多重复6次。

  1. 创建maven项目引入依赖
<dependencies>
    <!-- https://mvnrepository.com/artifact/org.apache.activemq/activemq-client -->
    <dependency>
        <groupId>org.apache.activemq</groupId>
        <artifactId>activemq-client</artifactId>
        <version>5.15.8</version>
    </dependency>
    <!-- https://mvnrepository.com/artifact/junit/junit -->
    <dependency>
        <groupId>junit</groupId>
        <artifactId>junit</artifactId>
        <version>4.12</version>
        <scope>test</scope>
    </dependency>
</dependencies>
  1. 编写测试方法模拟【无消息重发的正常情况】
package org.newmean;

import org.apache.activemq.ActiveMQConnectionFactory;
import org.junit.Test;

import javax.jms.*;

public class ActiveMQTest {
    //消息发送方-producter
    @Test
    public void test1() throws JMSException {
        //创建连接工厂对象
        ConnectionFactory connectionFactory = new       ActiveMQConnectionFactory("tcp://localhost:61616");
        //从工厂中获取一个连接对象
        Connection connection = connectionFactory.createConnection();
        //连接MQ服务
        connection.start();
        //获取session对象
        //参数说明 b 是否使用事务 i jms消息确认机制 1 2 3 0 用常量表示
        Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
        //通过session创建Topic
        Topic topic = session.createTopic("TestTopic");
        //通过session创建消息发送者
        MessageProducer producer = session.createProducer(topic);
        //通过session创建消息对象
        TextMessage message = session.createTextMessage("hello");
        //发送消息
        producer.send(message);
        //关闭资源
        producer.close();
        session.close();
        connection.close();
    }
    //消息接收方-consumer
    @Test
    public void test2() throws JMSException {
        //创建连接工厂对象
        ConnectionFactory connectionFactory = new ActiveMQConnectionFactory("tcp://localhost:61616");
        //从工厂中获取一个连接对象
        Connection connection = connectionFactory.createConnection();
        //连接MQ服务
        connection.start();
        //获取session对象
        //参数说明 b 是否使用事务 i jms消息确认机制 1 2 3 0 用常量表示
        Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
        //通过session创建Topic
        Topic topic = session.createTopic("TestTopic");
        //通过session创建消费者
        MessageConsumer consumer = session.createConsumer(topic);
        //指定消息监听器
        consumer.setMessageListener(new MessageListener() {
            //当我们监听的topic中存在消息,onMessage这个方法就会自动运行
            public void onMessage(Message message) {
                TextMessage textMessage = (TextMessage) message;
                try {
                    System.out.println("消费者接收到了消息:"+textMessage.getText());
                } catch (JMSException e) {
                    e.printStackTrace();
                }
            }
        });
        //因为要接收消息不能关闭,同时线程不能死掉
        while (true){

        }

    }
}

先启动test2方法发起订阅“TestTopic”消息,然后启动test1方法,这时消费者收到了消息。

  1. 消息重发模拟

    我们只需要更消息接收方的代码,改动如下:

//消息接收方-consumer
    @Test
    public void test2() throws JMSException {
        //创建连接工厂对象
        ConnectionFactory connectionFactory = new ActiveMQConnectionFactory("tcp://localhost:61616");
        //从工厂中获取一个连接对象
        Connection connection = connectionFactory.createConnection();
        //连接MQ服务
        connection.start();
        //获取session对象
        //参数说明 b 是否使用事务 i jms消息确认机制 1 2 3 0 用常量表示
        final Session session = connection.createSession(false, Session.CLIENT_ACKNOWLEDGE);
        //通过session创建Topic
        Topic topic = session.createTopic("TestTopic");
        //通过session创建消费者
        MessageConsumer consumer = session.createConsumer(topic);
        //指定消息监听器
        consumer.setMessageListener(new MessageListener() {
            //当我们监听的topic中存在消息,onMessage这个方法就会自动运行
            public void onMessage(Message message) {
                TextMessage textMessage = (TextMessage) message;
                try {
                    if(textMessage.getText().equals("nihao")){
                        System.out.println("消费者接收到了消息:"+textMessage.getText());
                        message.acknowledge();
                    }else {
                        System.out.println("消息处理失败了..");
                        session.recover();
                    }

                } catch (JMSException e) {
                    e.printStackTrace();
                }
            }
        });
        //因为要接收消息不能关闭,同时线程不能死掉
        while (true){

        }

    }

先启动test2方法发起订阅“TestTopic”消息,然后启动test1方法,这时消费者就会调用session.recover()方法让消息发布者重发消息默认6次,我们能够看到7条(第一次+重发六次)“消息处理失败了..”输出。

原文地址:https://www.cnblogs.com/nm666/p/10409984.html

时间: 2024-10-10 11:21:13

ActiveMQ Topic消息重发的相关文章

ActiveMQ中消息的重发与持久化保存

消息中间件解决方案续 上一篇中我们讲到了在Spring工程中基本的使用消息中间件,这里就不在继续赘述. 针对消息中间件,这篇讲解两个我们常遇到的问题. 问题1:如果我们的消息的接收过程中发生异常,怎么解决? 问题2:发布订阅模式(Topic)下如果消费端宕机引起的消息的丢失,怎么解决? 问题解决方案: 问题1暂时有两种解决方案:第一种是开启消息确认机制,第二种开启事务.下面会在点对点模式下进行演示. 问题2的解决方案:实现发布订阅消息的持久化保存. 根据上面的解决方案搭建工程如下:测试消息的重发

解决Springboot整合ActiveMQ发送和接收topic消息的问题

环境搭建 1.创建maven项目(jar) 2.pom.xml添加依赖 <parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>1.4.0.RELEASE</version> </parent> <dependencies> &l

RocketMQ(消息重发、重复消费、事务、消息模式)

RocketMQ基础:https://github.com/apache/rocketmq/tree/rocketmq-all-4.5.1/docs/cn 分布式消息系统作为实现分布式系统可扩展.可伸缩性的关键组件,需要具有高吞吐量.高可用等特点.而谈到消息系统的设计,就回避不了两个问题: 消息的顺序问题 消息的重复问题 RocketMQ作为阿里开源的一款高性能.高吞吐量的消息中间件,它是怎样来解决这两个问题的?RocketMQ 有哪些关键特性?其实现原理是怎样的? 关键特性以及其实现原理 一.

JMS开发步骤和持久化/非持久化Topic消息

------------------------------------------------ 开发一个JMS的基本步骤如下: 1.创建一个JMS connection factory 2.通过connection factory来创建JMS connection 3.启动JMS connection 4.通过connection创建JMS session 5.创建JMS destination 6.创建JMS producer 或者创建JMS message,并设置destination 7

ActiveMQ的消息持久化机制

为了避免意外宕机以后丢失信息,需要做到重启后可以恢复消息队列,消息系统一般都会采用持久化机制. ActiveMQ的消息持久化机制有JDBC,AMQ,KahaDB和LevelDB,无论使用哪种持久化方式,消息的存储逻辑都是一致的. 就是在发送者将消息发送出去后,消息中心首先将消息存储到本地数据文件.内存数据库或者远程数据库等,然后试图将消息发送给接收者,发送成功则将消息从存储中删除,失败则继续尝试. 消息中心启动以后首先要检查指定的存储位置,如果有未发送成功的消息,则需要把消息发送出去. 1. J

activemq的消息确认机制ACK

一.简介 消息消费者有没有接收到消息,需要有一种机制让消息提供者知道,这个机制就是消息确认机制. ACK(Acknowledgement)即确认字符,在数据通信中,接收站发给发送站的一种传输类控制字符.表示发来的数据已确认接收无误. 二.ACK_MODE有几类 我们在开发JMS应用程序的时候,会经常使用到上述ACK_MODE,其中"INDIVIDUAL_ACKNOWLEDGE "只有ActiveMQ支持,当然开发者也可以使用它. ACK_MODE描述了Consumer与broker确认

消息中间件--ActiveMQ&amp;JMS消息服务

### 消息中间件 ### ---------- **消息中间件** 1. 消息中间件的概述 2. 消息中间件的应用场景(查看大纲文档,了解消息队列的应用场景) * 异步处理 * 应用解耦 * 流量削峰 * 消息通信 ---------- ### JMS消息服务 ### ---------- **JMS的概述** 1. JMS消息服务的概述 2. JMS消息模型 * P2P模式 * Pub/Sub模式 3. 消息消费的方式 * 同步的方式---手动 * 异步的方式---listener监听 4.

ActiveMQ的学习(三)(ActiveMQ的消息事务和消息的确认机制)

ActiveMQ的消息事务 消息事务,是保证消息传递原子性的一个重要特性,和JDBC的事务特征类似. 一个事务性发送,其中一组消息要么能够全部保证到达服务器,要么都不到达服务器.生产者,消费者与消息服务器都支持事务性.ActiveMQ得事务主要偏向在生产者得应用. ActiveMQ消息事务流程图: 原生jms事务发送(生产者的事务发送) 不加事务得情况:(程序没有错误,10条消息会到达mq中) 不加事务得情况:(程序有错误,结果是发送成功3条,其余不成功---因为没有加事务) 加事务得情况:(程

ActiveMQ之消息指针

消息指针(Message cursor)是activeMQ里一个非常重要的核心类,它是提供某种优化消息存储的方法.消息中间件的实现一般都是当消费者准备好消费消息的时候,它会从持久化存储中一批一批的读取消息,并发送给消费者.消息指针维护着下一批待读取消息的相关位置信息.  消息游标: 当producer发送的持久化消息到达broker之后,broker首先会把它保存在持久存储中.接下来,如果发现当前有活跃的consumer,而且这个consumer消费消息的速度能跟上producer生产消息的速度