RocketMQ - Snet Docs

Apache RocketMQ 中间件

包: Snet.RocketMQ | 类: RocketMQOperate | 基类: MqAbstract<RocketMQOperate, RocketMQData.Basics>

基于 Apache RocketMQ 5.x 的高吞吐分布式发布/订阅消息中间件。客户端经 gRPC 协议连接 RocketMQ Proxy(默认端口 8081,而非 nameserver),支持主题动态订阅、多主题并发消费、ACL 认证与 SSL 开关。

概览

RocketMQOperate 通过 IMq(IProducer + IConsumer)提供消息的生产与消费。继承自 MqAbstract<RocketMQOperate, RocketMQData.Basics>,基于官方 RocketMQ.Client SDK 实现。

安装

dotnet add package Snet.RocketMQ

快速开始

using Snet.RocketMQ;

var rmq = await RocketMQOperate.InstanceAsync(new RocketMQData.Basics
{
    IpAddress = "127.0.0.1",
    Port = 8081,        // RocketMQ 5.x Proxy gRPC 端口
    ConsumerGroup = "snet"
});

await rmq.OnAsync();

// 先注册事件处理器,再消费——接收到的消息在这里
rmq.OnDataEventAsync += async (sender, e) =>
{
    if (e.Status)
        Console.WriteLine($"收到数据: {e.ResultData}");
    else
        Console.WriteLine($"消费失败: {e.Message}");
};

await rmq.ConsumeAsync("sensor-data");   // 首次订阅会懒创建消费者

await rmq.ProduceAsync("sensor-data", System.Text.Encoding.UTF8.GetBytes("hello rocketmq"));

// 后续:停止消费
// await rmq.UnConsumeAsync("sensor-data");
// await rmq.OffAsync();
// await rmq.DisposeAsync();

配置

参数 类型 默认值 说明
SN string 序列号 / 设备标识
IpAddress string "127.0.0.1" Broker / Proxy 地址
Port int 6688 连接端口——RocketMQ 5.x 客户端经 gRPC 连接 Proxy(默认 8081);标准 5.x 部署请设 8081
AccessKey string null ACL 认证 AccessKey(须与 SecretKey 成对配置)
SecretKey string null ACL 认证 SecretKey(须与 AccessKey 成对配置)
ConsumerGroup string "snet" 消费组
SslEnabled bool false 启用 SSL——RocketMQ 客户端默认开启 TLS;本地无 TLS 的 Broker 请关闭
ResponseType enum Content 接收数据的格式:Bytes / Content / ContentWithTopic(见下)
ACL

AccessKeySecretKey 必须同时配置——只配其中一个时 OnAsync 返回失败「AccessKey 与 SecretKey 需同时配置」。

SSL: RocketMQ.Client SDK 默认开启 TLS;SslEnabled = false(默认值)即关闭。仅在 Proxy 提供 TLS 时开启。

响应类型

ResponseType 决定消费消息通过 OnDataEventAsync 下发时的格式:

e.ResultData 说明
Bytes byte[] 原始消息体
Content string 消息体按 UTF-8 解码的文本
ContentWithTopic ResponseModel ResponseModel 反序列化——同时携带主题与内容

支持的操作

操作 方法 说明
连接 OnAsync() 构建生产者客户端;校验 ACL 凭据是否成对
断开 OffAsync(bool hardClose = false) 先关闭入口,摘除消费者/生产者引用,出锁后 Dispose
状态 GetStatusAsync() 本地标志:「已连接」/「关闭中」/「未连接」
发布 ProduceAsync(topic, byte[] content, token) 发送消息;消息 Key 设为 topic,同主题消息落入同一分区(保序)。消息 ID 在 OperateResult.ResultDataISendReceipt
订阅 ConsumeAsync(topic, token) 订阅主题。消费者在首次调用时懒创建(SDK 的 PushConsumer.Builder 要求非空初始订阅),以首个主题初始化;后续调用动态追加主题
取消订阅 UnConsumeAsync(topic, token) 移除主题订阅;最后一个主题移除后 Dispose 消费者。未订阅过的主题返回失败「此主题不存在」

消费行为

  • 消费者懒创建: 首次 ConsumeAsync 以该主题为初始订阅构建 PushConsumer;后续主题通过 Subscribe(topic, FilterExpression.SubAll) 动态添加。
  • 重复订阅: 已订阅的主题再次订阅返回失败「{topic} 此主题已订阅」。
  • 消息推送: 收到消息后经 OnDataEventAsync fire-and-forget 推送(不阻塞 SDK 消费线程)。事件推送抛异常时返回 FAILURE,RocketMQ 按重试策略重新投递。
  • 并发安全: 主题/订阅状态由异步锁保护;SDK 对象的 DisposeAsync(长网络等待)在锁外执行,不阻塞其他操作。
  • 消息顺序: ProduceAsync 将消息 Key 设为 topic,同主题消息路由到同一分区,保持顺序。

事件

事件 签名 说明
OnDataEvent EventHandler<EventDataResult> 收到消费消息时触发
OnDataEventAsync EventHandlerAsync<EventDataResult> OnDataEvent 的异步版本
OnInfoEvent EventHandler<EventInfoResult> 信息与状态消息
OnInfoEventAsync EventHandlerAsync<EventInfoResult> OnInfoEvent 的异步版本

相关文档