Snet 框架 -- 中间件概览
最后更新: 2026-07-21
入口文件: Snet.Core/abstract/MqAbstract.cs, Snet.Core/extend/CoreUnify.cs
目标框架: .NET 8 / .NET 9 / .NET 10
架构
+----------------------------+
| 应用层 |
+-----------+----------------+
|
+-----------v----------------+
| Snet.Core (抽象层) |
| MqAbstract<O, D> | -- 客户端中间件
| CoreUnify<O, D> | -- 服务端 / 代理端
+-----------+----------------+
|
+------------+------------+------------+------------+------------+
| | | | | |
+----v---+ +-----v----+ +-----v----+ +-----v----+ +-----v----+ +-----v----+
| MQTT | | Kafka | | RabbitMQ | | NetMQ | | Netty | | TEP |
| 客户端 | | | | | | (ZeroMQ) | | TCP+Sys | | 主端 |
| 代理端 | | | | | | 发布/订阅 | | 客户端 | | +从端 |
| +WS | | | | | | | | +服务端 | | |
+--------+ +----------+ +----------+ +----------+ +----------+ +----------+
中间件组件分为两类:
| 类别 | 基类 | 用途 | 示例 |
|---|---|---|---|
| 客户端 | MqAbstract<O, D> |
连接、生产、消费 | MQTT 客户端, Kafka, RabbitMQ, NetMQ, Netty 客户端 |
| 服务端 | CoreUnify<O, D> |
托管服务、接受连接 | MQTT 代理, Netty 服务端, TEP 从端 |
通用消息队列模式 (所有客户端中间件)
每个消息队列客户端遵循完全一致的架构:
namespace Snet.{Middleware}
{
// 1. 配置
public class {Middleware}Data
{
public class Basics
{
public string SN { get; set; }
public string IpAddress { get; set; } // (Kafka 使用 BootstrapServers)
public int Port { get; set; }
public ResponseType ResponseType { get; set; }
// ... 中间件特定字段
}
}
// 2. 实现
public class {Middleware}Operate : MqAbstract<{Middleware}Operate, {Middleware}Data.Basics>, IMq
{
public static async Task<{Middleware}Operate> InstanceAsync(Basics config, CancellationToken token)
// 必须实现的 6 个方法:
public override async Task<OperateResult> OnAsync(CancellationToken token)
public override async Task<OperateResult> OffAsync(bool hardClose, CancellationToken token)
public override async Task<OperateResult> GetStatusAsync(CancellationToken token)
public override async Task<OperateResult> ProduceAsync(string topic, byte[] content, CancellationToken token)
public override async Task<OperateResult> ConsumeAsync(string topic, CancellationToken token)
public override async Task<OperateResult> UnConsumeAsync(string topic, CancellationToken token)
}
}
通过 MqAbstract 继承 CoreUnify:获得单例管理、事件总线(OnDataEvent / OnInfoEvent)、日志、多语言。
统一用法 (所有客户端中间件)
// 第一步: 配置
var config = new MqttClientData.Basics
{
IpAddress = "localhost",
Port = 1883,
UserName = "admin",
Password = "admin"
};
// 第二步: 实例化
var client = await MqttClientOperate.InstanceAsync(config);
// 第三步: 打开连接
await client.OnAsync();
// 第四步: 订阅 / 消费
await client.ConsumeAsync("sensors/temperature");
client.OnDataEvent += (sender, e) =>
{
Console.WriteLine($"收到消息: {e.ResultData}");
};
// 第五步: 生产 / 发送
await client.ProduceAsync("sensors/temperature", "25.5°C");
// -- 或作为字节数组 --
await client.ProduceAsync("sensors/temperature", new byte[] { 0x01, 0x02 });
// 第六步: 清理资源
await client.UnConsumeAsync("sensors/temperature");
await client.OffAsync();
将
MqttClient替换为任意中间件名称(如Kafka、RabbitMQ、NetMQ、NettyClient)-- API 接口完全一致。
完整中间件目录
MQTT -- Snet.Mqtt
基于 MQTTnet,支持 MQTT 3.1.1 和 MQTT 5.0。
| 组件 | 基类 | 角色 | 主要特性 |
|---|---|---|---|
MqttClientOperate |
MqAbstract<MqttClientOperate, MqttClientData.Basics> |
MQTT 客户端 | QoS 0/1/2, 用户名/密码认证, ClientId, KeepAlive |
MqttServiceOperate |
CoreUnify<MqttServiceOperate, MqttServiceData.Basics> |
嵌入式 MQTT 代理 | 独立代理,接受客户端连接 |
MqttWebSocketServiceOperate |
CoreUnify<MqttWebSocketServiceOperate, MqttWebSocketServiceData.Basics> |
双协议 MQTT + WebSocket 代理 | 基于 Kestrel 运行,同一端口支持 MQTT 和 WS |
配置 (客户端):
var config = new MqttClientData.Basics
{
IpAddress = "broker.emqx.io",
Port = 1883,
UserName = "shunnet",
Password = "shunnet",
ClientID = "snet-client-001", // 不输入则自动生成
MessageExpirationTime = 86400000, // 毫秒 (24h)
QualityOfServiceLevel = MqttQualityOfServiceLevel.AtMostOnce
};
嵌入式代理:
var broker = new MqttServiceData.Basics { Port = 1883 };
var service = await MqttServiceOperate.InstanceAsync(broker);
await service.OnAsync();
// MQTT 代理现已运行在端口 1883
Kafka -- Snet.Kafka
基于 Confluent.Kafka,支持完整的 Kafka 生产者/消费者模型。
| 组件 | 基类 | 主要特性 |
|---|---|---|
KafkaOperate |
MqAbstract<KafkaOperate, KafkaData.Basics> |
生产者 + 消费者, SASL 认证, 自动偏移管理 |
配置:
var config = new KafkaData.Basics
{
BootstrapServers = "kafka-broker-1:9092,kafka-broker-2:9092",
SecurityProtocol = SecurityProtocol.SaslPlaintext,
SaslMechanism = SaslMechanism.Plain,
AutoOffsetReset = AutoOffsetReset.Latest,
ResponseType = ResponseType.Content
};
| 字段 | 选项 | 说明 |
|---|---|---|
SecurityProtocol |
Plaintext / SaslPlaintext / SaslSsl / Ssl |
传输安全级别 |
SaslMechanism |
Gssapi / Plain / ScramSha256 / ScramSha512 / OAuthBearer |
SASL 认证方式 |
AutoOffsetReset |
Latest / Earliest / Error |
消费起始位置 |
SaslKerberosServiceName |
字符串 (默认: "snet") |
使用 GSSAPI 时必需 |
RabbitMQ -- Snet.RabbitMQ
基于 RabbitMQ.Client,支持所有交换机类型和队列模式。
| 组件 | 基类 | 主要特性 |
|---|---|---|
RabbitMQOperate |
MqAbstract<RabbitMQOperate, RabbitMQData.Basics> |
基于交换机的发布/订阅, 消息 TTL, 用户/密码认证 |
配置:
var config = new RabbitMQData.Basics
{
ExChangeName = "snet-exchange",
IpAddress = "localhost",
Port = 5672,
UserName = "shunnet",
Password = "shunnet",
MessageExpirationTime = 86400000, // 24 小时 (毫秒)
ResponseType = ResponseType.Content
};
交换机类型: direct / fanout / headers / topic
ConsumeAsync(topic) 调用将以 topic 作为路由键,将队列绑定到配置的交换机。
NetMQ -- Snet.NetMQ
100% 纯 C# ZeroMQ 实现(无原生依赖)。发布/订阅模式。
| 组件 | 基类 | 主要特性 |
|---|---|---|
NetMQOperate |
MqAbstract<NetMQOperate, NetMQData.Basics> |
ZeroMQ 发布/订阅,零拷贝,高吞吐量 |
配置:
var pub = new NetMQData.Basics
{
UModel = UseModel.PubModel, // PubModel (绑定) 或 SubModel (连接)
Address = "tcp://127.0.0.1:8866", // ZeroMQ 地址格式
TimeOut = 1000 // 毫秒
};
var sub = new NetMQData.Basics
{
UModel = UseModel.SubModel,
Address = "tcp://127.0.0.1:8866"
};
UseModel |
行为 |
|---|---|
PubModel |
绑定 PUB 套接字 -- 一个发布者,多个订阅者 |
SubModel |
连接 SUB 套接字 -- 订阅某个发布者 |
Netty -- Snet.Netty
基于 DotNetty(Netty 的 C# 移植)。TCP 客户端/服务端,支持可选的 SSL/TLS。
| 组件 | 基类 | 主要特性 |
|---|---|---|
NettyClientOperate |
MqAbstract<NettyClientOperate, NettyClientData.Basics> |
TCP 客户端, SSL/TLS, 基于任务的管道 |
NettyServiceOperate |
CoreUnify<NettyServiceOperate, NettyServiceData.Basics> |
TCP 服务端, SSL/TLS, 多客户端支持 |
客户端配置:
var client = new NettyClientData.Basics
{
IpAddress = "localhost",
Port = 8899,
SslFilePath = "/path/to/cert.pfx", // 可选 SSL 证书
SslFilePassword = "cert-password", // 可选证书密码
TaskNumber = 5
};
服务端配置:
var server = new NettyServiceData.Basics { Port = 8899 };
var service = await NettyServiceOperate.InstanceAsync(server);
await service.OnAsync();
TEP -- Snet.TEP (TCP 扩展插件)
自定义二进制协议,12 字节帧标记。双角色设计(主端/从端)。
| 组件 | 基类 | 角色 | 主要特性 |
|---|---|---|---|
TepMasterOperate |
DaqAbstract<TepMasterOperate, TepMasterData.Basics> |
连接管理器 + 数据核心 | 管理从端连接,路由数据 |
TepSlaveOperate |
CoreUnify<TepSlaveOperate, TepSlaveData.Basics> |
从端客户端 | 连接主端,通过 设备名+用户名/密码 认证 |
主端配置:
var master = new TepMasterData.Basics
{
Port = 10086
// 管理来自多个从端的连接
};
var masterOperate = await TepMasterOperate.InstanceAsync(master);
await masterOperate.OnAsync();
// 通过统一的 Address 模型对已注册的从端进行读写操作
从端配置:
var slave = new TepSlaveData.Basics
{
IpAddress = "192.168.1.10",
Port = 10086,
DevName = "snet",
DevID = "10001",
UserName = "shunnet",
Password = "shunnet"
};
var slaveOperate = await TepSlaveOperate.InstanceAsync(slave);
await slaveOperate.OnAsync();
TEP 主端使用
DaqAbstract(而非MqAbstract),因为它实现了完整的 Read/Write/Subscribe 驱动契约,将已连接的从端视为数据源。
对比: DaqAbstract vs MqAbstract
| 方面 | DaqAbstract (协议驱动) | MqAbstract (消息中间件) |
|---|---|---|
| 数据流 | 读取 / 写入 | 生产 / 消费 |
| 寻址方式 | Address (寄存器、线圈、标签) | Topic (基于字符串) |
| 订阅方式 | SubscribeAsync / UnSubscribeAsync | ConsumeAsync / UnConsumeAsync |
| 内置 WebAPI | 支持 | 不支持 |
| 数据类型 | 强类型 (Int16, Float, Bool 等) | 原始字节或字符串 |
| 适用场景 | PLC、传感器、仪器仪表 | 消息队列、事件总线 |
日志系统 -- Snet.Log
基于 Serilog,支持自动文件组织和清理。
| 特性 | 详情 |
|---|---|
| 引擎 | Serilog (结构化日志) |
| 等级 | Verbose(详细), Debug(调试), Information(信息), Warning(警告), Error(错误), Fatal(致命) |
| 文件滚动 | 按小时滚动,按日期文件夹组织 |
| 自动清理 | 默认保留 30 天,后台任务执行 |
| 文件命名 | 帕斯卡命名转为短横线命名 (如 ModbusOperate 转为 modbus-operate) |
| 线程安全 | 使用 ConcurrentDictionary 缓存 Logger 实例 |
// 所有驱动和中间件均通过 CoreUnify 继承日志能力
// 日志文件自动写入:
// logs/{ClassName}/{yyyy-MM-dd}/{HH}.log
数据库访问 -- Snet.DB
双处理器架构,提供最大灵活性。
| 特性 | 详情 |
|---|---|
| 支持数据库 | SqlServer(SQL Server), MySql, Oracle, SQLite |
| 双处理器 | SqlSugar (完整 CRUD) + Dapper (只读 + 订阅) |
| 查询语法 | LINQ 表达式树 |
| 基类 | DBData.Basics : SubscribeData.SCData -- 支持相同的订阅引擎 |
配置:
var db = new DBData.Basics
{
ConnectStr = "Data Source=app.db", // SQLite 连接字符串
DBType = DBType.SQLite, // 数据库类型
HandlerType = DBHandlerType.Daq // Daq (完整) 或 Dapper (只读)
};
var dbOperate = await DBOperate.InstanceAsync(db);
await dbOperate.OnAsync();
// LINQ 查询
var result = await dbOperate.QueryAsync<UserEntity>(x => x.age > 18 && x.name == "test");
| 数据库类型 | 默认端口 | 连接字符串示例 |
|---|---|---|
SQLite |
本地文件 | Data Source=app.db |
SqlServer |
1433 | Server=.;Database=snet;User Id=sa;Password=xxx; |
MySql |
3306 | Server=localhost;Database=snet;Uid=root;Pwd=xxx; |
Oracle |
1521 | Data Source=localhost/snet;User Id=system;Password=xxx; |
传输层架构
+------------------+ +------------------+ +------------------+
| MQTT 客户端 | | Kafka 生产者 | | RabbitMQ 客户端 |
| (MQTTnet) | | (Confluent) | | (RabbitMQ.Client)|
+--------+---------+ +--------+---------+ +--------+---------+
| | |
+--------------------+------------------------+
|
+---------v----------+
| MqAbstract<O, D> |
| CoreUnify<O, D> |
| (单例 + 事件 + 日志)|
+---------+----------+
|
+---------v----------+
| OnDataEvent |
| (统一输出) |
+--------------------+
相关文档
- 协议驱动概览 -- 33 种 PLC/传感器驱动
- 协议驱动概览 (中文)
- Snet.Core 源码:
Snet.Core/abstract/MqAbstract.cs - 接口定义:
Snet.Model/interface/IMq.cs,Snet.Model/interface/IProducer.cs,Snet.Model/interface/IConsumer.cs
