RabbitMQ简单示例
Hello RabbitMQ
添加依赖
先创建好 Maven 工程,pom.xml 添入依赖:
<dependencies>
<!--rabbitmq 依赖客户端-->
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>5.8.0</version>
</dependency>
<!--操作文件流的一个依赖-->
<dependency>
<groupId>commons-io</groupId>
<artifactId>commons-io</artifactId>
<version>2.6</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<configuration>
<source>8</source>
<target>8</target>
</configuration>
</plugin>
</plugins>
</build>
消息生产者
/**
* 生产者:发消息
*/
public class Producer {
//对列名称
public static final String QUEUE_NAME="hello";
//发消息
public static void main(String[] args) throws IOException, TimeoutException {
//创建一个连接工厂
ConnectionFactory factory = new ConnectionFactory();
//工厂IP 连接RabbitMQ对列
factory.setHost("localhost");
//用户名
factory.setUsername("guest");
//密码
factory.setPassword("guest");
//创建连接
Connection connection = factory.newConnection();
//获取信道
Channel channel = connection.createChannel();
/**
* 生产一个对列
* 1.对列名称
* 2.对列里面的消息是否持久化,默认情况下,消息存储在内存中
* 3.该队列是否只供一个消费者进行消费,是否进行消息共享,true可以多个消费者消费 false:只能一个消费者消费
* 4.是否自动删除,最后一个消费者端开链接以后,该队列是否自动删除,true表示自动删除
* 5.其他参数
*/
channel.queueDeclare(QUEUE_NAME,false,false,false,null);
//发消息
String message = "Hello,RabbitMQ";
/**
* 发送一个消息
* 1.发送到哪个交换机
* 2.路由的key值是哪个本次是队列的名称
* 3.其他参数信息
* 4.发送消息的消息体
*/
channel.basicPublish("",QUEUE_NAME,null,message.getBytes());
System.out.println("消息发送完毕");
}
}
消息消费者
/**
* 消费者:接受消息
*/
public class Consumer {
//队列的名称
public static final String QUEUE_NAME="hello";
//接受消息
public static void main(String[] args) throws IOException, TimeoutException {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setUsername("guest");
factory.setPassword("guest");
Connection connection = factory.newConnection();
Channel channel = connection.createChannel();
//声明接收消息
DeliverCallback deliverCallback = (consumerTag,message) -> {
System.out.println(new String(message.getBody()));
};
//取消消息时的回调
CancelCallback cancelCallback = consumerTag ->{
System.out.println("消息消费被中断");
};
/**
* 消费者消费消息
* 1.消费哪个队列
* 2.消费成功之后是否要自动应答true:代表自动应答false:代表手动应答
* 3.消费者未成功消费的回调
* 4.消费者取消消费的回调
*/
channel.basicConsume(QUEUE_NAME,true,deliverCallback,cancelCallback);
}
}
Work Queues
Work Queues 是工作队列(又称任务队列)的主要思想是避免立即执行资源密集型任务,而不得不等待它完成。相反我们安排任务在之后执行。我们把任务封装为消息并将其发送到队列。在后台运行的工作进程将弹出任务并最终执行作业。当有多个工作线程时,这些工作线程将一起处理这些任务。
轮询消费
轮询消费消息指的是轮流消费消息,即每个工作队列都会获取一个消息进行消费,并且获取的次数按照顺序依次往下轮流。
案例中生产者叫做 Task,一个消费者就是一个工作队列,启动两个工作队列消费消息,这个两个工作队列会以轮询的方式消费消息。
轮询案例
封装工具类:RabbitMQUtils
/**
* 此类为连接工厂创建信道的工具类
*/
public class RabbitMQUtils {
public static Channel getChannel() throws IOException, TimeoutException {
//创建一个连接工厂
ConnectionFactory factory = new ConnectionFactory();
//工厂IP 连接RabbitMQ对列
factory.setHost("localhost");
//用户名
factory.setUsername("guest");
//密码
factory.setPassword("guest");
//创建连接
Connection connection = factory.newConnection();
//获取信道
Channel channel = connection.createChannel();
return channel;
}
}
创建两个工作队列,并且启动
/**
* 这是一个工作线程,相当于之间的Consumer
*/
public class Work01 {
//队列的名称
public static final String QUEUE_NAME="hello";
//接收消息
public static void main(String[] args) throws IOException, TimeoutException {
Channel channel = RabbitMQUtils.getChannel();
//消息的接受
DeliverCallback deliverCallback = (consumerTag,message) ->{
System.out.println("接收到的消息:"+new String(message.getBody()));
};
//消息接受被取消时,执行下面的内容
CancelCallback cancelCallback = consumerTag -> {
System.out.println(consumerTag+"消息被消费者取消消费接口回调逻辑");
};
//消息的接受
channel.basicConsume(QUEUE_NAME,true,deliverCallback,cancelCallback);
}
}
创建好一个工作队列,只需要以多线程方式启动两次该 main 函数即可,以 first、second 区别消息队列。
创建生产者
/**
* 生产者
*/
public class Task01 {
//队列名称
public static final String QUEUE_NAME="hello";
//发送大量消息
public static void main(String[] args) throws IOException, TimeoutException {
Channel channel = RabbitMQUtils.getChannel();
//队列的声明
channel.queueDeclare(QUEUE_NAME,false,false,false,null);
//发送消息
//从控制台当中接受信息
Scanner scanner = new Scanner(System.in);
while (scanner.hasNext()) {
String message = scanner.next();
channel.basicPublish("",QUEUE_NAME,null,message.getBytes());
System.out.println("消息发送完成:"+message);
}
}
}