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 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

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

  1. 您的姓名

  2. 公司信息与联系方式

  3. 操作系统及其版本

  4. 产品版本号

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

# 安装部署

# 部署启动

解压安装包

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

修改配置,修改对应配置文件的地址和端口

vi conf/nameserver.properties
vi conf/broker.conf
1
2

启动

bin/rocketmq start nameserver
bin/rocketmq start broker
1
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

# 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

# 异步发送

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

# 单向发送

Message msg = new Message("test-topic", "tag-a", "One-way message".getBytes());
producer.sendOneway(msg);
1
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

# 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

# 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

# 消费者

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

# 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

# 消费者

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

# 高级特性

# 顺序消息

// 发送顺序消息
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

# 事务消息

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

# 延时消息

Message msg = new Message("delay-topic", "tag-a", "content".getBytes());
msg.setDelayTimeLevel(3);  // 延时 10 秒
producer.send(msg);
1
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
编辑页面 (opens new window)
#开发手册

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

  • 浅色模式