官网地址:http://activemq.apache.org/
参考文章:http://my.oschina.net/nk2011/blog/366395
JMS支持两种消息发送和接收模型。一种称为P2P(Ponit to Point)模型,即采用点对点的方式发送消息。P2P模型是基于队列的,消息生产者发送消息到队列,消息消费者从队列中接收消息,队列的存在使得消息的异步传输称为可能,P2P模型在点对点的情况下进行消息传递时采用。

另一种称为Pub/Sub(Publish/Subscribe,即发布-订阅)模型,发布-订阅模型定义了如何向一个内容节点发布和订阅消息,这个内容节点称为topic(主题)。主题可以认为是消息传递的中介,消息发布这将消息发布到某个主题,而消息订阅者则从主题订阅消息。主题使得消息的订阅者与消息的发布者互相保持独立,不需要进行接触即可保证消息的传递,发布-订阅模型在消息的一对多广播时采用。

ActiveMQ的安装
下载最新的安装包apache-activemq-5.13.2-bin.tar.gz(此包linux下的,案例也是针对linux系统进行阐述,当然ActiveMQ也有win版的,这里就不赘述了),可以去官网下载,也可以在下方留言区留下你的邮箱,博主会发给你的~
下载之后解压: tar -zvxf apache-activemq-5.13.2-bin.tar.gz
ActiveMQ目录内容有:
-
bin目录包含ActiveMQ的启动脚本
-
conf目录包含ActiveMQ的所有配置文件
-
data目录包含日志文件和持久性消息数据
-
example: ActiveMQ的示例
-
lib: ActiveMQ运行所需要的lib
-
webapps: ActiveMQ的web控制台和一些相关的demo
运行命令:activemq start(在activemq/bin下运行)
1INFO: Loading '/users/shr/apache-activemq-5.13.2//bin/env' 2INFO: Using java '/users/shr/util/JavaDir/jdk/bin/java' 3INFO: Starting - inspect logfiles specified in logging.properties and log4j.properties to get details 4INFO: pidfile created : '/users/shr/apache-activemq-5.13.2//data/activemq.pid' (pid '986')
查看activemq是否运行命令:ps -aux | grep activemq
1shr 986 1.2 9.7 1281720 201936 pts/5 Sl 19:43 0:17 /users/shr/util/JavaDir/jdk/bin/java -Xms64M -Xmx1G -Djava.util.logging.config.file=logging.properties -Djava.security.auth.login.config=/users/shr/apache-activemq-5.13.2//conf/login.config -Dcom.sun.management.jmxremote -Djava.awt.headless=true -Djava.io.tmpdir=/users/shr/apache-activemq-5.13.2//tmp -Dactivemq.classpath=/users/shr/apache-activemq-5.13.2//conf:/users/shr/apache-activemq-5.13.2//../lib/: -Dactivemq.home=/users/shr/apache-activemq-5.13.2/ -Dactivemq.base=/users/shr/apache-activemq-5.13.2/ -Dactivemq.conf=/users/shr/apache-activemq-5.13.2//conf -Dactivemq.data=/users/shr/apache-activemq-5.13.2//data -jar /users/shr/apache-activemq-5.13.2//bin/activemq.jar start 2shr 1501 0.0 0.0 5176 724 pts/5 S+ 20:06 0:00 grep activemq
关闭命令: activemq stop
1INFO: Loading '/users/shr/apache-activemq-5.13.2//bin/env' 2INFO: Using java '/users/shr/util/JavaDir/jdk/bin/java' 3INFO: Waiting at least 30 seconds for regular process termination of pid '986' : 4Java Runtime: Oracle Corporation 1.7.0_79 /users/shr/util/JavaDir/jdk1.7.0_79/jre 5 Heap sizes: current=63232k free=62218k max=932096k 6 JVM args: -Xms64M -Xmx1G -Djava.util.logging.config.file=logging.properties -Djava.security.auth.login.config=/users/shr/apache-activemq-5.13.2//conf/login.config -Dactivemq.classpath=/users/shr/apache-activemq-5.13.2//conf:/users/shr/apache-activemq-5.13.2//../lib/: -Dactivemq.home=/users/shr/apache-activemq-5.13.2/ -Dactivemq.base=/users/shr/apache-activemq-5.13.2/ -Dactivemq.conf=/users/shr/apache-activemq-5.13.2//conf -Dactivemq.data=/users/shr/apache-activemq-5.13.2//data 7Extensions classpath: 8 [/users/shr/apache-activemq-5.13.2/lib,/users/shr/apache-activemq-5.13.2/lib/camel,/users/shr/apache-activemq-5.13.2/lib/optional,/users/shr/apache-activemq-5.13.2/lib/web,/users/shr/apache-activemq-5.13.2/lib/extra] 9ACTIVEMQ_HOME: /users/shr/apache-activemq-5.13.2 10ACTIVEMQ_BASE: /users/shr/apache-activemq-5.13.2 11ACTIVEMQ_CONF: /users/shr/apache-activemq-5.13.2/conf 12ACTIVEMQ_DATA: /users/shr/apache-activemq-5.13.2/data 13Connecting to pid: 986 14..Stopping broker: localhost 15.. TERMINATED
ActiveMQ的默认服务端口为61616,这个可以在conf/activemq.xml配置文件中修改:
1<transportConnectors> 2 <!-- DOS protection, limit concurrent connections to 1000 and frame size to 100MB --> 3 <transportConnector name="openwire" uri="tcp://0.0.0.0:61616?maximumConnections=1000&wireFormat.maxFrameSize=104857600"/> 4</transportConnectors>
案例
在下载的apache-activemq-5.13.2-bin.tar.gz包中解压有一个jar包:activemq-all-5.13.2.jar,引入这个jar到你的项目中即可开始编写案例代码。
博主的activemq服务器地址为10.10.195.187,这个在下面代码中会有体现。
按照JMS的规范,我们首先需要获得一个JMS connection factory.,通过这个connection factory来创建connection.在这个基础之上我们再创建session, destination, producer和consumer。因此主要的几个步骤如下:
-
获得JMS connection factory. 通过我们提供特定环境的连接信息来构造factory。
-
利用factory构造JMS connection
-
启动connection
-
通过connection创建JMS session.
-
指定JMS destination.
-
创建JMS producer或者创建JMS message并提供destination.
-
创建JMS consumer或注册JMS message listener.
-
发送和接收JMS message.
-
关闭所有JMS资源,包括connection, session, producer, consumer等。
下面来看代码举例(P2P式)。
通过Java实现的基于ActiveMQ的请求提交:
1package com.zzh.activemq; 2 3import java.io.Serializable; 4import java.util.HashMap; 5 6import javax.jms.Connection; 7import javax.jms.ConnectionFactory; 8import javax.jms.DeliveryMode; 9import javax.jms.Destination; 10import javax.jms.MessageProducer; 11import javax.jms.ObjectMessage; 12import javax.jms.Session; 13 14import org.apache.activemq.ActiveMQConnection; 15import org.apache.activemq.ActiveMQConnectionFactory; 16 17public class RequestSubmit 18{ 19 //消息发送者 20 private MessageProducer producer; 21 //一个发送或者接受消息的线程 22 private Session session; 23 24 public void init() throws Exception 25 { 26 //ConnectionFactory连接工厂,JMS用它创建连接 27 ConnectionFactory connectionFactory = new ActiveMQConnectionFactory( 28 ActiveMQConnection.DEFAULT_USER, 29 ActiveMQConnection.DEFAULT_PASSWORD, 30 "tcp://10.10.195.187:61616"); 31 //Connection:JMS客户端到JMS Provider的连接,从构造工厂中得到连接对象 32 Connection connection = connectionFactory.createConnection(); 33 //启动 34 connection.start(); 35 //获取连接操作 36 session = connection.createSession(Boolean.TRUE, Session.AUTO_ACKNOWLEDGE); 37 Destination destinatin = session.createQueue("RequestQueue"); 38 //得到消息生成(发送)者 39 producer = session.createProducer(destinatin); 40 //设置不持久化 41 producer.setDeliveryMode(DeliveryMode.NON_PERSISTENT); 42 } 43 44 public void submit(HashMap<Serializable,Serializable> requestParam) throws Exception 45 { 46 ObjectMessage message = session.createObjectMessage(requestParam); 47 producer.send(message); 48 session.commit(); 49 } 50 51 public static void main(String[] args) throws Exception{ 52 RequestSubmit submit = new RequestSubmit(); 53 submit.init(); 54 HashMap<Serializable,Serializable> requestParam = new HashMap<Serializable,Serializable>(); 55 requestParam.put("郏高阳", "zzh"); 56 submit.submit(requestParam); 57 } 58}
创建Session时有两个非常重要的参数,第一个boolean类型的参数用来表示是否采用事务消息。如果是事务消息,对于的参数设置为true,此时消息的提交自动有comit处理,消息的回滚则自动由rollback处理。加入消息不是事务的,则对应的该参数设置为false,此时分为三种情况:
-
Session.AUTO_ACKNOWLEDGE表示Session会自动确认所接收到的消息。
-
Session.CLIENT_ACKNOWLEDGE表示由客户端程序通过调用消息的确认方法来确认所接收到的消息。
-
Session.DUPS_OK_ACKNOWLEDGE使得Session将“懒惰”地确认消息,即不会立即确认消息,这样有可能导致消息重复投递。
提供Java实现的基于ActiveMQ的请求处理:
1package com.zzh.activemq; 2 3import java.io.Serializable; 4import java.util.HashMap; 5import java.util.Map; 6 7import javax.jms.Connection; 8import javax.jms.ConnectionFactory; 9import javax.jms.Destination; 10import javax.jms.MessageConsumer; 11import javax.jms.ObjectMessage; 12import javax.jms.Session; 13 14import org.apache.activemq.ActiveMQConnection; 15import org.apache.activemq.ActiveMQConnectionFactory; 16 17public class RequestProcessor 18{ 19 public void requestHandler(HashMap<Serializable,Serializable> requestParam) throws Exception 20 { 21 System.out.println("requestHandler....."+requestParam.toString()); 22 for(Map.Entry<Serializable, Serializable> entry : requestParam.entrySet()) 23 { 24 System.out.println(entry.getKey()+":"+entry.getValue()); 25 } 26 } 27 28 public static void main(String[] args) throws Exception 29 { 30 ConnectionFactory connectionFactory = new ActiveMQConnectionFactory( 31 ActiveMQConnection.DEFAULT_USER, 32 ActiveMQConnection.DEFAULT_PASSWORD, 33 "tcp://10.10.195.187:61616"); 34 Connection connection = connectionFactory.createConnection(); 35 connection.start(); 36 Session session = connection.createSession(Boolean.FALSE, Session.AUTO_ACKNOWLEDGE); 37 Destination destination = session.createQueue("RequestQueue"); 38 //消息消费(接收)者 39 MessageConsumer consumer = session.createConsumer(destination); 40 41 RequestProcessor processor = new RequestProcessor(); 42 43 while(true) 44 { 45 ObjectMessage message = (ObjectMessage) consumer.receive(1000); 46 if(null != message) 47 { 48 System.out.println(message); 49 HashMap<Serializable,Serializable> requestParam = (HashMap<Serializable,Serializable>) message.getObject(); 50 processor.requestHandler(requestParam); 51 } 52 else 53 { 54 break; 55 } 56 } 57 } 58}
输出结果:
1ActiveMQObjectMessage {commandId = 6, responseRequired = false, messageId = ID:hidden-PC-58748-1460550507055-1:1:1:1:1, originalDestination = null, originalTransactionId = null, producerId = ID:hidden-PC-58748-1460550507055-1:1:1:1, destination = queue://RequestQueue, transactionId = TX:ID:hidden-PC-58748-1460550507055-1:1:1, expiration = 0, timestamp = 1460550507333, arrival = 0, brokerInTime = 1460550505969, brokerOutTime = 1460550509143, correlationId = null, replyTo = null, persistent = false, type = null, priority = 4, groupID = null, groupSequence = 0, targetConsumerId = null, compressed = false, userID = null, content = org.apache.activemq.util.ByteSequence@74a456bb, marshalledProperties = null, dataStructure = null, redeliveryCounter = 0, size = 0, properties = null, readOnlyProperties = true, readOnlyBody = true, droppable = false, jmsXGroupFirstForConsumer = false} 2requestHandler.....{郏高阳=zzh} 3郏高阳:zzh
可以通过页面查看队列的使用情况,在浏览器中输入http://10.10.195.187:8161/admin/queues.jsp,用户名和密码都是:admin,看到以下页面:

这个是在jetty服务器下跑的,可以修改conf/jetty.xml来修改相关jetty配置。
上面的例子是关于P2P模式的,不过有个不妥之处,就是没有资源的释放。下面举一个Pub/Sub模式的。
通过JMS创建ActiveMQ的topic,并给topic发送消息:
1import javax.jms.Connection; 2import javax.jms.ConnectionFactory; 3import javax.jms.DeliveryMode; 4import javax.jms.JMSException; 5import javax.jms.MessageProducer; 6import javax.jms.ObjectMessage; 7import javax.jms.Session; 8import javax.jms.TextMessage; 9import javax.jms.Topic; 10 11import org.apache.activemq.ActiveMQConnection; 12import org.apache.activemq.ActiveMQConnectionFactory; 13import org.apache.camel.Produce; 14 15public class TopicRequest 16{ 17 //消息发送者 18 private MessageProducer producer; 19 //一个发送或者接受消息的线程 20 private Session session; 21 //Connection:JMS客户端到JMS Provider的连接 22 private Connection connection; 23 24 public void init() throws Exception 25 { 26 //ConnectionFactory连接工厂,JMS用它创建连接 27 ConnectionFactory connectionFactory = new ActiveMQConnectionFactory( 28 ActiveMQConnection.DEFAULT_USER, 29 ActiveMQConnection.DEFAULT_PASSWORD, 30 "tcp://10.10.195.187:61616"); 31 //从构造工厂中得到连接对象 32 connection = connectionFactory.createConnection(); 33 //启动 34 connection.start(); 35 //获取连接操作 36 session = connection.createSession(Boolean.FALSE, Session.AUTO_ACKNOWLEDGE); 37 Topic topic = session.createTopic("MessageTopic"); 38 producer = session.createProducer(topic); 39 //设置不持久化 40 producer.setDeliveryMode(DeliveryMode.NON_PERSISTENT); 41 } 42 43 public void submit(String mess) throws Exception 44 { 45 TextMessage message = session.createTextMessage(); 46 message.setText(mess); 47 producer.send(message); 48 } 49 50 public void close() 51 { 52 try 53 { 54 if(session != null) 55 session.close(); 56 if(producer != null) 57 producer.close(); 58 if(connection !=null ) 59 connection.close(); 60 } 61 catch (JMSException e) 62 { 63 e.printStackTrace(); 64 } 65 } 66 67 public static void main(String[] args) throws Exception 68 { 69 TopicRequest topicRequest = new TopicRequest(); 70 topicRequest.init(); 71 topicRequest.submit("I'm first"); 72 topicRequest.close(); 73 } 74}
消息发送到对应的topic后,需要将listener注册到需要订阅的topic上,以便能够接收该topic的消息:
1import javax.jms.Connection; 2import javax.jms.ConnectionFactory; 3import javax.jms.JMSException; 4import javax.jms.Message; 5import javax.jms.MessageConsumer; 6import javax.jms.MessageListener; 7import javax.jms.Session; 8import javax.jms.TextMessage; 9import javax.jms.Topic; 10 11import org.apache.activemq.ActiveMQConnection; 12import org.apache.activemq.ActiveMQConnectionFactory; 13 14public class TopicReceive 15{ 16 private MessageConsumer consumer; 17 private Session session; 18 19 public void init() throws Exception 20 { 21 ConnectionFactory connectionFactory = new ActiveMQConnectionFactory( 22 ActiveMQConnection.DEFAULT_USER, 23 ActiveMQConnection.DEFAULT_PASSWORD, 24 "tcp://10.10.195.187:61616"); 25 Connection connection = connectionFactory.createConnection(); 26 connection.start(); 27 session = connection.createSession(Boolean.FALSE, Session.AUTO_ACKNOWLEDGE); 28 Topic topic = session.createTopic("MessageTopic"); 29 consumer = session.createConsumer(topic); 30 31 consumer.setMessageListener(new MessageListener(){ 32 @Override 33 public void onMessage(Message message) 34 { 35 TextMessage tm = (TextMessage) message; 36 System.out.println(tm); 37 try 38 { 39 System.out.println(tm.getText()); 40 } 41 catch (JMSException e) 42 { 43 e.printStackTrace(); 44 } 45 } 46 }); 47 } 48 49 public static void main(String[] args) throws Exception 50 { 51 TopicReceive receive = new TopicReceive(); 52 receive.init(); 53 } 54}
输出结果:
1ActiveMQTextMessage {commandId = 5, responseRequired = false, messageId = ID:hidden-PC-50073-1460597487065-1:1:1:1:1, originalDestination = null, originalTransactionId = null, producerId = ID:hidden-PC-50073-1460597487065-1:1:1:1, destination = topic://MessageTopic, transactionId = null, expiration = 0, timestamp = 1460597487308, arrival = 0, brokerInTime = 1460597487297, brokerOutTime = 1460597487298, correlationId = null, replyTo = null, persistent = false, type = null, priority = 4, groupID = null, groupSequence = 0, targetConsumerId = null, compressed = false, userID = null, content = org.apache.activemq.util.ByteSequence@2e4d3abf, marshalledProperties = null, dataStructure = null, redeliveryCounter = 0, size = 0, properties = null, readOnlyProperties = true, readOnlyBody = true, droppable = false, jmsXGroupFirstForConsumer = false, text = I'm first} 2I'm first