java全家桶之64:java rabbitmq
·
以下是Java中使用RabbitMQ的完整技术指南,涵盖核心概念、代码实现及最佳实践:
🔧 一、环境准备与依赖配置
-
安装RabbitMQ服务器
-
需先安装Erlang(RabbitMQ基于Erlang开发),再安装对应版本的RabbitMQ 。
-
启动服务后,管理界面默认地址:
http://localhost:15672(账号/密码:guest/guest)。
-
-
Java项目依赖
Maven项目中添加客户端库:
<dependency> <groupId>com.rabbitmq</groupId> <artifactId>amqp-client</artifactId> <version>5.14.2</version> <!-- 使用最新版本 --> </dependency>
⚙️ 二、核心操作步骤
1. 创建连接与通道
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost"); // 服务器地址
factory.setPort(5672); // 端口(默认5672)
factory.setUsername("guest"); // 用户名
factory.setPassword("guest"); // 密码
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
// 后续操作在此执行
}
说明:
-
Channel是核心操作接口,用于声明队列、发送/接收消息 。 -
使用
try-with-resources自动关闭连接,避免资源泄漏 。
2. 声明队列
channel.queueDeclare("my_queue", true, false, false, null);
参数详解:
| 参数 | 含义 | 默认建议值 |
|---|---|---|
| 队列名 | 自定义标识(如:order_queue) |
必填 |
| 持久化 | 重启后队列是否保留 | true(生产环境) |
| 排他性 | 仅限当前连接使用 | false |
| 自动删除 | 无消费者时自动删除 | false |
| 额外参数 | 如消息TTL、死信队列等 | null |
3. 发送消息
String message = "订单已创建";
// 发送到默认交换机(Direct类型)
channel.basicPublish("", "my_queue", null, message.getBytes());
高级发送(持久化+确认机制):
channel.basicPublish(
"",
"my_queue",
MessageProperties.PERSISTENT_TEXT_PLAIN, // 消息持久化
message.getBytes()
);
channel.waitForConfirms(); // 等待Broker确认
4. 消费消息
DeliverCallback callback = (consumerTag, delivery) -> {
String msg = new String(delivery.getBody(), "UTF-8");
System.out.println("处理消息: " + msg);
// 手动确认(避免自动确认导致消息丢失)
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
};
channel.basicConsume("my_queue", false, callback, consumerTag -> {});
关键配置:
-
autoAck=false:关闭自动确认,确保业务处理成功后才确认消息 。 -
basicAck():手动发送确认信号,Broker收到后删除消息 。
🔁 三、交换机(Exchange)与路由
交换机类型与使用
| 类型 | 路由规则 | 代码示例 |
|---|---|---|
| Direct | 精确匹配路由键(routingKey) |
channel.exchangeDeclare("logs", "direct") |
| Topic | 通配符匹配(*.error) |
channel.exchangeDeclare("logs", "topic") |
| Fanout | 广播到所有绑定队列 | channel.exchangeDeclare("logs", "fanout") |
| Headers | 基于消息头匹配 | 少用,复杂 |
绑定队列示例:
channel.exchangeDeclare("order_exchange", "direct");
channel.queueDeclare("order_queue", true, false, false, null);
channel.queueBind("order_queue", "order_exchange", "order.create"); // 路由键绑定
🛡️ 四、可靠性保障机制
-
消息持久化
-
队列持久化:
channel.queueDeclare(..., true, ...) -
消息持久化:
MessageProperties.PERSISTENT_TEXT_PLAIN。
-
-
消费者确认模式
-
手动确认(
basicAck)替代自动确认,防止消息处理失败丢失 。
-
-
生产者确认机制
channel.confirmSelect(); // 开启确认模式 channel.addConfirmListener((seq, multiple) -> { System.out.println("消息已送达"); }, (seq, multiple) -> { System.out.println("消息送达失败"); });
🌱 五、Spring Boot集成(简化版)
-
依赖与配置
<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency>spring: rabbitmq: host: localhost port: 5672 username: guest password: guest -
发送与接收
// 发送方 @Service public class OrderSender { @Autowired private RabbitTemplate rabbitTemplate; public void send(String message) { rabbitTemplate.convertAndSend("order_exchange", "order.create", message); } } // 接收方 @Component public class OrderListener { @RabbitListener(queues = "order_queue") public void process(String message) { System.out.println("收到订单: " + message); } }
⚠️ 六、常见问题与优化
-
连接管理
-
复用连接,避免频繁创建(连接工厂是线程安全的)。
-
-
性能瓶颈
-
多消费者监听同一队列实现负载均衡 。
-
-
消息积压
-
监控队列长度(管理界面),动态增加消费者 。
-
-
死信队列
Map<String, Object> args = new HashMap<>(); args.put("x-dead-letter-exchange", "dead_letter_exchange"); channel.queueDeclare("order_queue", true, false, false, args);将处理失败的消息路由到死信队列 。
📚 完整代码示例
生产者:
public class Producer {
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
// ... 设置连接参数
try (Connection conn = factory.newConnection();
Channel channel = conn.createChannel()) {
channel.queueDeclare("test_queue", true, false, false, null);
channel.basicPublish("", "test_queue",
MessageProperties.PERSISTENT_TEXT_PLAIN,
"测试消息".getBytes());
}
}
}
消费者:
public class Consumer {
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
// ... 设置连接参数
Connection conn = factory.newConnection();
Channel channel = conn.createChannel();
channel.basicConsume("test_queue", false, (tag, delivery) -> {
String msg = new String(delivery.getBody());
System.out.println("消费: " + msg);
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
}, tag -> {});
}
}
通过以上步骤,可快速实现Java与RabbitMQ的集成。建议结合官方文档调整参数以适配生产环境需求 。
更多推荐

所有评论(0)