开发手册
# 前言
本文档为金蝶Apusic分布式消息队列for RabbitMQ(Apusic Distributed Message Queue,简称:ADMQ for RabbitMQ)的开发使用说明,帮助用户快速学习如何使用金蝶Apusic分布式消息队列for RabbitMQ进行开发。
# 适用对象
本文档适用于IT信息化业务负责人、研发经理、软件项目经理、软件架构师、运维工程师。
# 相关文档
了解更多ADMQ for RabbitMQ产品相关的信息,请参阅以下ADMQ for RabbitMQ产品手册文档集:
| 序号 | 手册文档 | 说明 |
|---|---|---|
| 1 | 金蝶Apusic分布式消息队列for RabbitMQ 快速使用手册 | 简单介绍了如何快速上手使用ADMQ for RabbitMQ 。 |
| 2 | 金蝶Apusic分布式消息队列for RabbitMQ 安装手册 | 详细介绍如何在各操作系统上安装ADMQ for RabbitMQ,以及ADMQ for RabbitMQ服务启停等操作。 |
| 3 | 金蝶Apusic分布式消息队列for RabbitMQ 消息引擎用户手册 | 详细介绍 ADMQ for RabbitMQ 消息引擎相关功能的使用、配置、管理及配套工具的使用方法。 |
| 4 | 金蝶Apusic分布式消息队列for RabbitMQ 管控台用户手册 | 详细介绍ADMQ for RabbitMQ管控台相关功能的使用和操作说明。 |
| 5 | 金蝶Apusic分布式消息队列for RabbitMQ 开发手册 | 详细介绍基于各开发语言进行ADMQ for RabbitMQ客户端应用开发的说明。 |
| 6 | 金蝶Apusic分布式消息队列for RabbitMQ 迁移手册 | 详细介绍从RabbitMQ迁移到ADMQ for RabbitMQ的说明。 |
| 7 | 金蝶Apusic分布式消息队列for RabbitMQ 运维手册 | 详细介绍ADMQ for RabbitMQ的监控、运维、安全加固等运维说明。 |
| 8 | 金蝶Apusic分布式消息队列for RabbitMQ 性能优化手册 | 详细介绍ADMQ for RabbitMQ性能调优的说明。 |
# 技术支持
ADMQ for RabbitMQ产品提供全面的技术支持服务,您可以通过以下方式获得技术支持:
- 网址:www.apusic.com
- 电话:400-855-5800
- 邮箱:support@apusic.com
- 金蝶云社区:https://vip.kingdee.com/?productId=73&productLineId=14&lang=zh-CN
您在取得技术支持时,请提供如下信息:
您的姓名
公司信息与联系方式
操作系统及其版本
产品版本号
出现异常及错误的日志、截图等详细信息
# 开发准备
# 连接参数
| 参数 | 说明 | 示例 |
|---|---|---|
| host | 服务器地址 | localhost |
| port | 端口 | 5672(AMQP)/ 5671(AMQP over SSL) |
| username | 用户名 | guest |
| password | 密码 | guest |
| virtualHost | 虚拟主机 | / |
# 客户端依赖
# Java
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>5.21.0</version>
</dependency>
1
2
3
4
5
2
3
4
5
# Python
pip install pika
1
# Go
go get github.com/rabbitmq/amqp091-go
1
# Node.js
npm install amqplib
1
# Java 开发示例
# 生产者
import com.rabbitmq.client.*;
public class Producer {
private static final String EXCHANGE_NAME = "demo.exchange";
private static final String ROUTING_KEY = "demo.key";
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setPort(5672);
factory.setUsername("guest");
factory.setPassword("guest");
factory.setVirtualHost("/");
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
// 声明交换机
channel.exchangeDeclare(EXCHANGE_NAME, "direct", true);
// 发送消息
String message = "Hello ADMQ for RabbitMQ!";
channel.basicPublish(EXCHANGE_NAME, ROUTING_KEY,
MessageProperties.PERSISTENT_TEXT_PLAIN,
message.getBytes("UTF-8"));
System.out.println("Sent: " + message);
}
}
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
# 消费者
import com.rabbitmq.client.*;
public class Consumer {
private static final String QUEUE_NAME = "demo.queue";
private static final String EXCHANGE_NAME = "demo.exchange";
private static final String ROUTING_KEY = "demo.key";
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setPort(5672);
factory.setUsername("guest");
factory.setPassword("guest");
factory.setVirtualHost("/");
Connection connection = factory.newConnection();
Channel channel = connection.createChannel();
// 声明队列
channel.queueDeclare(QUEUE_NAME, true, false, false, null);
// 绑定队列到交换机
channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, ROUTING_KEY);
// 创建消费者
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), "UTF-8");
System.out.println("Received: " + message);
// 手动确认
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
};
// 消费消息(关闭自动确认)
channel.basicConsume(QUEUE_NAME, false, deliverCallback, consumerTag -> {});
}
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
# 事务消息
channel.txSelect(); // 开启事务
try {
channel.basicPublish(EXCHANGE, ROUTING_KEY, null, message1.getBytes());
channel.basicPublish(EXCHANGE, ROUTING_KEY, null, message2.getBytes());
channel.txCommit(); // 提交事务
} catch (Exception e) {
channel.txRollback(); // 回滚事务
}
1
2
3
4
5
6
7
8
9
2
3
4
5
6
7
8
9
# Confirm 模式
channel.confirmSelect(); // 开启 Confirm 模式
channel.basicPublish(EXCHANGE, ROUTING_KEY, null, message.getBytes());
// 同步等待确认
if (channel.waitForConfirms()) {
System.out.println("Message confirmed");
}
// 异步确认
channel.addConfirmListener(new ConfirmListener() {
@Override
public void handleAck(long deliveryTag, boolean multiple) {
System.out.println("Message ack: " + deliveryTag);
}
@Override
public void handleNack(long deliveryTag, boolean multiple) {
System.out.println("Message nack: " + deliveryTag);
}
});
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
# Python 开发示例
# 生产者
import pika
connection = pika.BlockingConnection(
pika.ConnectionParameters(host='localhost', port=5672,
credentials=pika.PlainCredentials('guest', 'guest')))
channel = connection.channel()
channel.exchange_declare(exchange='demo.exchange', exchange_type='direct', durable=True)
channel.basic_publish(exchange='demo.exchange',
routing_key='demo.key',
body='Hello ADMQ for RabbitMQ!',
properties=pika.BasicProperties(delivery_mode=2))
print("Sent message")
connection.close()
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
# 消费者
import pika
connection = pika.BlockingConnection(
pika.ConnectionParameters(host='localhost', port=5672,
credentials=pika.PlainCredentials('guest', 'guest')))
channel = connection.channel()
channel.queue_declare(queue='demo.queue', durable=True)
channel.queue_bind(queue='demo.queue', exchange='demo.exchange', routing_key='demo.key')
def callback(ch, method, properties, body):
print(f"Received {body.decode()}")
ch.basic_ack(delivery_tag=method.delivery_tag)
channel.basic_consume(queue='demo.queue', on_message_callback=callback, auto_ack=False)
print('Waiting for messages...')
channel.start_consuming()
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
# Go 开发示例
# 生产者
package main
import (
"fmt"
"github.com/rabbitmq/amqp091-go"
"log"
)
func main() {
conn, err := amqp091.Dial("amqp://guest:guest@localhost:5672/")
if err != nil {
log.Fatal(err)
}
defer conn.Close()
ch, err := conn.Channel()
if err != nil {
log.Fatal(err)
}
defer ch.Close()
err = ch.ExchangeDeclare("demo.exchange", "direct", true, false, false, false, nil)
if err != nil {
log.Fatal(err)
}
body := "Hello ADMQ for RabbitMQ!"
err = ch.Publish("demo.exchange", "demo.key", false, false,
amqp091.Publishing{
DeliveryMode: amqp091.Persistent,
ContentType: "text/plain",
Body: []byte(body),
})
if err != nil {
log.Fatal(err)
}
fmt.Println("Sent:", body)
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
# 消费者
package main
import (
"fmt"
"github.com/rabbitmq/amqp091-go"
"log"
)
func main() {
conn, err := amqp091.Dial("amqp://guest:guest@localhost:5672/")
if err != nil {
log.Fatal(err)
}
defer conn.Close()
ch, err := conn.Channel()
if err != nil {
log.Fatal(err)
}
defer ch.Close()
q, err := ch.QueueDeclare("demo.queue", true, false, false, false, nil)
if err != nil {
log.Fatal(err)
}
err = ch.QueueBind(q.Name, "demo.key", "demo.exchange", false, nil)
if err != nil {
log.Fatal(err)
}
msgs, err := ch.Consume(q.Name, "", false, false, false, false, nil)
if err != nil {
log.Fatal(err)
}
fmt.Println("Waiting for messages...")
for msg := range msgs {
fmt.Printf("Received: %s\n", msg.Body)
msg.Ack(false)
}
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
# Node.js 开发示例
# 生产者
const amqp = require('amqplib');
async function produce() {
const conn = await amqp.connect('amqp://guest:guest@localhost:5672');
const ch = await conn.createChannel();
await ch.assertExchange('demo.exchange', 'direct', { durable: true });
const msg = 'Hello ADMQ for RabbitMQ!';
ch.publish('demo.exchange', 'demo.key', Buffer.from(msg), { persistent: true });
console.log('Sent:', msg);
await ch.close();
await conn.close();
}
produce().catch(console.error);
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
# 消费者
const amqp = require('amqplib');
async function consume() {
const conn = await amqp.connect('amqp://guest:guest@localhost:5672');
const ch = await conn.createChannel();
await ch.assertQueue('demo.queue', { durable: true });
await ch.bindQueue('demo.queue', 'demo.exchange', 'demo.key');
console.log('Waiting for messages...');
ch.consume('demo.queue', (msg) => {
console.log('Received:', msg.content.toString());
ch.ack(msg);
}, { noAck: false });
}
consume().catch(console.error);
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
# 高级特性
# 延迟消息
使用死信交换机 + TTL 实现延迟队列:
// 创建死信交换机和队列
channel.exchangeDeclare("dlx.exchange", "direct", true);
channel.queueDeclare("dlx.queue", true, false, false, null);
channel.queueBind("dlx.queue", "dlx.exchange", "dlx.key");
// 创建延迟队列(设置 TTL 和死信参数)
Map<String, Object> args = new HashMap<>();
args.put("x-message-ttl", 60000); // 60秒延迟
args.put("x-dead-letter-exchange", "dlx.exchange");
args.put("x-dead-letter-routing-key", "dlx.key");
channel.queueDeclare("delay.queue", true, false, false, args);
// 发送消息到延迟队列
channel.basicPublish("", "delay.queue", null, message.getBytes());
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
2
3
4
5
6
7
8
9
10
11
12
13
14
15
# 优先级队列
Map<String, Object> args = new HashMap<>();
args.put("x-max-priority", 10);
channel.queueDeclare("priority.queue", true, false, false, args);
// 发送高优先级消息
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
.priority(9)
.build();
channel.basicPublish("", "priority.queue", props, message.getBytes());
1
2
3
4
5
6
7
8
9
2
3
4
5
6
7
8
9
# 批量消费
// 设置 QoS,每次预取 100 条消息
channel.basicQos(100);
// 批量确认
channel.basicAck(deliveryTag, true); // multiple=true 确认所有小于等于该 deliveryTag 的消息
1
2
3
4
5
2
3
4
5
# 最佳实践
# 连接管理
- 使用连接池:避免频繁创建和关闭连接
- 合理设置心跳:防止网络异常导致连接假死
- 异常处理:捕获异常并进行重连
# 消息可靠性
- 使用手动确认:确保消息被正确处理后确认
- 消息持久化:重要消息设置 delivery_mode=2
- 队列持久化:重要队列设置 durable=true
- 使用 Confirm 模式:确保消息成功发送到交换机
# 性能优化
- 批量操作:批量发送消息、批量确认
- 合理设置 QoS:避免消费者积压过多消息
- 使用多通道:单个连接上创建多个通道并发处理
- 选择合适的队列类型:Quorum Queue 适合高可用场景,Stream Queue 适合高吞吐场景
# 异常处理
// 连接断开的回调
factory.setExceptionHandler(new DefaultExceptionHandler() {
@Override
public void handleConnectionRecoveryException(Connection conn, Throwable exception) {
// 记录日志并进行重连
}
});
// 自动恢复
factory.setAutomaticRecoveryEnabled(true);
factory.setNetworkRecoveryInterval(5000); // 5秒重试间隔
factory.setTopologyRecoveryEnabled(true); // 自动恢复拓扑结构
1
2
3
4
5
6
7
8
9
10
11
12
2
3
4
5
6
7
8
9
10
11
12
编辑页面 (opens new window)