欢迎您访问程序员文章站本站旨在为大家提供分享程序员计算机编程知识!
您现在的位置是: 首页

ActiveMQ消息中间件

程序员文章站 2022-07-01 15:32:06
...

JMS简介

        JMS(Java Messaging Service)是Java平台上有关面向消息中间件的技术规范,它便于消息系统中的Java应用程序进行消息交换,并且通过提供标准的产生、发送、接收消息的接口简化企业应用的开发。

       JMS本身只定义了一系列的接口规范,是一种与厂商无关的 API,用来访问消息收发系统。它类似于 JDBC(java Database Connectivity):这里,JDBC 是可以用来访问许多不同关系数据库的 API,而 JMS 则提供同样与厂商无关的访问方法,以访问消息收发服务。许多      厂商目前都支持 JMS,包括 IBM 的 MQSeries、BEA的 Weblogic JMS service和 Progress 的 SonicMQ,这只是几个例子。 JMS 使您能够通过消息收发服务(有时称为消息中介程序或路由器)从一个 JMS 客户机向另一个 JML 客户机发送消息。消息是 JMS 中的一种类型对象,由两部分组成:报头和消息主体。报头由路由信息以及有关该消息的元数据组成。消息主体则携带着应用程序的数据或有效负载。根据有效负载 的类型来划分,可以将消息分为几种类型,它们分别携带:简单文本 (TextMessage)、可序列化的对象 (ObjectMessage)、属性集合 (MapMessage)、字节流 (BytesMessage)、原始值流 (StreamMessage),还有无有效负载的消息 (Message)。

ActiveMQ简介

ActiveMQ 是Apache出品,最流行的,能力强劲的开源消息总线。ActiveMQ 是一个完全支持JMS1.1和J2EE 1.4规范的 JMS Provider实现,尽管JMS规范出台已经是很久的事情了,但是JMS在当今的J2EE应用中间仍然扮演着特殊的地位。

对于消息的传递有两种类型:

一种是点对点的,即一个生产者和一个消费者一一对应;

ActiveMQ消息中间件

 

另一种是发布/ 订阅模式,即一个生产者产生消息并进行发送后,可以由多个消费者进行接收。

ActiveMQ消息中间件

JMS 定义了五种不同的消息正文格式,以及调用的消息类型,允许你发送并接收以一

些不同形式的数据,提供现有消息格式的一些级别的兼容性。

· StreamMessage -- Java 原始值的数据流

· MapMessage--一套名称-值对

· TextMessage--一个字符串对象

· ObjectMessage--一个序列化的 Java 对象

· BytesMessage--一个字节的数据流

ActiveMQ下载与安装

官方网站下载

ActiveMQ入门小DEMO(点对点)

点对点的模式主要建立在一个队列上面,当连接一个列队的时候,发送端不需要知道接收端是否正在接收,可以直接向ActiveMQ发送消息,发送的消息,将会先进入队列中,如果有接收端在监听,则会发向接收端,如果没有接收端接收,则会保存在activemq服务器,直到接收端接收消息,点对点的消息模式可以有多个发送端,多个接收端,但是一条消息,只会被一个接收端给接收到,哪个接收端先连上ActiveMQ,则会先接收到,而后来的接收端则接收不到那条消息。

消息生产者

(1)创建工程activemqDemo ,引入依赖

	<dependency>
	  		<groupId>org.apache.activemq</groupId>
	  		<artifactId>activemq-all</artifactId>
	  		<version>5.11.2</version>
	</dependency>

(2)创建类QueueProducer

