package learn.amq; import javax.jms.Connection; import javax.jms.ConnectionFactory; import javax.jms.Destination; import javax.jms.JMSException; import javax.jms.MessageProducer; import javax.jms.Session; import javax.jms.TextMessage; import org.apache.activemq.ActiveMQConnectionFactory; /** * Date: 2016-12-27 * * @author */ public class MsgProducer { private static final String BROKER_URL = "failover://tcp://127.0.0.1:61616"; public static void main(String[] args) throws JMSException, InterruptedException { //建立链接工厂 ConnectionFactory connectionFactory = new ActiveMQConnectionFactory(BROKER_URL); //得到链接 Connection conn = connectionFactory.createConnection(); //start conn.start(); //建立Session,此方法第一个参数表示会话是否在事务中执行,第二个参数设定会话的应答模式 Session session = conn.createSession(false, Session.AUTO_ACKNOWLEDGE); //建立队列 Destination dest = session.createQueue("test-queue"); //建立消息生产者 MessageProducer producer = session.createProducer(dest); for (int i = 0; i < 100; i++) { //初始化一个mq消息 TextMessage message = session.createTextMessage("这是第 " + i + " 条消息!"); //发送消息 producer.send(message); System.out.println("send message:消息" + i); //暂停3秒 Thread.sleep(100); } //关闭mq链接 conn.close(); } } import javax.jms.Connection; import javax.jms.ConnectionFactory; import javax.jms.Destination; import javax.jms.JMSException; import javax.jms.Message; import javax.jms.MessageConsumer; import javax.jms.MessageListener; import javax.jms.Session; import javax.jms.TextMessage; import org.apache.activemq.ActiveMQConnectionFactory; /** * Date: 2016-12-27 * * @author */ public class MsgConsumer implements MessageListener { private static final String BROKER_URL="failover://tcp://127.0.0.1:61616"; public static void main(String[] args) throws JMSException { //建立链接工厂 ConnectionFactory connectionFactory=new ActiveMQConnectionFactory(BROKER_URL); //得到链接 Connection conn = connectionFactory.createConnection(); //start conn.start(); //建立Session,此方法第一个参数表示会话是否在事务中执行,第二个参数设定会话的应答模式 Session session = conn.createSession(false, Session.AUTO_ACKNOWLEDGE); //建立队列 Destination dest = session.createQueue("test-queue"); //建立消息生产者 MessageConsumer consumer = session.createConsumer(dest); //初始化MessageListener MsgConsumer msgConsumer = new MsgConsumer(); //给消费者设定监听对象 consumer.setMessageListener(msgConsumer); } /** * 消费者须要实现MessageListener接口 * 接口有一个onMessage(Message message)须要在此方法中作消息的处理 */ @Override public void onMessage(Message msg) { TextMessage txtMessage = (TextMessage)msg; try { System.out.println("get message:" + txtMessage.getText()); } catch (JMSException e) { e.printStackTrace(); } } }
依赖amq-core.jar java
compile 'org.apache.activemq:activemq-core:5.7.0'
先启动amq, 浏览器输入http://localhost:8161/ 进入amq控制台。点击Manage ActiveMQ broker,admin/admin 进入管理后台apache
先运行生产者,关停,在运行消费者,观察pending messages, messages enqueued, message dequeued的数量变化。浏览器
注意:每次重启amq,messages enqueued, message dequeued会置0,pending messages表示待消费的消息(包括上次关闭前未消费的消息)。messages enqueued, message dequeued表示本次开启amq后的入队(生产)、出队(消费)消息。numbers of consumers表示消费者数量。开启消费者时为1, 关闭消费者时为0session