ActiveMQ官網下載地址:http://activemq.apache.org/download.htmlhtml
ActiveMQ 提供了Windows 和Linux、Unix 等幾個版本,樓主這裏選擇了Linux 版本下進行開發。linux
下載完安裝包,解壓以後的目錄:web
從它的目錄來講,仍是很簡單的: apache
進入到ActiveMQ 安裝目錄的Bin 目錄,linux 下輸入 ./activemq start 啓動activeMQ 服務。後端
輸入命令以後,會提示咱們建立了一個進程IP 號,這時候說明服務已經成功啓動了。瀏覽器
ActiveMQ默認啓動時,啓動了內置的jetty服務器,提供一個用於監控ActiveMQ的admin應用。
admin:http://127.0.0.1:8161/admin/tomcat
咱們在瀏覽器打開連接以後輸入帳號密碼(這裏和tomcat 服務器相似)安全
默認帳號:admin服務器
密碼:adminsession
到這裏爲止,ActiveMQ 服務端就啓動完畢了。
ActiveMQ 在linux 下的終止命令是 ./activemq stop
項目目錄結構:
上述在官網下載ActiveMq 的時候,咱們能夠在目錄下看到一個jar包:
這個jar 包就是咱們須要在項目中進行開發中使用到的相關依賴。
public class Producter { //ActiveMq 的默認用戶名 private static final String USERNAME = ActiveMQConnection.DEFAULT_USER; //ActiveMq 的默認登陸密碼 private static final String PASSWORD = ActiveMQConnection.DEFAULT_PASSWORD; //ActiveMQ 的連接地址 private static final String BROKEN_URL = ActiveMQConnection.DEFAULT_BROKER_URL; AtomicInteger count = new AtomicInteger(0); //連接工廠 ConnectionFactory connectionFactory; //連接對象 Connection connection; //事務管理 Session session; ThreadLocal<MessageProducer> threadLocal = new ThreadLocal<>(); public void init(){ try { //建立一個連接工廠 connectionFactory = new ActiveMQConnectionFactory(USERNAME,PASSWORD,BROKEN_URL); //從工廠中建立一個連接 connection = connectionFactory.createConnection(); //開啓連接 connection.start(); //建立一個事務(這裏經過參數能夠設置事務的級別) session = connection.createSession(true,Session.SESSION_TRANSACTED); } catch (JMSException e) { e.printStackTrace(); } } public void sendMessage(String disname){ try { //建立一個消息隊列 Queue queue = session.createQueue(disname); //消息生產者 MessageProducer messageProducer = null; if(threadLocal.get()!=null){ messageProducer = threadLocal.get(); }else{ messageProducer = session.createProducer(queue); threadLocal.set(messageProducer); } while(true){ Thread.sleep(1000); int num = count.getAndIncrement(); //建立一條消息 TextMessage msg = session.createTextMessage(Thread.currentThread().getName()+ "productor:我是大帥哥,我如今正在生產東西!,count:"+num); System.out.println(Thread.currentThread().getName()+ "productor:我是大帥哥,我如今正在生產東西!,count:"+num); //發送消息 messageProducer.send(msg); //提交事務 session.commit(); } } catch (JMSException e) { e.printStackTrace(); } catch (InterruptedException e) { e.printStackTrace(); } } }
public class Comsumer { private static final String USERNAME = ActiveMQConnection.DEFAULT_USER; private static final String PASSWORD = ActiveMQConnection.DEFAULT_PASSWORD; private static final String BROKEN_URL = ActiveMQConnection.DEFAULT_BROKER_URL; ConnectionFactory connectionFactory; Connection connection; Session session; ThreadLocal<MessageConsumer> threadLocal = new ThreadLocal<>(); AtomicInteger count = new AtomicInteger(); public void init(){ try { connectionFactory = new ActiveMQConnectionFactory(USERNAME,PASSWORD,BROKEN_URL); connection = connectionFactory.createConnection(); connection.start(); session = connection.createSession(false,Session.AUTO_ACKNOWLEDGE); } catch (JMSException e) { e.printStackTrace(); } } public void getMessage(String disname){ try { Queue queue = session.createQueue(disname); MessageConsumer consumer = null; if(threadLocal.get()!=null){ consumer = threadLocal.get(); }else{ consumer = session.createConsumer(queue); threadLocal.set(consumer); } while(true){ Thread.sleep(1000); TextMessage msg = (TextMessage) consumer.receive(); if(msg!=null) { msg.acknowledge(); System.out.println(Thread.currentThread().getName()+": Consumer:我是消費者,我正在消費Msg"+msg.getText()+"--->"+count.getAndIncrement()); }else { break; } } } catch (JMSException e) { e.printStackTrace(); } catch (InterruptedException e) { e.printStackTrace(); } } }
public class TestMq { public static void main(String[] args){ Producter producter = new Producter(); producter.init(); TestMq testMq = new TestMq(); try { Thread.sleep(1000); } catch (InterruptedException e) { e.printStackTrace(); } //Thread 1 new Thread(testMq.new ProductorMq(producter)).start(); //Thread 2 new Thread(testMq.new ProductorMq(producter)).start(); //Thread 3 new Thread(testMq.new ProductorMq(producter)).start(); //Thread 4 new Thread(testMq.new ProductorMq(producter)).start(); //Thread 5 new Thread(testMq.new ProductorMq(producter)).start(); } private class ProductorMq implements Runnable{ Producter producter; public ProductorMq(Producter producter){ this.producter = producter; } @Override public void run() { while(true){ try { producter.sendMessage("Jaycekon-MQ"); Thread.sleep(10000); } catch (InterruptedException e) { e.printStackTrace(); } } } } }
運行結果:
INFO | Successfully connected to tcp://localhost:61616 Thread-6productor:我是大帥哥,我如今正在生產東西!,count:0 Thread-4productor:我是大帥哥,我如今正在生產東西!,count:1 Thread-2productor:我是大帥哥,我如今正在生產東西!,count:3 Thread-5productor:我是大帥哥,我如今正在生產東西!,count:2 Thread-3productor:我是大帥哥,我如今正在生產東西!,count:4 Thread-6productor:我是大帥哥,我如今正在生產東西!,count:5 Thread-3productor:我是大帥哥,我如今正在生產東西!,count:6 Thread-5productor:我是大帥哥,我如今正在生產東西!,count:7 Thread-2productor:我是大帥哥,我如今正在生產東西!,count:8 Thread-4productor:我是大帥哥,我如今正在生產東西!,count:9 Thread-6productor:我是大帥哥,我如今正在生產東西!,count:10 Thread-3productor:我是大帥哥,我如今正在生產東西!,count:11 Thread-5productor:我是大帥哥,我如今正在生產東西!,count:12 Thread-2productor:我是大帥哥,我如今正在生產東西!,count:13 Thread-4productor:我是大帥哥,我如今正在生產東西!,count:14 Thread-6productor:我是大帥哥,我如今正在生產東西!,count:15 Thread-3productor:我是大帥哥,我如今正在生產東西!,count:16 Thread-5productor:我是大帥哥,我如今正在生產東西!,count:17 Thread-2productor:我是大帥哥,我如今正在生產東西!,count:18 Thread-4productor:我是大帥哥,我如今正在生產東西!,count:19
public class TestConsumer { public static void main(String[] args){ Comsumer comsumer = new Comsumer(); comsumer.init(); TestConsumer testConsumer = new TestConsumer(); new Thread(testConsumer.new ConsumerMq(comsumer)).start(); new Thread(testConsumer.new ConsumerMq(comsumer)).start(); new Thread(testConsumer.new ConsumerMq(comsumer)).start(); new Thread(testConsumer.new ConsumerMq(comsumer)).start(); } private class ConsumerMq implements Runnable{ Comsumer comsumer; public ConsumerMq(Comsumer comsumer){ this.comsumer = comsumer; } @Override public void run() { while(true){ try { comsumer.getMessage("Jaycekon-MQ"); Thread.sleep(10000); } catch (InterruptedException e) { e.printStackTrace(); } } } } }
運行結果:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
|
INFO | Successfully connected to tcp:
//localhost:61616
Thread-2: Consumer:我是消費者,我正在消費MsgThread-5productor:我是大帥哥,我如今正在生產東西!,count:4--->0
Thread-3: Consumer:我是消費者,我正在消費MsgThread-4productor:我是大帥哥,我如今正在生產東西!,count:36--->1
Thread-4: Consumer:我是消費者,我正在消費MsgThread-3productor:我是大帥哥,我如今正在生產東西!,count:38--->2
Thread-5: Consumer:我是消費者,我正在消費MsgThread-6productor:我是大帥哥,我如今正在生產東西!,count:37--->3
Thread-2: Consumer:我是消費者,我正在消費MsgThread-6productor:我是大帥哥,我如今正在生產東西!,count:2--->4
Thread-3: Consumer:我是消費者,我正在消費MsgThread-5productor:我是大帥哥,我如今正在生產東西!,count:40--->5
Thread-4: Consumer:我是消費者,我正在消費MsgThread-6productor:我是大帥哥,我如今正在生產東西!,count:42--->6
Thread-5: Consumer:我是消費者,我正在消費MsgThread-4productor:我是大帥哥,我如今正在生產東西!,count:41--->7
Thread-2: Consumer:我是消費者,我正在消費MsgThread-3productor:我是大帥哥,我如今正在生產東西!,count:1--->8
Thread-3: Consumer:我是消費者,我正在消費MsgThread-2productor:我是大帥哥,我如今正在生產東西!,count:44--->9
Thread-4: Consumer:我是消費者,我正在消費MsgThread-4productor:我是大帥哥,我如今正在生產東西!,count:46--->10
Thread-5: Consumer:我是消費者,我正在消費MsgThread-5productor:我是大帥哥,我如今正在生產東西!,count:45--->11
Thread-2: Consumer:我是消費者,我正在消費MsgThread-2productor:我是大帥哥,我如今正在生產東西!,count:3--->12
Thread-3: Consumer:我是消費者,我正在消費MsgThread-3productor:我是大帥哥,我如今正在生產東西!,count:48--->13
Thread-4: Consumer:我是消費者,我正在消費MsgThread-5productor:我是大帥哥,我如今正在生產東西!,count:50--->14
Thread-5: Consumer:我是消費者,我正在消費MsgThread-2productor:我是大帥哥,我如今正在生產東西!,count:49--->15
Thread-4: Consumer:我是消費者,我正在消費MsgThread-2productor:我是大帥哥,我如今正在生產東西!,count:54--->16
Thread-2: Consumer:我是消費者,我正在消費MsgThread-5productor:我是大帥哥,我如今正在生產東西!,count:6--->17
Thread-3: Consumer:我是消費者,我正在消費MsgThread-6productor:我是大帥哥,我如今正在生產東西!,count:52--->18
Thread-5: Consumer:我是消費者,我正在消費MsgThread-3productor:我是大帥哥,我如今正在生產東西!,count:53--->19
Thread-4: Consumer:我是消費者,我正在消費MsgThread-3productor:我是大帥哥,我如今正在生產東西!,count:58--->20
|
查看運行結果,咱們能夠作ActiveMQ 服務端:http://127.0.0.1:8161/admin/ 裏面的Queues 中查看咱們生產的消息。