public class QueueProducer {
	public static void main(String[] args) throws JMSException {
		//1.创建连接工厂
		ActiveMQConnectionFactory connectionFactory=new ActiveMQConnectionFactory("tcp://192.168.25.129:61616");
		//2.获取连接
		Connection connection = connectionFactory.createConnection();
		//3.启动连接
		connection.start();
		//4.获取session
		Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
		//5.创建队列对象
		Queue queue = session.createQueue("test-queue");
		//6.创建消息生产者对象
		MessageProducer producer = session.createProducer(queue);
		//7.创建消息
		TextMessage textMessage = session.createTextMessage("欢迎来到神奇的品优购世界!");
		//8.使用生成者发送消息
		producer.send(textMessage);
		//9.关闭资源
		producer.close();
		session.close();
		connection.close();
	}
}

上述代码中第4步创建session  的两个参数:

第1个参数 是否使用事务

第2个参数 消息的确认模式

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

消息消费者

public class QueueConsumer {
	public static void main(String[] args) throws JMSException, IOException {
		//1.创建连接工厂
		ActiveMQConnectionFactory connectionFactory=new ActiveMQConnectionFactory("tcp://192.168.25.129:61616");		
		//2.获取连接
		Connection connection = connectionFactory.createConnection();
		//3.启动连接
		connection.start();
		//4.获取session
		Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
		//5.创建队列对象
		Queue queue = session.createQueue("test-queue");
		//6.创建消息消费者
		MessageConsumer consumer=session.createConsumer(queue);
		//7.接受消息
		consumer.setMessageListener(new MessageListener() {			
			@Override
			public void onMessage(Message message) {
				//获取文本消息对象
				TextMessage textMessage=(TextMessage)message;				
				try {
					String text = textMessage.getText();//提取文本
					System.out.println(text);
				} catch (JMSException e) {
					// TODO Auto-generated catch block
					e.printStackTrace();
				}				
			}
		});
		//8.等待键盘输入
		System.in.read();
		//9.关闭资源		
		consumer.close();
		session.close();
		connection.close();
	}
}

ActiveMQ入门小DEMO(发布/订阅) 

消息生产者

public class TopicProducer {
	public static void main(String[] args) throws JMSException {		
		//1.创建连接工厂
		ActiveMQConnectionFactory connectionFactory=new ActiveMQConnectionFactory("tcp://192.168.25.129:61616");		
		//2.创建连接
		Connection connection = connectionFactory.createConnection();		
		//3.启动连接
		connection.start();
		//4.创建会话
		Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
		//5.创建一个订阅
		Topic topic = session.createTopic("test-topic");
		//6.创建消息生产者		
		MessageProducer producer = session.createProducer(topic);
		//7.创建消息
		TextMessage textMessage = session.createTextMessage("您的手机该续费了!");
		//8.发送消息		
		producer.send(textMessage);
		//9.关闭资源
		producer.close();
		session.close();
		connection.close();
	}
}

消息消费者 

public class TopicConsumer {
	public static void main(String[] args) throws JMSException, IOException {
		//1.创建连接工厂
		ActiveMQConnectionFactory connectionFactory=new ActiveMQConnectionFactory("tcp://192.168.25.129:61616");		
		//2.获取连接
		Connection connection = connectionFactory.createConnection();
		//3.启动连接
		connection.start();
		//4.获取session
		Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
		//5.创建主题对象
		Topic topic = session.createTopic("test-topic");
		//6.创建消息消费者
		MessageConsumer consumer=session.createConsumer(topic);
		//7.接受消息
		consumer.setMessageListener(new MessageListener() {			
			@Override
			public void onMessage(Message message) {
				//获取文本消息对象
				TextMessage textMessage=(TextMessage)message;				
				try {
					String text = textMessage.getText();//提取文本
					System.out.println(text);
				} catch (JMSException e) {
					// TODO Auto-generated catch block
					e.printStackTrace();
				}				
			}
		});
		//8.等待键盘输入
		System.in.read();
		//9.关闭资源		
		consumer.close();
		session.close();
		connection.close();
	}
}

Spring整合ActiveMQ 

点对点消息

