📨 MqAbstract 基类
MqAbstract<O, D> 是所有消息中间件实现的基类。它继承自 CoreUnify<O, D> 并实现 IMq。
public abstract class MqAbstract<O, D> : CoreUnify<O, D>, IMq
where O : class
where D : class
与 DaqAbstract 的核心区别
MqAbstract 面向消息队列场景(发布/订阅),而非设备读写:
| 方面 | DaqAbstract (协议驱动) | MqAbstract (消息中间件) |
|---|---|---|
| 数据流 | Read / Write | Produce / Consume |
| 寻址方式 | Address (寄存器地址) | Topic (基于字符串) |
| 订阅方式 | SubscribeAsync / UnSubscribeAsync | ConsumeAsync / UnConsumeAsync |
| 内置 WebAPI | ✅ 支持 | ❌ 不支持 |
| 数据类型 | 强类型 (Int16, Float, Bool 等) | 原始字节或字符串 |
| 适用场景 | PLC、传感器、仪器仪表 | MQTT、Kafka、RabbitMQ 等消息队列 |
6个抽象方法 + 1个虚方法
每个中间件实现必须实现 6 个抽象方法,另有 1 个虚方法提供默认实现:
1. OnAsync — 打开连接
public abstract Task<OperateResult> OnAsync(CancellationToken token = default);
建立与消息代理/服务器的连接。对于 MQTT 客户端,此操作连接 Broker;对于 NetMQ,此操作绑定/连接 ZeroMQ 套接字。
2. OffAsync — 关闭连接
public abstract Task<OperateResult> OffAsync(bool hardClose = false, CancellationToken token = default);
关闭与消息代理的连接。hardClose = true 执行强制关闭;false 尝试优雅关闭。
3. ProduceAsync (字节) — 发布原始字节 ⭐ 抽象
public abstract Task<OperateResult> ProduceAsync(string topic, byte[] content, CancellationToken token = default);
这是子类必须实现的核心发布方法。 将字节数据发布到指定主题。
4. ProduceAsync (字符串) — 发布字符串 🔧 虚方法
public virtual Task<OperateResult> ProduceAsync(string topic, string content, Encoding? encoding = null, CancellationToken token = default);
字符串发布的便捷虚方法。默认实现为:将字符串按指定编码(默认 UTF8)转为字节数组后,委托给抽象的 ProduceAsync(string, byte[], CancellationToken)。
子类可以重写此方法以实现自定义字符串处理逻辑。
5. ConsumeAsync — 开始消费
public abstract Task<OperateResult> ConsumeAsync(string topic, CancellationToken token = default);
订阅指定主题并开始接收消息。消息通过 OnDataEvent / OnDataEventAsync 事件分发。
6. UnConsumeAsync — 停止消费
public abstract Task<OperateResult> UnConsumeAsync(string topic, CancellationToken token = default);
取消订阅指定主题,停止接收消息。
7. GetStatusAsync — 连接状态
public abstract Task<OperateResult> GetStatusAsync(CancellationToken token = default);
返回当前与消息代理的连接状态。
同步便捷方法
所有异步方法均有对应的同步包装(来自 MqAbstract 实现,非抽象):
| 同步方法 | 委托的异步方法 |
|---|---|
On() |
OnAsync().GetAwaiter().GetResult() |
Off(bool hardClose = false) |
OffAsync(hardClose).GetAwaiter().GetResult() |
Produce(string topic, string content, Encoding? encoding = null) |
ProduceAsync(topic, content, encoding).GetAwaiter().GetResult() |
Produce(string topic, byte[] content) |
ProduceAsync(topic, content).GetAwaiter().GetResult() |
Consume(string topic) |
ConsumeAsync(topic).GetAwaiter().GetResult() |
UnConsume(string topic) |
UnConsumeAsync(topic).GetAwaiter().GetResult() |
GetStatus() |
GetStatusAsync().GetAwaiter().GetResult() |
继承自 CoreUnify 的能力
通过 CoreUnify<O, D> 继承获得:
| 能力 | 方法/属性 |
|---|---|
| 单例管理 | InstanceAsync(D?) / Instance(D?) — 最多 255 实例 |
| 事件系统 | OnDataEvent / OnDataEventAsync、OnInfoEvent / OnInfoEventAsync、OnLanguageEvent / OnLanguageEventAsync |
| 日志 | LogOperateSet(...) / LogOperateGet() — 基于 Serilog |
| 多语言 | GetLanguage() / SetLanguage(LanguageType) / GetLanguageValue(key, model) |
| 参数反射 | GetParam(bool) / GetAutoAllocatingParam() |
| 实例创建 | CreateInstance<T>(param) / CreateInstance(string json) |
Dispose 行为
// 同步释放:先强制关闭连接,再释放单例
public override void Dispose()
{
Off(true);
base.Dispose();
}
// 异步释放
public override async ValueTask DisposeAsync()
{
await OffAsync(true).ConfigureAwait(false);
await base.DisposeAsync().ConfigureAwait(false);
}
具体实现
所有中间件操作类均继承自 MqAbstract<O, D>:
| 类 | 包 | 协议 | 特点 |
|---|---|---|---|
MqttClientOperate |
Snet.Mqtt | MQTT 3.1.1/5.0 | QoS 0/1/2、Retain、Message Expiry |
KafkaOperate |
Snet.Kafka | Apache Kafka | SASL 认证、AdminClient、多 Topic 并发 |
RabbitMQOperate |
Snet.RabbitMQ | AMQP 0-9-1 | 4 种交换机类型、持久化、手动 ACK |
NetMQOperate |
Snet.NetMQ | ZeroMQ | 无 Broker、<1ms 延迟、Pub/Sub 模式 |
NettyClientOperate |
Snet.Netty | 原始 TCP (DotNetty) | 自定义帧协议、SSL/TLS |
注意:
MqttServiceOperate、MqttWebSocketServiceOperate和NettyServiceOperate是服务端类,直接继承CoreUnify而非MqAbstract。它们不是客户端消费者/生产者,因此不使用IMq接口。
使用示例
// 以 MQTT 客户端为例 — 所有 MqAbstract 子类 API 完全一致
using Snet.Mqtt.client;
var mqtt = new MqttClientOperate(new MqttClientData.Basics
{
IpAddress = "broker.emqx.io",
Port = 1883,
UserName = "shunnet",
Password = "shunnet"
});
// 1. 连接
await mqtt.OnAsync();
// 2. 绑定事件再消费
mqtt.OnDataEventAsync += async (sender, e) =>
{
if (e.Status)
Console.WriteLine($"收到: {e.ResultData}");
};
// 3. 订阅主题
await mqtt.ConsumeAsync("sensors/temperature");
// 4. 发布消息(字符串 — 虚方法,默认 UTF8 编码)
await mqtt.ProduceAsync("sensors/temperature", "25.3 °C");
// 5. 发布消息(字节 — 抽象方法)
byte[] payload = BitConverter.GetBytes(25.3f);
await mqtt.ProduceAsync("sensors/temperature", payload);
// 6. 取消订阅
await mqtt.UnConsumeAsync("sensors/temperature");
// 7. 断开连接
await mqtt.OffAsync();
继承链
CoreUnify<O, D> — 单例、事件、计时、日志、语言
└── MqAbstract<O, D> — Produce/Consume(消息中间件)
├── MqttClientOperate — MQTT 客户端
├── KafkaOperate — Kafka 客户端
├── RabbitMQOperate — RabbitMQ 客户端
├── NetMQOperate — NetMQ/ZeroMQ 客户端
└── NettyClientOperate — Netty TCP 客户端
