- Middleware - Overview - Snet Docs

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 替换为任意中间件名称(如 KafkaRabbitMQNetMQNettyClient)-- 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