消息生产者

  1. 创建工程activemq_spring_producer,在POM文件中引入依赖
  2. 在src/main/resources下创建spring配置文件  applicationContext-activemq-producer.xml
 <!-- 真正可以产生Connection的ConnectionFactory,由对应的 JMS服务厂商提供-->  
	<bean id="targetConnectionFactory" class="org.apache.activemq.ActiveMQConnectionFactory">  
	    <property name="brokerURL" value="tcp://192.168.25.129:61616"/>  
	</bean>          
    <!-- Spring用于管理真正的ConnectionFactory的ConnectionFactory -->  
	<bean id="connectionFactory" class="org.springframework.jms.connection.SingleConnectionFactory">  
	<!-- 目标ConnectionFactory对应真实的可以产生JMS Connection的ConnectionFactory -->  
	    <property name="targetConnectionFactory" ref="targetConnectionFactory"/>  
	</bean>  	   
    <!-- Spring提供的JMS工具类,它可以进行消息发送、接收等 -->  
	<bean id="jmsTemplate" class="org.springframework.jms.core.JmsTemplate">  
	    <!-- 这个connectionFactory对应的是我们定义的Spring提供的那个ConnectionFactory对象 -->  
	    <property name="connectionFactory" ref="connectionFactory"/>  
	</bean>      
    <!--这个是队列目的地,点对点的  文本信息-->  
	<bean id="queueTextDestination" class="org.apache.activemq.command.ActiveMQQueue">  
	    <constructor-arg value="queue_text"/>  
	</bean>    
</beans>

(3)创建消息生产者类

public class QueueProducer {

	@Autowired
	private JmsTemplate jmsTemplate;
	
	@Autowired
	private Destination  queueTextDestination;
	
	/**
	 * 发送文本消息
	 */
	public void sendTextMessage(){				
		jmsTemplate.send(queueTextDestination,new MessageCreator() {			
			public Message createMessage(Session session) throws JMSException {
				return session.createTextMessage("spring与activeMQ整合--文本消息");
			}
		});		
	}	
}

消息消费者

  1. 创建工程activemq_spring_consumer,在POM文件中引入依赖 (同上一个工程)
  2. 编写监听类
public class MyMessageListener implements MessageListener{
	public void onMessage(Message message) {
		TextMessage textMessage=(TextMessage)message;
		String text;
		try {
			text = textMessage.getText();
			System.out.println(text);
		} catch (JMSException e) {
			e.printStackTrace();
		}		
	}
}

(3)创建配置文件 applicationContext-activemq-consumer-queue.xml

 <!-- 真正可以产生Connection的ConnectionFactory,由对应的 JMS服务厂商提供-->  
	<bean id="targetConnectionFactory" class="org.apache.activemq.ActiveMQConnectionFactory">  
	    <property name="brokerURL" value="tcp://192.168.25.129:61616"/>  
	</bean>	   
    <!-- Spring用于管理真正的ConnectionFactory的ConnectionFactory -->  
	<bean id="connectionFactory" class="org.springframework.jms.connection.SingleConnectionFactory">  
	<!-- 目标ConnectionFactory对应真实的可以产生JMS Connection的ConnectionFactory -->  
	    <property name="targetConnectionFactory" ref="targetConnectionFactory"/>  
	</bean>  	
    <!--这个是队列目的地,点对点的  文本信息-->  
	<bean id="queueTextDestination" class="org.apache.activemq.command.ActiveMQQueue">  
	    <constructor-arg value="queue_text"/>  
	</bean>    
	<!-- 我的监听类 -->
	<bean id="myMessageListener" class="cn.itcast.demo.MyMessageListener"></bean>
	<!-- 消息监听容器 -->
	<bean class="org.springframework.jms.listener.DefaultMessageListenerContainer">
		<property name="connectionFactory" ref="connectionFactory" />
		<property name="destination" ref="queueTextDestination" />
		<property name="messageListener" ref="myMessageListener" />
	</bean>	
</beans>

 

 

相关标签: ActiveMQ