开发手册
# 前言
本文档为金蝶Apusic分布式消息队列for RocketMQ(Apusic Distributed Message Queue for RocketMQ,简称:ADMQ for RocketMQ)的开发使用说明,帮助用户快速学习如何使用金蝶Apusic分布式消息队列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
您在取得技术支持时,请提供如下信息:
您的姓名
公司信息与联系方式
操作系统及其版本
产品版本号
出现异常及错误的日志、截图等详细信息
# 安装部署
# 部署启动
解压安装包
mkdir -p /apusic
cd /apusic
# 解压安装包
#tar zxvf ADMQ-V2.0.6-RocketMQ-构建日期.tar.gz -C /apusic/
tar zxvf ADMQ-V2.0.6.534-RocketMQ-20260625.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
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
修改配置,修改对应配置文件的地址和端口
vi conf/nameserver.properties
vi conf/broker.conf
1
2
2
启动
bin/rocketmq start nameserver
bin/rocketmq start broker
1
2
2
# 连接参数
| 参数 | 说明 | 示例 |
|---|---|---|
| namesrvAddr | NameServer 地址 | localhost:9876 |
| groupName | 生产者/消费者组名 | test-group |
| accessKey | ACL 访问密钥(可选) | test-user |
| secretKey | ACL 访问密钥(可选) | test-password |
# 客户端开发
# 客户端依赖
# Java
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-client</artifactId>
<version>4.9.8</version> (如果是5.x版本,可以用5.1.4)
</dependency>
1
2
3
4
5
2
3
4
5
# Go
go get github.com/apache/rocketmq-clients/golang/v2
1
# Python
pip install rocketmq-client-python
1
# Java 开发示例
# 生产者
# 同步发送
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
public class SyncProducer {
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"));
SendResult result = producer.send(msg);
System.out.println("SendResult: " + result);
producer.shutdown();
}
}
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
# 异步发送
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
public class AsyncProducer {
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, new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
System.out.println("Send success: " + sendResult);
}
@Override
public void onException(Throwable e) {
System.err.println("Send failed: " + e.getMessage());
}
});
Thread.sleep(5000);
producer.shutdown();
}
}
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
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
# 单向发送
Message msg = new Message("test-topic", "tag-a", "One-way message".getBytes());
producer.sendOneway(msg);
1
2
2
# 消费者
# push模式消费
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.MessageExt;
import java.util.List;
public class PushConsumer {
public static void main(String[] args) throws Exception {
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
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
# pull模式消费
import org.apache.rocketmq.client.consumer.DefaultLitePullConsumer;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
public class PullConsumer {
public static void main(String[] args) throws Exception {
DefaultLitePullConsumer consumer = new DefaultLitePullConsumer("test-group");
consumer.setNamesrvAddr("localhost:9876");
consumer.subscribe("test-topic", "*");
consumer.start();
while (true) {
List<MessageExt> msgs = consumer.poll();
for (MessageExt msg : msgs) {
System.out.println("Received: " + new String(msg.getBody()));
}
}
}
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
# 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!"))
msg.SetTag("tag-a")
msg.SetKeys("test-key")
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
32
33
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
# 消费者
package main
import (
"context"
"fmt"
"github.com/apache/rocketmq-clients/golang/v2"
)
func main() {
consumer, err := golang.NewSimpleConsumer(
golang.WithEndpoints("localhost:9876"),
golang.WithConsumerGroup("test-group"),
golang.WithSubscriptionExpressions(map[string]*golang.FilterExpression{
"test-topic": golang.SUB_ALL,
}),
)
if err != nil {
panic(err)
}
err = consumer.Start()
if err != nil {
panic(err)
}
defer consumer.Shutdown()
for {
msgs, err := consumer.Receive(context.TODO(), 1, 30)
if err != nil {
fmt.Printf("Receive error: %v\n", err)
continue
}
for _, msg := range msgs {
fmt.Printf("Received: %s\n", string(msg.GetBody()))
err = consumer.Ack(context.TODO(), msg)
if err != nil {
fmt.Printf("Ack error: %v\n", err)
}
}
}
}
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
# Python 开发示例
# 生产者
from rocketmq.client import Producer, Message
producer = Producer('test-group')
producer.set_name_server_address('localhost:9876')
producer.start()
msg = Message('test-topic')
msg.set_tags('tag-a')
msg.set_keys('test-key')
msg.set_body('Hello ADMQ for RocketMQ!')
result = producer.send_sync(msg)
print('Message sent: {result.msg_id}')
producer.shutdown()
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
# 消费者
from rocketmq.client import PushConsumer, ConsumeStatus
def callback(msg):
print(f'Received: {msg.body.decode("utf-8")}')
return ConsumeStatus.CONSUME_SUCCESS
consumer = PushConsumer('test-group')
consumer.set_name_server_address('localhost:9876')
consumer.subscribe('test-topic', '*', callback)
consumer.start()
print('Consumer started')
import time
time.sleep(60)
consumer.shutdown()
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
# 高级特性
# 顺序消息
// 发送顺序消息
Message msg = new Message("order-topic", "tag-a", "order-1", "content".getBytes());
SendResult result = producer.send(msg, new MessageQueueSelector() {
@Override
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
Long orderId = (Long) arg;
long index = orderId % mqs.size();
return mqs.get((int) index);
}
}, orderId);
// 消费顺序消息
consumer.registerMessageListener(new MessageListenerOrderly() {
@Override
public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs, ConsumeOrderlyContext context) {
for (MessageExt msg : msgs) {
System.out.println("Received: " + new String(msg.getBody()));
}
return ConsumeOrderlyStatus.SUCCESS;
}
});
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
# 事务消息
TransactionMQProducer producer = new TransactionMQProducer("tx-group");
producer.setNamesrvAddr("localhost:9876");
producer.setTransactionListener(new TransactionListener() {
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
// 执行本地事务
executeLocalBusiness(msg);
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// 回查本地事务状态
if (isTransactionSuccess(msg)) {
return LocalTransactionState.COMMIT_MESSAGE;
}
return LocalTransactionState.ROLLBACK_MESSAGE;
}
});
producer.start();
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
# 延时消息
Message msg = new Message("delay-topic", "tag-a", "content".getBytes());
msg.setDelayTimeLevel(3); // 延时 10 秒
producer.send(msg);
1
2
3
2
3
# 批量消费
consumer.setConsumeMessageBatchMaxSize(10); // 每次消费 10 条消息
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
ConsumeConcurrentlyContext context) {
// 批量处理消息
for (MessageExt msg : msgs) {
processMessage(msg);
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
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)