RabbitMQ实现延迟消息
前提条件
确保RabbitMQ已安装并启用了RabbitMQ Delayed Message插件。如果尚未启用,可以按照以下步骤操作:
(图片来源网络,侵删)
-
下载插件:
- 从RabbitMQ社区插件页面下载rabbitmq_delayed_message_exchange插件。
-
安装插件:
- 将插件文件(.ez文件)放置在RabbitMQ插件目录中,通常为/usr/lib/rabbitmq/lib/rabbitmq_server-/plugins。
-
启用插件:
rabbitmq-plugins enable rabbitmq_delayed_message_exchange
Maven依赖
在你的Maven项目的pom.xml中添加RabbitMQ客户端库的依赖:
com.rabbitmq
amqp-client
5.13.0
生产者(Producer)代码
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.BuiltinExchangeType;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.util.HashMap;
import java.util.Map;
public class Producer {
private final static String EXCHANGE_NAME = "delayed_exchange";
private final static String QUEUE_NAME = "delayed_queue";
private final static String ROUTING_KEY = "delayed_key";
public static void main(String[] argv) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
// 声明延迟交换机
Map args = new HashMap();
args.put("x-delayed-type", BuiltinExchangeType.DIRECT.getType());
channel.exchangeDeclare(EXCHANGE_NAME, "x-delayed-message", true, false, args);
// 声明队列
channel.queueDeclare(QUEUE_NAME, true, false, false, null);
// 绑定队列到延迟交换机
channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, ROUTING_KEY);
String message = "Hello World with delay!";
int delay = 5000; // 延迟时间,以毫秒为单位
// 设置消息属性,包括延迟时间
Map headers = new HashMap();
headers.put("x-delay", delay);
AMQP.BasicProperties.Builder props = new AMQP.BasicProperties.Builder()
.headers(headers)
.deliveryMode(2); // 使消息持久化
// 发布消息
channel.basicPublish(EXCHANGE_NAME, ROUTING_KEY, props.build(), message.getBytes("UTF-8"));
System.out.println(" [x] Sent '" + message + "' with delay " + delay + " ms");
}
}
}
消费者(Consumer)代码
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;
public class Consumer {
private final static String QUEUE_NAME = "delayed_queue";
public static void main(String[] argv) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
// 声明队列
channel.queueDeclare(QUEUE_NAME, true, false, false, null);
System.out.println(" [*] Waiting for messages. To exit press CTRL+C");
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), "UTF-8");
System.out.println(" [x] Received '" + message + "'");
};
// 监听队列并处理消息
channel.basicConsume(QUEUE_NAME, true, deliverCallback, consumerTag -> { });
}
}
}
说明
-
生产者(Producer)代码:
- 声明了一个延迟交换机,类型为x-delayed-message,并设置了x-delayed-type参数为direct。
- 声明了一个队列并将其绑定到延迟交换机。
- 设置消息的x-delay头部,指定延迟时间。
- 将消息发布到延迟交换机。
-
消费者(Consumer)代码:
- 声明了一个队列。
- 使用DeliverCallback监听队列并处理收到的消息。
通过上述步骤和代码示例,你可以在Java中实现RabbitMQ的延迟消息功能。
免责声明:我们致力于保护作者版权,注重分享,被刊用文章因无法核实真实出处,未能及时与作者取得联系,或有版权异议的,请联系管理员,我们会立即处理! 部分文章是来自自研大数据AI进行生成,内容摘自(百度百科,百度知道,头条百科,中国民法典,刑法,牛津词典,新华词典,汉语词典,国家院校,科普平台)等数据,内容仅供学习参考,不准确地方联系删除处理! 图片声明:本站部分配图来自人工智能系统AI生成,觅知网授权图片,PxHere摄影无版权图库和百度,360,搜狗等多加搜索引擎自动关键词搜索配图,如有侵权的图片,请第一时间联系我们,邮箱:ciyunidc@ciyunshuju.com。本站只作为美观性配图使用,无任何非法侵犯第三方意图,一切解释权归图片著作权方,本站不承担任何责任。如有恶意碰瓷者,必当奉陪到底严惩不贷!
