centos安装
`1、安装erlang 以root身份执行下面命令 yum install erlang 2、安装 rabbitmq-server 打开RabbitMQ的下载页面,http://www.rabbitmq.com/download.html ,选择对应平台的二进制发行包下载;目前使用的是CentOS ,属于与RHEL/Fedora相兼容的版本,下载针对RHEL的二进制版本(Binary)即可: 本例中RabbitMQ的版本是3.5.1,下载得到文件rabbitmq-server-3.5.1-1.noarch.rpm 命令如下: wget http://www.rabbitmq.com/releases/rabbitmq-server/v3.5.1/rabbitmq-server-3.5.1-1.noarch.rpm
安装RabbitMQ Server
rpm --import http://www.rabbitmq.com/rabbitmq-signing-key-public.asc
yum install rabbitmq-server-3.5.1-1.noarch.rpm
3、启动RabbitMQ
配置为守护进程随系统自动启动,root权限下执行:
chkconfig rabbitmq-server on
启动rabbitMQ服务
/sbin/service rabbitmq-server start 如果报如下异常:
Starting rabbitmq-server (via systemctl): Job for rabbitmq-server.service failed. See 'systemctl status rabbitmq-server.service' and 'journalctl -xn' for details. [FAILED] 尝试下面的操作: 禁用 SELinux ,修改 /etc/selinux/config SELINUX=disabled 修改后重启系统
4、安装Web管理界面插件 终端输入:
rabbitmq-plugins enable rabbitmq_management 安装成功后会显示如下内容
The following plugins have been enabled: mochiweb webmachine rabbitmq_web_dispatch amqp_client rabbitmq_management_agent rabbitmq_management Plugin configuration has changed. Restart RabbitMQ for changes to take effect. 5、登录Web管理界面 安装好插件并开启服务后,可以浏览器输入xxxx:15672,账号密码全输入guest即可登录。
rabbitmq的web管理界面无法使用guest用户登录
这里需要注意下,从3.3.1版本开始,RabbitMQ默认不允许远程ip登录,即只能使用localhost登录。如果希望远程登录,请添加用户权限,方法见我另一篇文章设置RabbitMQ远程ip登录。
rabbitmq例子
消费者,生产者都继承自该类
1public abstract class EndPoint { 2 protected Channel channel; 3 protected Connection connection; 4 protected String endPointName; 5 6 7 public EndPoint(String endPointName) throws IOException{ 8 //queue名称 9 this.endPointName = endPointName; 10 ConnectionFactory factory = new ConnectionFactory(); 11 factory.setHost("10.2.223.71"); 12 //注意这里的port 为5672 13 factory.setPort(5672); 14 factory.setUsername("oyth"); 15 factory.setPassword("oyth"); 16 connection = factory.newConnection(); 17 channel = connection.createChannel(); 18 channel.queueDeclare(endPointName, false, false, false, null); 19 } 20 21 public void close() throws IOException{ 22 this.channel.close(); 23 this.connection.close(); 24 }
生产者
1public class Producer extends EndPoint { 2 3 public Producer(String endPointName) throws IOException { 4 super(endPointName); 5 } 6 7 public void sendMessage(Serializable object) throws IOException{ 8 channel.basicPublish("",endPointName,null, SerializationUtils.serialize(object)); 9 }
消费者
1public class QueueConsumer extends EndPoint implements Runnable,Consumer { 2 3 public QueueConsumer(String endPointName) throws IOException{ 4 super(endPointName); 5 } 6 7 @Override 8 public void handleConsumeOk(String consumerTag) { 9 System.out.println("Consumer "+consumerTag +" registered"); 10 } 11 12 @Override 13 public void handleCancelOk(String s) { 14 15 } 16 17 @Override 18 public void handleCancel(String s) throws IOException { 19 20 } 21 22 @Override 23 public void handleShutdownSignal(String s, ShutdownSignalException e) { 24 25 } 26 27 @Override 28 public void handleRecoverOk(String s) { 29 30 } 31 32 @Override 33 public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties basicProperties, byte[] body) throws IOException { 34 Map map = (HashMap) SerializationUtils.deserialize(body); 35 System.out.println("Message Number "+ map.get("message number") + " received."); 36 } 37 38 /** 39 * When an object implementing interface <code>Runnable</code> is used 40 * to create a thread, starting the thread causes the object's 41 * <code>run</code> method to be called in that separately executing 42 * thread. 43 * <p/> 44 * The general contract of the method <code>run</code> is that it may 45 * take any action whatsoever. 46 * 47 * @see Thread#run() 48 */ 49 @Override 50 public void run() { 51 try { 52 //start consuming messages. Auto acknowledge messages. 53 channel.basicConsume(endPointName, true,this); 54 } catch (IOException e) { 55 e.printStackTrace(); 56 } 57 } 58}
运行代码
1public class TestMq { 2 3 public TestMq() throws Exception{ 4 QueueConsumer queueConsumer = new QueueConsumer("queue"); 5 Thread thread = new Thread(queueConsumer); 6 thread.start(); 7 8 Producer producer = new Producer("queue"); 9 for (int i =0;i<100;i++){ 10 HashMap msg = new HashMap(); 11 msg.put("message number",i); 12 producer.sendMessage(msg); 13 System.out.println("Message Number "+ i +" sent."); 14 } 15 } 16 17 public static void main(String[] args) throws Exception{ 18 new TestMq(); 19 } 20}