用最简单的方式理解和使用ActivityMQ(基础入门)

用最简单的方式理解和使用ActivityMQ(基础入门)
强烈推介IDEA2021.1.3破解激活,IntelliJ IDEA 注册码,2021.1.3IDEA 激活码  

大家好,我是架构君,一个会写代码吟诗的架构师。今天说一说用最简单的方式理解和使用ActivityMQ(基础入门),希望能够帮助大家进步!!!

基于JMS框架的消息中间件:ActivityMQ
优点:异步(无需等待)

下面分六步简单介绍
先简单概括ActivityMQ消息中间件
一、JMS基本概念
二、消息模式
三、其他消息中间件
四、Activity应用
五、整合springBoot使用消息中间件
六、其他注意点

正常情况下消息传递是通过请求响应现象完成的,但这种行为是同步的,容易产生阻塞,或者请求超时(出现重复提交的情况——通常情况下可以用token令牌解决)。
在这里插入图片描述

同步的前提下:
A项目:
如果出现B项目未响应 1、可以使用补偿机制(自动尝试连接3次),不行就会将信息放在日志中(补偿表)2、定时健康检查(哪些接口未更新导致无法连接)3、定时使用job(但这个不是实时的)。
B项目:
会出现幂等问题(多次重复响应请求)

在总结上面问题之后,就出现了我们的消息中间件ActivityMQ异步地传递消息:

在这里插入图片描述
(下图知识点很重要,请详细阅读)
在这里插入图片描述

在这里插入图片描述

总结:基于JMS的ActivityMQ框架解决了消息异步的问题,他有两种模式一种是点对点通信,一种是订阅模式,两者区别是前者生产者发布消息之后,消费者每次使用消息就删除消息,后者消费者使用消息之后消息不会被删除,下面为大家详细讲解JMS消息传递的相关概念~

一:JMS基本概念

  1.  JMS的目标
    
     为企业级的应用提供一种智能的消息系统,JMS定义了一整套的企业级的消息概念与工具,
    
     尽可能最小化的Java语言概念去构建最大化企业消息应用。统一已经存在的企业级消息系
    
     统功能。
    
  2.  JMS提供者
    
     JMS提供者是指那些完全完成JMS功能与管理功能的JMS消息厂商,理论上JMS提供者完成 
    
     JMS消息产品必须是100%的纯Java语言实现,可以运行在跨平台的架构与操作系统上,当前
    
     一些JMS厂商包括IBM,Oracle, JBoss社区 (JBoss Community), Apache 社区(ApacheCommunity)。
    
  3.  JMS应用程序, 一个完整的JMS应用应该实现以下功能:
    
      JMS 客户端 – Java语言开发的接受与发送消息的程序
    
      非JMS客户端 – 基于消息系统的本地API实现而不是JMS
    
      消息 – 应用程序用来相互交流信息的载体
    
       被管理对象–预先配置的JMS对象,JMS管理员创建,被客户端运用。如链接工厂,主题等
    
      JMS提供者–完成JMS功能与管理功能的消息系统
    

二:JMS的消息模式

2.1 点对点的消息模式(Point to Point Messaging)

2.1.1 P2P模式图

在这里插入图片描述

2.1.2 涉及到的概念

  • 1.消息队列(Queue)
  • 2.发送者(Sender)
  • 3.接收者(Receiver)
  • 4.每个消息都被发送到一个特定的队列,接收者从队列中获取消息。队列保留着消息,直到他们被消费或超时。

2.1.3 P2P的特点

  • 1.每个消息只有一个消费者(Consumer)(即一旦被消费,消息就不再在消息队列中)
  • 2.发送者和接收者之间在时间上没有依赖性,也就是说当发送者发送了消息之后,不管接收者有没有正在运行,它不会影响到消息被发送到队列
  • 3.接收者在成功接收消息之后需向队列应答成功
  • 4.如果消息未被持久化,那么在MQ关闭的时候队列中的消息缓存就会被清空
    如果你希望发送的每个消息都应该被成功处理的话,那么你需要P2P模式。

2.1.4 应用场景

A用户与B用户发送消息

