快速使用手册
# 前言
本文档为金蝶Apusic分布式消息队列for RocketMQ(Apusic Distributed Message Queue for RocketMQ,简称:ADMQ for RocketMQ)的快速使用手册,提供产品介绍,安装部署产品的基本步骤及验证方式。
# 适用对象
本文档适用于IT信息化业务负责人、研发经理、软件项目经理、软件架构师、运维工程师。
# 相关文档
了解更多ADMQ for RocketMQ产品相关的信息,请参阅以下ADMQ for RocketMQ产品手册文档集:
| 序号 | 手册文档 | 说明 |
|---|---|---|
| 1 | 金蝶Apusic分布式消息队列for RocketMQ 快速使用手册 | 简单介绍了如何快速上手使用ADMQ for RocketMQ 。 |
| 2 | 金蝶Apusic分布式消息队列for RocketMQ 安装手册 | 详细介绍如何在各操作系统上安装ADMQ for RocketMQ,以及ADMQ for RocketMQ服务启停等操作。 |
| 3 | 金蝶Apusic分布式消息队列for RocketMQ 消息引擎用户手册 | 详细介绍 ADMQ for RocketMQ 消息引擎相关功能的使用、配置、管理及配套工具的使用方法。 |
| 4 | 金蝶Apusic分布式消息队列for RocketMQ 管控台用户手册 | 详细介绍ADMQ for RocketMQ管控台相关功能的使用和操作说明。 |
| 5 | 金蝶Apusic分布式消息队列for RocketMQ 开发手册 | 详细介绍基于各开发语言进行ADMQ for RocketMQ客户端应用开发的说明。 |
| 6 | 金蝶Apusic分布式消息队列for RocketMQ 迁移手册 | 详细介绍从RocketMQ迁移到ADMQ for RocketMQ的说明。 |
| 7 | 金蝶Apusic分布式消息队列for RocketMQ 运维手册 | 详细介绍ADMQ for RocketMQ的监控、运维、安全加固等运维说明。 |
| 8 | 金蝶Apusic分布式消息队列for RocketMQ 性能优化手册 | 详细介绍ADMQ for RocketMQ性能调优的说明。 |
# 技术支持
ADMQ for RocketMQ产品提供全面的技术支持服务,您可以通过以下方式获得技术支持:
- 网址:www.apusic.com
- 电话:400-855-5800
- 邮箱:support@apusic.com
- 金蝶云社区:https://vip.kingdee.com/?productId=73&productLineId=14&lang=zh-CN
您在取得技术支持时,请提供如下信息:
您的姓名
公司信息与联系方式
操作系统及其版本
产品版本号
出现异常及错误的日志、截图等详细信息
# 快速部署
# 环境准备
| 组件 | 最低要求 |
|---|---|
| CPU | 4 核 |
| 内存 | 16 GB |
| 磁盘 | 200 GB SSD |
| 操作系统 | CentOS 7+ / Ubuntu 18.04+ |
# 一键部署
mkdir -p /apusic
cd /apusic
# 1.解压安装包
#tar zxvf ADMQ-V2.0.6-RocketMQ-构建日期.tar.gz -C /apusic/
tar zxvf ADMQ-V2.0.6.534-RocketMQ-20260624.tar.gz -C /apusic/
cd /apusic
# 建议nameserver broker组件文件夹结构分开
unzip rocketmq-V5.3.4-all.zip
mv rocketmq-V5.3.4-all nameserver
cd /apusic/nameserver
unzip rocketmq-V5.3.4-all.zip
mv rocketmq-V5.3.4-all broker
cd /apusic/broker
# 2. 启动 NameServer
bin/rocketmq start nameserver
# 3. 启动 Broker
bin/rocketmq start broker
# 4. 验证状态
bin/mqadmin clusterList -n localhost:9876
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
# 管控台部署
# 1. 解压管控台
tar zxvf admq-manager-V2.0.6.tar.gz -C /apusic/
cd /apusic/admq-manager-V2.0.6
# 2. 启动管控台
bin/admq-service start
1
2
3
4
5
6
2
3
4
5
6
# 快速验证
# 访问管控台
地址:http://<管控台IP>:12305;https://<管控台IP>:12306
默认账号:admq / 11111111
1
2
2
# 创建测试资源
创建主题
- 进入「资源管理」-「主题管理」页面
- 点击「新建」,输入主题名称(如:test-topic)
- 配置读写队列数(如:3)
- 确认创建
创建订阅组
- 进入「资源管理」-「订阅组管理」页面
- 点击「新建」,输入订阅组名称(如:group)
- 配置订阅组信息
- 确认创建
查看消息发送
进入「消息查询」页面
选择主题:test-topic
可查看该 Topic 的消息
# 发送和接收消息
# 使用命令行工具
# 发送消息
bin/mqadmin sendMessage -n localhost:9876 -t test-topic -p "Hello ADMQ for RocketMQ!"
# 查看消息
bin/mqadmin queryMsgByKey -n localhost:9876 -t test-topic -k test-key
1
2
3
4
5
2
3
4
5
# 使用管控台发送消息
- 进入「资源管理」-「Topic 管理」页面
- 选择 test-topic,点击「发送消息」
- 填写消息内容、Tag、Key 等信息
- 点击发送
# 使用开源SDK实现消息收发(Java SDK为例)
# 概览
目前RocketMQ兼容开源特性,根据底层通信协议的差异主要支持两个系列的客户端SDK,分别是Remoting协议和gRPC协议 关于 Remoting 协议 SDK 和 gRPC 协议 SDK 的对比参考如下:
| 对比项 | Remoting 协议SDK | gRPC 协议SDK |
|---|---|---|
| 多语言支持 | Java为主,其他语言为第三方仓库实现 | Java/C/C++/.NET/Go/Rust |
| 接口范围 | Producer PushConsumer PullConsumer LitePullConsumer Admin | Producer PushConsumer(仅Java) SimpleConsumer |
| 兼容版本 | 兼容4.x、5.x版本服务端 | 仅支持5.x 版本服务端 |
# 生产消息
创建消息生产者
// 实例化消息生产者Producer
DefaultMQProducer producer = new DefaultMQProducer(
namespace,
groupName,
new AclClientRPCHook(new SessionCredentials(accessKey, secretKey)) // ACL权限
);
// 设置NameServer的地址
producer.setNamesrvAddr(nameserver);
// 启动Producer实例
producer.start();
1
2
3
4
5
6
7
8
9
10
2
3
4
5
6
7
8
9
10
发送消息 发送消息由多种方式,同步发送、异步发送、单向发送等。
- 同步消息
for (int i = 0; i < 10; i++) {
// 创建消息实例,设置topic和消息内容
Message msg = new Message(topic_name, "TAG", ("Hello RocketMQ " + i).getBytes(RemotingHelper.DEFAULT_CHARSET));
// 发送消息
SendResult sendResult = producer.send(msg);
System.out.printf("%s%n", sendResult);
}
1
2
3
4
5
6
7
2
3
4
5
6
7
- 异步消息
// 设置发送失败后不重试
producer.setRetryTimesWhenSendAsyncFailed(0);
// 设置发送消息的数量
int messageCount = 10;
final CountDownLatch countDownLatch = new CountDownLatch(messageCount);
for (int i = 0; i < messageCount; i++) {
try {
final int index = i;
// 创建消息实体,设置topic和消息内容
Message msg = new Message(topic_name, "TAG", ("Hello rocketMq " + index).getBytes(RemotingHelper.DEFAULT_CHARSET));
producer.send(msg, new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
// 消息发送成功逻辑
countDownLatch.countDown();
System.out.printf("%-10d OK %s %n", index, sendResult.getMsgId());
}
@Override
public void onException(Throwable e) {
// 消息发送失败逻辑
countDownLatch.countDown();
System.out.printf("%-10d Exception %s %n", index, e);
e.printStackTrace();
}
});
} catch (Exception e) {
e.printStackTrace();
}
}
countDownLatch.await(5, TimeUnit.SECONDS);
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
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
- 单向发送
for (int i = 0; i < 10; i++) {
// 创建消息实例,设置topic和消息内容
Message msg = new Message(topic_name, "TAG", ("Hello RocketMQ " + i).getBytes(RemotingHelper.DEFAULT_CHARSET));
// 发送单向消息
producer.sendOneway(msg);
}
1
2
3
4
5
6
2
3
4
5
6
# 消费消息(Remoting协议为例)
创建消费者 支持 push 和 pull 两种消费模式
- push消费者
// 实例化消费者
DefaultMQPushConsumer pushConsumer = new DefaultMQPushConsumer(
namespace,
groupName,
new AclClientRPCHook(new SessionCredentials(accessKey, secretKey))); //ACL权限
// 设置NameServer的地址
pushConsumer.setNamesrvAddr(nameserver);
1
2
3
4
5
6
7
2
3
4
5
6
7
- pull消费者
// 实例化消费者
DefaultLitePullConsumer pullConsumer = new DefaultLitePullConsumer(
namespace,
groupName,
new AclClientRPCHook(new SessionCredentials(accessKey, secretKey)));
// 设置NameServer的地址
pullConsumer.setNamesrvAddr(nameserver);
// 设置从第一个偏移量开始消费
pullConsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
1
2
3
4
5
6
7
8
9
2
3
4
5
6
7
8
9
订阅消息 根据消费模式不同,订阅方式也有所区别。
- push订阅
// 订阅topic
pushConsumer.subscribe(topic_name, "*");
// 注册回调实现类来处理从broker拉取回来的消息
pushConsumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
// 消息处理逻辑
System.out.printf("%s Receive New Messages: %s %n", Thread.currentThread().getName(), msgs);
// 标记该消息已经被成功消费, 根据消费情况,返回处理状态
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
// 启动消费者实例
pushConsumer.start();
1
2
3
4
5
6
7
8
9
10
11
2
3
4
5
6
7
8
9
10
11
- pull订阅
// 订阅topic
pullConsumer.subscribe(topic_name, "*");
// 启动消费者实例
pullConsumer.start();
try {
System.out.printf("Consumer Started.%n");
while (true) {
// 拉取消息
List<MessageExt> messageExts = pullConsumer.poll();
System.out.printf("%s%n", messageExts);
}
} finally {
pullConsumer.shutdown();
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
2
3
4
5
6
7
8
9
10
11
12
13
14
# 客户端快速接入
# Java 客户端示例
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
public class QuickStart {
public static void main(String[] args) throws Exception {
// 生产者
DefaultMQProducer producer = new DefaultMQProducer("test-group");
producer.setNamesrvAddr("localhost:9876");
producer.start();
// 发送消息
Message msg = new Message("test-topic", "tag-a", "test-key",
"Hello ADMQ for RocketMQ!".getBytes("UTF-8"));
producer.send(msg);
System.out.println("Message sent");
producer.shutdown();
// 消费者
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("test-group");
consumer.setNamesrvAddr("localhost:9876");
consumer.subscribe("test-topic", "*");
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
ConsumeConcurrentlyContext context) {
for (MessageExt msg : msgs) {
System.out.println("Received: " + new String(msg.getBody()));
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
consumer.start();
System.out.println("Consumer started");
}
}
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
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
# Go 客户端示例
package main
import (
"context"
"fmt"
"github.com/apache/rocketmq-clients/golang/v2"
)
func main() {
// 生产者
producer, err := golang.NewProducer(
golang.WithEndpoints("localhost:9876"),
golang.WithConsumerGroup("test-group"),
)
if err != nil {
panic(err)
}
err = producer.Start()
if err != nil {
panic(err)
}
defer producer.Shutdown()
msg := golang.NewMessage("test-topic", []byte("Hello ADMQ for RocketMQ!"))
resp, err := producer.Send(context.TODO(), msg)
if err != nil {
panic(err)
}
fmt.Printf("Message sent: %s\n", resp[0].MsgId)
}
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
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
编辑页面 (opens new window)