创建项目
创建一个maven项目rabbitmq
pom.xml
添加rabbitmq和hutool的依赖。
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>edu.hpu</groupId> <artifactId>rabbitmq</artifactId> <version>0.0.1-SNAPSHOT</version> <packaging>jar</packaging> <name>rabbitmq</name> <url>http://maven.apache.org</url> <properties> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> </properties> <build> <plugins> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-compiler-plugin</artifactId> <version>2.0.2</version> <configuration> <source>1.8</source> <target>1.8</target> </configuration> </plugin> </plugins> </build> <dependencies> <dependency> <groupId>com.rabbitmq</groupId> <artifactId>amqp-client</artifactId> <version>3.6.5</version> </dependency> <dependency> <groupId>cn.hutool</groupId> <artifactId>hutool-all</artifactId> <version>4.3.1</version> </dependency> </dependencies> </project>
工具类
RabbitMQUtil,用于判断服务器是否启动。
package edu.hpu.util; import javax.swing.JOptionPane; import cn.hutool.core.util.NetUtil; public class RabbitMQUtil { public static void main(String[] args) { checkServer(); } public static void checkServer() { if(NetUtil.isUsableLocalPort(15672)) { JOptionPane.showMessageDialog(null, "RabbitMQ 服务器未启动 "); System.exit(1); } } }
Fanout模式实例
消息生产类
TestFanoutProducter,
package edu.hpu; import java.io.IOException; import java.util.concurrent.TimeoutException; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import edu.hpu.util.RabbitMQUtil; /* * 消息生成者 * Fanout模式 */ public class TestFanoutProducter { public final static String EXCHANGE_NAME="fanout_exchange"; public static void main(String[] args) throws IOException, TimeoutException { RabbitMQUtil.checkServer(); //创建连接工厂 ConnectionFactory factory = new ConnectionFactory(); //设置RabbitMQ相关信息 factory.setHost("localhost"); //创建一个新的连接 Connection connection = factory.newConnection(); //创建一个通道 Channel channel = connection.createChannel(); channel.exchangeDeclare(EXCHANGE_NAME, "fanout"); for (int i = 0; i < 100; i++) { String message = "direct 消息 " +i; //发送消息到队列中 channel.basicPublish(EXCHANGE_NAME, "", null, message.getBytes("UTF-8")); System.out.println("发送消息: " + message); } //关闭通道和连接 channel.close(); connection.close(); } }
消息消费类
TestFanoutConsumer,
package edu.hpu; import java.io.IOException; import java.util.concurrent.TimeoutException; import com.rabbitmq.client.AMQP; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.Consumer; import com.rabbitmq.client.DefaultConsumer; import com.rabbitmq.client.Envelope; import cn.hutool.core.util.RandomUtil; import edu.hpu.util.RabbitMQUtil; public class TestFanoutConsumer { public final static String EXCHANGE_NAME="fanout_exchange"; public static void main(String[] args) throws IOException, TimeoutException { //为当前消费者取随机名 String name = "consumer-"+ RandomUtil.randomString(5); //判断服务器是否启动 RabbitMQUtil.checkServer(); // 创建连接工厂 ConnectionFactory factory = new ConnectionFactory(); //设置RabbitMQ地址 factory.setHost("localhost"); //创建一个新的连接 Connection connection = factory.newConnection(); //创建一个通道 Channel channel = connection.createChannel(); //交换机声明(参数为:交换机名称;交换机类型) channel.exchangeDeclare(EXCHANGE_NAME,"fanout"); //获取一个临时队列 String queueName = channel.queueDeclare().getQueue(); //队列与交换机绑定(参数为:队列名称;交换机名称;routingKey忽略) channel.queueBind(queueName,EXCHANGE_NAME,""); System.out.println(name +" 等待接受消息"); //DefaultConsumer类实现了Consumer接口,通过传入一个频道, // 告诉服务器我们需要那个频道的消息,如果频道中有消息,就会执行回调函数handleDelivery Consumer consumer = new DefaultConsumer(channel) { @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { String message = new String(body, "UTF-8"); System.out.println(name + " 接收到消息 '" + message + "'"); } }; //自动回复队列应答 -- RabbitMQ中的消息确认机制 channel.basicConsume(queueName, true, consumer); } }
运行
运行消息消费类两次,产生两个消息消费者,运行消息生产者一次,产生消息生产者。
Direct模式
消息生产类
TestDriectProducer,
package edu.hpu; import java.io.IOException; import java.util.concurrent.TimeoutException; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import edu.hpu.util.RabbitMQUtil; public class TestDriectProducer { public final static String QUEUE_NAME="direct_queue"; public static void main(String[] args) throws IOException, TimeoutException { RabbitMQUtil.checkServer(); //创建连接工厂 ConnectionFactory factory = new ConnectionFactory(); //设置RabbitMQ相关信息 factory.setHost("localhost"); //创建一个新的连接 Connection connection = factory.newConnection(); //创建一个通道 Channel channel = connection.createChannel(); for (int i = 0; i < 100; i++) { String message = "direct 消息 " +i; //发送消息到队列中 channel.basicPublish("", QUEUE_NAME, null, message.getBytes("UTF-8")); System.out.println("发送消息: " + message); } //关闭通道和连接 channel.close(); connection.close(); } }
消息消费类
TestDriectCustomer,
package edu.hpu; import java.io.IOException; import java.util.concurrent.TimeoutException; import com.rabbitmq.client.AMQP; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.Consumer; import com.rabbitmq.client.DefaultConsumer; import com.rabbitmq.client.Envelope; import cn.hutool.core.util.RandomUtil; import edu.hpu.util.RabbitMQUtil; public class TestDriectCustomer { private final static String QUEUE_NAME = "direct_queue"; public static void main(String[] args) throws IOException, TimeoutException { //为当前消费者取随机名 String name = "consumer-"+ RandomUtil.randomString(5); //判断服务器是否启动 RabbitMQUtil.checkServer(); // 创建连接工厂 ConnectionFactory factory = new ConnectionFactory(); //设置RabbitMQ地址 factory.setHost("localhost"); //创建一个新的连接 Connection connection = factory.newConnection(); //创建一个通道 Channel channel = connection.createChannel(); //声明要关注的队列 channel.queueDeclare(QUEUE_NAME, false, false, true, null); System.out.println(name +" 等待接受消息"); //DefaultConsumer类实现了Consumer接口,通过传入一个频道, // 告诉服务器我们需要那个频道的消息,如果频道中有消息,就会执行回调函数handleDelivery Consumer consumer = new DefaultConsumer(channel) { @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { String message = new String(body, "UTF-8"); System.out.println(name + " 接收到消息 '" + message + "'"); } }; //自动回复队列应答 -- RabbitMQ中的消息确认机制 channel.basicConsume(QUEUE_NAME, true, consumer); } }
启动
结束Fanout模式的进程,运行消息消费类两次,产生两个消息消费者,运行消息生产者一次,产生消息生产者。消息消费者分食消息生产者产生的消息。
Topic模式
消息生产类
TestTopicProducter,四个路由:“usa.news”, “usa.weather”, “europe.news”, “europe.weather” 上发布 “美国新闻”, “美国天气”, “欧洲新闻”, “欧洲天气”。
package edu.hpu; import java.io.IOException; import java.util.concurrent.TimeoutException; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import edu.hpu.util.RabbitMQUtil; public class TestTopicProducter { public final static String EXCHANGE_NAME="topics_exchange"; public static void main(String[] args) throws IOException, TimeoutException { RabbitMQUtil.checkServer(); //创建连接工厂 ConnectionFactory factory = new ConnectionFactory(); //设置RabbitMQ相关信息 factory.setHost("localhost"); //创建一个新的连接 Connection connection = factory.newConnection(); //创建一个通道 Channel channel = connection.createChannel(); channel.exchangeDeclare(EXCHANGE_NAME, "topic"); String[] routing_keys = new String[] { "usa.news", "usa.weather", "europe.news", "europe.weather" }; String[] messages = new String[] { "美国新闻", "美国天气", "欧洲新闻", "欧洲天气" }; for (int i = 0; i < routing_keys.length; i++) { String routingKey = routing_keys[i]; String message = messages[i]; channel.basicPublish(EXCHANGE_NAME, routingKey, null, message .getBytes()); System.out.printf("发送消息到路由:%s, 内容是: %s%n ", routingKey,message); } //关闭通道和连接 channel.close(); connection.close(); } }
消息消费类
两个消费类,分别根据符号匹配接受不同的消息。
TestTopicCustomer4USA
接收usa.* 消息,
package edu.hpu; import java.io.IOException; import java.util.concurrent.TimeoutException; import com.rabbitmq.client.AMQP; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.Consumer; import com.rabbitmq.client.DefaultConsumer; import com.rabbitmq.client.Envelope; import edu.hpu.util.RabbitMQUtil; public class TestTopicCustomer4USA { public final static String EXCHANGE_NAME="topics_exchange"; public static void main(String[] args) throws IOException, TimeoutException { //为当前消费者取名称 String name = "consumer-usa"; //判断服务器是否启动 RabbitMQUtil.checkServer(); // 创建连接工厂 ConnectionFactory factory = new ConnectionFactory(); //设置RabbitMQ地址 factory.setHost("localhost"); //创建一个新的连接 Connection connection = factory.newConnection(); //创建一个通道 Channel channel = connection.createChannel(); //交换机声明(参数为:交换机名称;交换机类型) channel.exchangeDeclare(EXCHANGE_NAME,"topic"); //获取一个临时队列 String queueName = channel.queueDeclare().getQueue(); //接受 USA 信息 channel.queueBind(queueName, EXCHANGE_NAME, "usa.*"); System.out.println(name +" 等待接受消息"); //DefaultConsumer类实现了Consumer接口,通过传入一个频道, // 告诉服务器我们需要那个频道的消息,如果频道中有消息,就会执行回调函数handleDelivery Consumer consumer = new DefaultConsumer(channel) { @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { String message = new String(body, "UTF-8"); System.out.println(name + " 接收到消息 '" + message + "'"); } }; //自动回复队列应答 -- RabbitMQ中的消息确认机制 channel.basicConsume(queueName, true, consumer); } }
TestTopicCustomer4News
接收*.news 消息,
package edu.hpu; import java.io.IOException; import java.util.concurrent.TimeoutException; import com.rabbitmq.client.AMQP; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.Consumer; import com.rabbitmq.client.DefaultConsumer; import com.rabbitmq.client.Envelope; import edu.hpu.util.RabbitMQUtil; public class TestTopicCustomer4News { public final static String EXCHANGE_NAME="topics_exchange"; public static void main(String[] args) throws IOException, TimeoutException { //为当前消费者取名称 String name = "consumer-news"; //判断服务器是否启动 RabbitMQUtil.checkServer(); // 创建连接工厂 ConnectionFactory factory = new ConnectionFactory(); //设置RabbitMQ地址 factory.setHost("localhost"); //创建一个新的连接 Connection connection = factory.newConnection(); //创建一个通道 Channel channel = connection.createChannel(); //交换机声明(参数为:交换机名称;交换机类型) channel.exchangeDeclare(EXCHANGE_NAME,"topic"); //获取一个临时队列 String queueName = channel.queueDeclare().getQueue(); //接受 USA 信息 channel.queueBind(queueName, EXCHANGE_NAME, "*.news"); System.out.println(name +" 等待接受消息"); //DefaultConsumer类实现了Consumer接口,通过传入一个频道, // 告诉服务器我们需要那个频道的消息,如果频道中有消息,就会执行回调函数handleDelivery Consumer consumer = new DefaultConsumer(channel) { @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { String message = new String(body, "UTF-8"); System.out.println(name + " 接收到消息 '" + message + "'"); } }; //自动回复队列应答 -- RabbitMQ中的消息确认机制 channel.basicConsume(queueName, true, consumer); } }
运行
运行两个消息生产类,运行消息生产类。
参考:
【1】、http://how2j.cn/k/message/message-rabbitmq-fanout/2034.html
【2】、http://how2j.cn/k/message/message-rabbitmq-direct/2032.html
【3】、http://how2j.cn/k/message/message-rabbitmq-topic/2033.html
【4】、https://blog.csdn.net/hao134838/article/details/71710067