2.2 发布订阅模式(publish – subscribe Mode)

2.2.1 Pub/Sub模式图

在这里插入图片描述

2.2.2涉及到的概念

  • 主题(Topic)
  • 发布者(Publisher)
  • 订阅者(Subscriber)
    客户端将消息发送到主题。多个发布者将消息发送到Topic,系统将这些消息传递给多个订阅者。

2.2.3 Pub/Sub的特点

每个消息可以有多个消费者
发布者和订阅者之间有时间上的依赖性。针对某个主题(Topic)的订阅者,它必须创建一个订阅者之后,才能消费发布者的消息,而且为了消费消息,订阅者必须保持运行的状态。
为了缓和这样严格的时间相关性,JMS允许订阅者创建一个可持久化的订阅。这样,即使订阅者没有被激活(运行),它也能接收到发布者的消息。
如果你希望发送的消息可以不被做任何处理、或者被一个消息者处理、或者可以被多个消费者处理的话,那么可以采用Pub/Sub模型
消息的消费
在JMS中,消息的产生和消息是异步的。对于消费来说,JMS的消息者可以通过两种方式来消费消息。
○ 同步
订阅者或接收者调用receive方法来接收消息,receive方法在能够接收到消息之前(或超时之前)将一直阻塞
○ 异步
订阅者或接收者可以注册为一个消息监听器。当消息到达之后,系统自动调用监听器的onMessage方法。

2.2.4 应用场景:

用户注册、订单修改库存、日志存储
画图演示
在这里插入图片描述

在这里插入图片描述

在这里插入图片描述

2.3 消息可靠机制

ActiveMQ消息签收机制:
客戶端成功接收一条消息的标志是一条消息被签收,成功应答。
消息的签收情形分两种:

  • 1、带事务的session
    如果session带有事务,并且事务成功提交,则消息被自动签收。如果事务回滚,则消息会被再次传送。

  • 2、不带事务的session
    不带事务的session的签收方式,取决于session的配置。
    Activemq支持一下三種模式:
    Session.AUTO_ACKNOWLEDGE 消息自动签收

    Session.CLIENT_ACKNOWLEDGE 客戶端调用acknowledge方法手动签收
    即 textMessage.acknowledge();//表示手动签收完成

    Session.DUPS_OK_ACKNOWLEDGE 不是必须签收,消息可能会重复发送。在第二次重新传送消息的时候,消息
    只有在被确认之后,才认为已经被成功地消费了。消息的成功消费通常包含三个阶段:客户接收消息、客户处理消息和消息被确认。 在事务性会话中,当一个事务被提交的时候,确认自动发生。在非事务性会话中,消息何时被确认取决于创建会话时的应答模式(acknowledgement mode)。

场景1


生产者不开启session,客户端必须有手动签收模式
Session session = createConnection.createSession(Boolean.FALSE, Session.CLIENT_ACKNOWLEDGE);
消费者不开启session,客户端必须有手动签收模式
TextMessage textMessage = (TextMessage) createConsumer.receive();
//接受消息
textMessage.acknowledge();

场景2

生产者不开启session,客户端自动签收模式 
Session session = createConnection.createSession(Boolean.FALSE, Session.AUTO_ACKNOWLEDGE);
消费者不开启session,自动签收消息
Session session = createConnection.createSession(Boolean.FALSE, Session.AUTO_ACKNOWLEDGE);

场景3

事物消息 生产者以事物形式,必须要将消息提交事物,才可以提交到队列中。
Session session = createConnection.createSession(Boolean.TRUE, Session.AUTO_ACKNOWLEDGE);
session.commit();
消费者
Session session = createConnection.createSession(Boolean.TRUE, Session.AUTO_ACKNOWLEDGE);
session.commit();

总结消息 可靠机制为:
1.自动签收:(没有实物补偿机制)万一数据出错,就没有后续手段了
2.事务签收:1.生产者完成消息,必须提交给队列。2.消费者没有提交实物,默认没有进行消费,默认进行重试(即生产者和消费者在完成操作之后都要commit才算完成事务)
3.手动签收:消费者没有签收,默认没有进行消费(只在消费者断进行限制,即只在消费者端需要确认签收)

