概览 - Snet Docs

Snet 框架 -- 中间件概览

入口文件: Snet.Core/abstract/MqAbstract.cs, Snet.Core/extend/CoreUnify.cs 目标框架: .NET 8 / .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, RocketMQ, 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, 60 秒 KeepAlive 心跳, 自动重连(5 次,5–30 秒退避)
MqttServiceOperate CoreUnify<MqttServiceOperate, MqttServiceData.Basics> 嵌入式 MQTT 代理 独立代理,接受客户端连接
MqttWebSocketServiceOperate CoreUnify<MqttWebSocketServiceOperate, MqttWebSocketServiceData.Basics> 双协议 MQTT + WebSocket 代理 基于 Kestrel 运行,同进程托管;MQTT(Port=6688)与 WebSocket(WsPort=8866)各占一个端口

配置 (客户端):

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.Plaintext,
    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 时必需
说明

仅 GSSAPI/Kerberos 实际可用(Plain/SCRAM 在现有 Basics 中无凭据注入路径——没有 SASL 用户名/密码字段)。


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 作为路由键,将队列绑定到配置的交换机。


RocketMQ -- Snet.RocketMQ

基于 RocketMQ.Client(Apache RocketMQ 5.x),客户端经 gRPC 连接 RocketMQ Proxy(默认端口 8081)。

组件 基类 主要特性
RocketMQOperate MqAbstract<RocketMQOperate, RocketMQData.Basics> 生产者 + 消费者,消费者懒创建,ACL 认证,SSL 开关

配置:

var config = new RocketMQData.Basics
{
    IpAddress = "127.0.0.1",
    Port = 8081,          // RocketMQ 5.x Proxy gRPC 端口
    ConsumerGroup = "snet",
    ResponseType = ResponseType.Content
};

完整指南见 RocketMQ


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 扩展插件)

TEP 为独立文档分组,详见 TEP 协议概览

自定义二进制协议,4 字节包头/包尾帧标记(头部共 11 字节:头+命令+方向+长度)。双角色设计(主端/从端)。

组件 基类 角色 主要特性
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
双处理器 SqlSugarCore (完整 CRUD) + Dapper (采集 + 订阅)
查询语法 LINQ 表达式树
基类 DBData.Basics : SubscribeData.SCData -- 支持相同的订阅引擎

配置:

using Snet.DB;

var db = new DBData.Basics
{
    ConnectStr = "Data Source=app.db",  // SQLite 连接字符串
    DBType = DBData.DBType.SQLite,            // 数据库类型
    HandlerType = DBData.DBHandlerType.Default // Default = SqlSugarCore 增删改查;Daq = Dapper 采集
};

var dbOperate = await DBOperate.InstanceAsync(db);
await dbOperate.OnAsync();

// UserEntity 是你自己的 SqlSugar 实体类型
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      |
                    |   (统一输出)        |
                    +--------------------+

相关文档

  • 协议驱动概览 -- 34 种 PLC/传感器驱动
  • Snet.Core 源码: Snet.Core/abstract/MqAbstract.cs
  • 接口定义: Snet.Model/interface/IMq.cs, Snet.Model/interface/IProducer.cs, Snet.Model/interface/IConsumer.cs