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
AccessKey 与 SecretKey 必须同时配置——只配其中一个时 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.ResultData(ISendReceipt) |
| 订阅 | ConsumeAsync(topic, token) |
订阅主题。消费者在首次调用时懒创建(SDK 的 PushConsumer.Builder 要求非空初始订阅),以首个主题初始化;后续调用动态追加主题 |
| 取消订阅 | UnConsumeAsync(topic, token) |
移除主题订阅;最后一个主题移除后 Dispose 消费者。未订阅过的主题返回失败「此主题不存在」 |
消费行为
- 消费者懒创建: 首次
ConsumeAsync以该主题为初始订阅构建PushConsumer;后续主题通过Subscribe(topic, FilterExpression.SubAll)动态添加。 - 重复订阅: 已订阅的主题再次订阅返回失败「
{topic}此主题已订阅」。 - 消息推送: 收到消息后经
OnDataEventAsyncfire-and-forget 推送(不阻塞 SDK 消费线程)。事件推送抛异常时返回FAILURE,RocketMQ 按重试策略重新投递。 - 并发安全: 主题/订阅状态由异步锁保护;SDK 对象的
DisposeAsync(长网络等待)在锁外执行,不阻塞其他操作。 - 消息顺序:
ProduceAsync将消息 Key 设为 topic,同主题消息路由到同一分区,保持顺序。
事件
| 事件 | 签名 | 说明 |
|---|---|---|
OnDataEvent |
EventHandler<EventDataResult> |
收到消费消息时触发 |
OnDataEventAsync |
EventHandlerAsync<EventDataResult> |
OnDataEvent 的异步版本 |
OnInfoEvent |
EventHandler<EventInfoResult> |
信息与状态消息 |
OnInfoEventAsync |
EventHandlerAsync<EventInfoResult> |
OnInfoEvent 的异步版本 |