三、其他各类MQ产品

在介绍怎么使用之前为大家简单介绍下其他产品

RabbitMQ

是使用Erlang编写的一个开源的消息队列,本身支持很多的协议:AMQP,XMPP, SMTP, STOMP,也正是如此,使的它变的非常重量级,更适合于企业级的开发。同时实现了一个经纪人(Broker)构架,这意味着消息在发送给客户端时先在中心队列排队。对路由(Routing),负载均衡(Load balance)或者数据持久化都有很好的支持。

Redis

是一个Key-Value的NoSQL数据库,开发维护很活跃,虽然它是一个Key-Value数据库存储系统,但它本身支持MQ功能,所以完全可以当做一个轻量级的队列服务来使用。对于RabbitMQ和Redis的入队和出队操作,各执行100万次,每10万次记录一次执行时间。测试数据分为128Bytes、512Bytes、1K和10K四个不同大小的数据。实验表明:入队时,当数据比较小时Redis的性能要高于RabbitMQ,而如果数据大小超过了10K,Redis则慢的无法忍受;出队时,无论数据大小,Redis都表现出非常好的性能,而RabbitMQ的出队性能则远低于Redis。

ZeroMQ

号称最快的消息队列系统,尤其针对大吞吐量的需求场景。ZMQ能够实现RabbitMQ不擅长的高级/复杂的队列,但是开发人员需要自己组合多种技术框架,技术上的复杂度是对这MQ能够应用成功的挑战。ZeroMQ具有一个独特的非中间件的模式,你不需要安装和运行一个消息服务器或中间件,因为你的应用程序将扮演了这个服务角色。你只需要简单的引用ZeroMQ程序库,可以使用NuGet安装,然后你就可以愉快的在应用程序之间发送消息了。但是ZeroMQ仅提供非持久性的队列,也就是说如果down机,数据将会丢失。其中,Twitter的Storm中使用ZeroMQ作为数据流的传输。

ActiveMQ

是Apache下的一个子项目。 类似于ZeroMQ,它能够以代理人和点对点的技术实现队列。同时类似于RabbitMQ,它少量代码就可以高效地实现高级应用场景。RabbitMQ、ZeroMQ、ActiveMQ均支持常用的多种语言客户端 C++、Java、.Net,、Python、 Php、 Ruby等。

Jafka/Kafka

Kafka是Apache下的一个子项目,是一个高性能跨语言分布式Publish/Subscribe消息队列系统,而Jafka是在Kafka之上孵化而来的,即Kafka的一个升级版。具有以下特性:快速持久化,可以在O(1)的系统开销下进行消息持久化;高吞吐,在一台普通的服务器上既可以达到10W/s的吞吐速率;完全的分布式系统,Broker、Producer、Consumer都原生自动支持分布式,自动实现复杂均衡;支持Hadoop数据并行加载,对于像Hadoop的一样的日志数据和离线分析系统,但又要求实时处理的限制,这是一个可行的解决方案。Kafka通过Hadoop的并行加载机制来统一了在线和离线的消息处理,这一点也是本课题所研究系统所看重的。Apache Kafka相对于ActiveMQ是一个非常轻量级的消息系统,除了性能非常好之外,还是一个工作良好的分布式系统。

其他一些队列列表HornetQ、Apache Qpid、Sparrow、Starling、Kestrel、Beanstalkd、Amazon SQS就不再一一分析。

以上总结起来就是rocketmq具备事务性,适合分布式系统,kafka具备日志功能适合大数据方面使用,activity只具备消息队列是消息中间件中比较容易使用的一个

四、 ActiveMQ使用

4.1 、window或者linux下 ActiveMQ安装

ActiveMQ部署其实很简单,和所有Java一样,要跑java程序就必须先安装JDK并配置好环境变量,这个很简单。
然后解压下载的apache-activemq-5.10-20140603.133406-78-bin.zip压缩包到一个目录,得到解压后的目录结构如下图:
在这里插入图片描述

