Apusic文档中心
首页
  • 应用服务器 AAS
  • 负载均衡器 ALB
  • 分布式消息队列 ADMQ
  • 分布式缓存 AMDC
  • 分布式配置中心 ADCC
  • Java开发工具包软件 AJDK
  • 搜索引擎 ASE
  • 中间件云平台 ACP
  • 统一管理平台 AUMP
  • 云原生中间件管理 ACMP
  • DevOps平台 ADOP
  • 许可授权中心 ACLS
  • Copilot智能问答系统 ACS
  • 监控平台 AMP
  • 智能日志 AILP
  • 应用性能管理 AAPM
  • 智能告警 AAlarm
  • 主数据管理 AMDM
  • 数据交换平台 ADXP
  • 企业服务总线 AESB
  • 数据智脑 ADPR
  • 服务治理 ASGP
  • 统一身份管理 AIDM
  • 标准模板
  • Markdown教程 (opens new window)
  • VuePress官方社区 (opens new window)
  • 帮助
贡献文档 (opens new window)
首页
  • 应用服务器 AAS
  • 负载均衡器 ALB
  • 分布式消息队列 ADMQ
  • 分布式缓存 AMDC
  • 分布式配置中心 ADCC
  • Java开发工具包软件 AJDK
  • 搜索引擎 ASE
  • 中间件云平台 ACP
  • 统一管理平台 AUMP
  • 云原生中间件管理 ACMP
  • DevOps平台 ADOP
  • 许可授权中心 ACLS
  • Copilot智能问答系统 ACS
  • 监控平台 AMP
  • 智能日志 AILP
  • 应用性能管理 AAPM
  • 智能告警 AAlarm
  • 主数据管理 AMDM
  • 数据交换平台 ADXP
  • 企业服务总线 AESB
  • 数据智脑 ADPR
  • 服务治理 ASGP
  • 统一身份管理 AIDM
  • 标准模板
  • Markdown教程 (opens new window)
  • VuePress官方社区 (opens new window)
  • 帮助
贡献文档 (opens new window)
文档中心
  • 金蝶Apusic应用服务器

  • 金蝶Apusic负载均衡器

  • 金蝶Apusic分布式消息队列

    • 产品白皮书
    • 产品更新说明
    • 统一管控台

    • V2.0.6(最新)

    • V2.0.6_for_kafka

    • V2.0.6_for_rabbitmq

      • 产品简介
      • 功能与技术规格
      • 安装手册
      • 快速使用手册
      • 管控台用户手册
      • 消息引擎用户手册
      • 开发手册
      • 迁移手册
      • 运维手册
      • 性能优化手册
      • 极简运维手册
    • V2.0.6_for_rocketmq

    • V2.0.6_for_MQTT

    • V2.0.5

    • V2.0.4

    • V2.0.3

  • 金蝶Apusic分布式缓存

  • 金蝶Apusic分布式配置中心

  • 金蝶Apusic Java开发工具包软件

  • 金蝶Apusic全文检索

开发手册

# 前言

本文档为金蝶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

您在取得技术支持时,请提供如下信息:

  1. 您的姓名

  2. 公司信息与联系方式

  3. 操作系统及其版本

  4. 产品版本号

  5. 出现异常及错误的日志、截图等详细信息

# 开发准备

# 连接参数

参数 说明 示例
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

# 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

# 消费者

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

# 事务消息

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

# 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

# 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

# 消费者

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

# 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

# 消费者

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

# 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

# 消费者

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

# 高级特性

# 延迟消息

使用死信交换机 + 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

# 优先级队列

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

# 批量消费

// 设置 QoS,每次预取 100 条消息
channel.basicQos(100);

// 批量确认
channel.basicAck(deliveryTag, true);  // multiple=true 确认所有小于等于该 deliveryTag 的消息
1
2
3
4
5

# 最佳实践

# 连接管理

  1. 使用连接池:避免频繁创建和关闭连接
  2. 合理设置心跳:防止网络异常导致连接假死
  3. 异常处理:捕获异常并进行重连

# 消息可靠性

  1. 使用手动确认:确保消息被正确处理后确认
  2. 消息持久化:重要消息设置 delivery_mode=2
  3. 队列持久化:重要队列设置 durable=true
  4. 使用 Confirm 模式:确保消息成功发送到交换机

# 性能优化

  1. 批量操作:批量发送消息、批量确认
  2. 合理设置 QoS:避免消费者积压过多消息
  3. 使用多通道:单个连接上创建多个通道并发处理
  4. 选择合适的队列类型: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
编辑页面 (opens new window)
#开发手册

← 消息引擎用户手册 迁移手册→

  • 浅色模式