MqAbstract 基类 - Snet Docs

📨 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 / OnDataEventAsyncOnInfoEvent / OnInfoEventAsyncOnLanguageEvent / 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

注意MqttServiceOperateMqttWebSocketServiceOperateNettyServiceOperate 是服务端类,直接继承 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 客户端

参见