在windows系统就进入bin双击启动activitymq,在linux系统就进入bin输入: ./activitymq start 启动 即可

启动系统之后 默认端口为127.0.0.1:8161,账号密码都为admin,即可查看activity页面。如果是虚拟机下的linux系统则输入虚拟机IP+8161即可

4.2实现点对点通信

  • 引入pom文件依赖
<dependencies>
		<dependency>
			<groupId>org.apache.activemq</groupId>
			<artifactId>activemq-core</artifactId>
			<version>5.7.0</version>
		</dependency>
	</dependencies>
  • 生产者
    主要步骤为:
  • 获取MQ工程
  • 创建连接
  • 启动连接
  • 创建会话工厂
  • 创建队列
  • 存储消息

启动生产者之后可以在127.0.0.1:8161查看消息队里

public class Producter {
	public static void main(String[] args) throws JMSException {
		// ConnectionFactory :连接工厂,JMS 用它创建连接
		ConnectionFactory connectionFactory = new ActiveMQConnectionFactory(ActiveMQConnection.DEFAULT_USER,
				ActiveMQConnection.DEFAULT_PASSWORD, "tcp://127.0.0.1:61616");
		// JMS 客户端到JMS Provider 的连接
		Connection connection = connectionFactory.createConnection();
		connection.start();
		// Session: 一个发送或接收消息的线程
		Session session = connection.createSession(Boolean.falst, Session.AUTO_ACKNOWLEDGE);
		// Destination :消息的目的地;消息发送给谁.
		// 获取session注意参数值my-queue是Query的名字
		Destination destination = session.createQueue("my-queue");
		// MessageProducer:消息生产者
		MessageProducer producer = session.createProducer(destination);
		// 设置不持久化(就是设置为点对点通信,如果持久化就是发布订阅模式)
		producer.setDeliveryMode(DeliveryMode.NON_PERSISTENT);
		// 发送一条消息
		for (int i = 1; i <= 5; i++) {
			sendMsg(session, producer, i);
		}
		connection.close();
	}
	/**
	 * 在指定的会话上,通过指定的消息生产者发出一条消息
	 * 
	 * @param session
	 *            消息会话
	 * @param producer
	 *            消息生产者
	 */
	public static void sendMsg(Session session, MessageProducer producer, int i) throws JMSException {
		// 创建一条文本消息
		TextMessage message = session.createTextMessage("Hello ActiveMQ!" + i);
		// 通过消息生产者发出消息
		producer.send(message);
	}
}
  • 消费者
public class JmsReceiver {
	public static void main(String[] args) throws JMSException {
		// ConnectionFactory :连接工厂,JMS 用它创建连接
		ConnectionFactory connectionFactory = new ActiveMQConnectionFactory(ActiveMQConnection.DEFAULT_USER,
				ActiveMQConnection.DEFAULT_PASSWORD, "tcp://127.0.0.1:61616");
		// JMS 客户端到JMS Provider 的连接
		Connection connection = connectionFactory.createConnection();
		connection.start();
		// Session: 一个发送或接收消息的线程
		Session session = connection.createSession(Boolean.TRUE, Session.AUTO_ACKNOWLEDGE);
		// Destination :消息的目的地;消息发送给谁.
		// 获取session注意参数值xingbo.xu-queue是一个服务器的queue,须在在ActiveMq的console配置
		Destination destination = session.createQueue("my-queue");
		// 消费者,消息接收者
		MessageConsumer consumer = session.createConsumer(destination);
		while (true) {
			TextMessage message = (TextMessage) consumer.receive();
			if (null != message) {
				System.out.println("收到消息:" + message.getText());
			} else
				break;
		}
		session.close();
		connection.close();
	}
}

在这里插入图片描述

Number Of Consumers 消费者 这个是消费者端的消费者数量

Number Of Pending Messages 等待消费的消息 这个是当前未出队列的数量。可以理解为总接收数-总出队列数
Messages Enqueued 进入队列的消息 进入队列的总数量,包括出队列的。 这个数量只增不减
Messages Dequeued 出了队列的消息 可以理解为是消费这消费掉的数量

4.4 、发布订阅

  • 生产者:
public class TOPSend {

	private static String BROKERURL = "tcp://127.0.0.1:61616";
	private static String TOPIC = "my-topic";

	public static void main(String[] args) throws JMSException {
		start();
	}

	static public void start() throws JMSException {
		System.out.println("生产者已经启动....");
		// 创建ActiveMQConnectionFactory 会话工厂
		ActiveMQConnectionFactory activeMQConnectionFactory = new ActiveMQConnectionFactory(
				ActiveMQConnection.DEFAULT_USER, ActiveMQConnection.DEFAULT_PASSWORD, BROKERURL);
		Connection connection = activeMQConnectionFactory.createConnection();
		// 启动JMS 连接
		connection.start();
		Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
		MessageProducer producer = session.createProducer(null);
		producer.setDeliveryMode(DeliveryMode.PERSISTENT);
		send(producer, session);
		System.out.println("发送成功!");
		connection.close();
	}

	static public void send(MessageProducer producer, Session session) throws JMSException {
		for (int i = 1; i <= 5; i++) {
			System.out.println("我是消息" + i);
			TextMessage textMessage = session.createTextMessage("我是消息" + i);
			Destination destination = session.createTopic(TOPIC);
			producer.send(destination, textMessage);
		}
	}

}

  • 消费者:
public class TopReceiver {
	private static String BROKERURL = "tcp://127.0.0.1:61616";
	private static String TOPIC = "my-topic";

	public static void main(String[] args) throws JMSException {
		start();
	}

	static public void start() throws JMSException {
		System.out.println("消费点启动...");
		// 创建ActiveMQConnectionFactory 会话工厂
		ActiveMQConnectionFactory activeMQConnectionFactory = new ActiveMQConnectionFactory(
				ActiveMQConnection.DEFAULT_USER, ActiveMQConnection.DEFAULT_PASSWORD, BROKERURL);
		Connection connection = activeMQConnectionFactory.createConnection();
		// 启动JMS 连接
		connection.start();
		// 不开消息启事物,消息主要发送消费者,则表示消息已经签收
		Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
		// 创建一个队列
		Topic topic = session.createTopic(TOPIC);
		MessageConsumer consumer = session.createConsumer(topic);
		// consumer.setMessageListener(new MsgListener());
		while (true) {
			TextMessage textMessage = (TextMessage) consumer.receive();
			if (textMessage != null) {
				System.out.println("接受到消息:" + textMessage.getText());
				// textMessage.acknowledge();// 手动签收
				// session.commit();
			} else {
				break;
			}
		}
		connection.close();
	}

}

4.5 、SpringBoot整合ActiveMQ

  • 生产者:

4.5.1 引入 maven依赖

<parent>
		<groupId>org.springframework.boot</groupId>
		<artifactId>spring-boot-starter-parent</artifactId>
		<version>1.5.4.RELEASE</version>
		<relativePath/> <!-- lookup parent from repository -->
	</parent>
	<properties>
		<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
		<project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
		<java.version>1.8</java.version>
	</properties>
	<dependencies>
		<dependency>
			<groupId>org.springframework.boot</groupId>
			<artifactId>spring-boot-starter</artifactId>
		</dependency>
		<!-- spring boot web支持:mvc,aop... -->
		<dependency>
			<groupId>org.springframework.boot</groupId>
			<artifactId>spring-boot-starter-web</artifactId>
		</dependency>
		<dependency>
			<groupId>org.springframework.boot</groupId>
			<artifactId>spring-boot-starter-test</artifactId>
			<scope>test</scope>
		</dependency>
		<dependency>
			<groupId>org.springframework.boot</groupId>
			<artifactId>spring-boot-starter-activemq</artifactId>
		</dependency>
	</dependencies>
	<build>
		<plugins>
			<plugin>
				<groupId>org.springframework.boot</groupId>
				<artifactId>spring-boot-maven-plugin</artifactId>
			</plugin>
		</plugins>
	</build>

4.5.2 引入 application.yml配置

spring:
  activemq:
    broker-url: tcp://127.0.0.1:61616
    user: admin
    password: admin
queue: springboot-queue
server:
  port: 8080

4.5.3 创建QueueConfig

@Configuration
public class QueueConfig {
	@Value("${queue}")
	private String queue;

	@Bean
	public Queue logQueue() {
		return new ActiveMQQueue(queue);
	}
}

4.5.4 创建Producer

@Component
@EnableScheduling
public class Producer {
	@Autowired
	private JmsMessagingTemplate jmsMessagingTemplate;
	@Autowired
	private Queue queue;

	@Scheduled(fixedDelay = 5000)
	public void send() {
		jmsMessagingTemplate.convertAndSend(queue, "测试消息队列" + System.currentTimeMillis());
	}
}

4.5.5 启动

@SpringBootApplication
@EnableScheduling
public class App {
	public static void main(String[] args) {
		SpringApplication.run(App.class, args);
	}
}
  • 消费者:

4.5.1 引入 maven依赖

<parent>
		<groupId>org.springframework.boot</groupId>
		<artifactId>spring-boot-starter-parent</artifactId>
		<version>1.5.4.RELEASE</version>
		<relativePath/> <!-- lookup parent from repository -->
	</parent>
	<properties>
		<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
		<project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
		<java.version>1.8</java.version>
	</properties>
	<dependencies>
		<dependency>
			<groupId>org.springframework.boot</groupId>
			<artifactId>spring-boot-starter</artifactId>
		</dependency>
		<!-- spring boot web支持:mvc,aop... -->
		<dependency>
			<groupId>org.springframework.boot</groupId>
			<artifactId>spring-boot-starter-web</artifactId>
		</dependency>
		<dependency>
			<groupId>org.springframework.boot</groupId>
			<artifactId>spring-boot-starter-test</artifactId>
			<scope>test</scope>
		</dependency>
		<dependency>
			<groupId>org.springframework.boot</groupId>
			<artifactId>spring-boot-starter-activemq</artifactId>
		</dependency>
	</dependencies>
	<build>
		<plugins>
			<plugin>
				<groupId>org.springframework.boot</groupId>
				<artifactId>spring-boot-maven-plugin</artifactId>
			</plugin>
		</plugins>
	</build>

4.5.2 引入 YML配置

spring:
  activemq:
    broker-url: tcp://127.0.0.1:61616
    user: admin
    password: admin
queue: springboot-queue
server:
  port: 8081

4.5.3 创建Consumer

@Component
public class Consumer {

	@JmsListener(destination = "${queue}")
	public void receive(String msg) {
		System.out.println("监听器收到msg:" + msg);
	}

}

4.5.4 启动

@SpringBootApplication
public class App {
	public static void main(String[] args) {
		SpringApplication.run(App.class, args);
	}
}

六、其他注意点

使用消息中间注意事项

1.消费者代码不要抛出异常,否则activqmq默认有重试机制。
2.如果代码发生异常,需要发布版本才可以解决的问题,不要使用重试机制,采用日志记录方式,定时Job进行补偿。
3.如果不需要发布版本解决的问题,可以采用重试机制进行补偿。

消费者如果保证消息幂等性,不被重复消费。

产生原因:网络延迟传输中,会造成进行MQ重试中,在重试过程中,可能会造成重复消费。
解决办法:
1.使用全局MessageID 判断消费方使用同一个,解决幂等性。(textmessage.getJMSMessageID,或者getJMSTime获取时间戳)
2.使用JMS可靠消息机制(避免第三次应答,应该在前两次就完成。

判断什么情况下需要重试

  • 第三方连接异常(jdbc异常):不重试
  • 数据转换异常:重试
    总结:需要发布版本才能解决的问题不重试(就是代码出问题的)

MQ协议

消息中间件实现集群-》activityMQ集群+zookeeper进行管理
如果是消费者集群就不需要担心幂等性问题,因为MQ只有一个,他自己知道会分发给谁,若出现高并发,它会自动分批发送!

在这里插入图片描述

在这里插入图片描述

本文来源GroovRain,由架构君转载发布,观点不代表Java架构师必看的立场,转载请标明来源出处:https://javajgs.com/archives/25368

发表评论