📘 数据库协议驱动
命名空间: Snet.DB | 类: DBOperate | 基类: DaqAbstract<DBOperate, DBData.Basics> | 接口: IDaq | 包: Snet.DB
通过 ADO.NET(Dapper,Daq 模式)/ SqlSugarCore(Default 模式)连接关系型数据库(SQL Server、MySQL、Oracle、SQLite)。与其余 DAQ 驱动一样,DBOperate 是一个标准 IDaq 驱动——用 OnAsync 连接,然后读取/订阅。由 HandlerType 选择两种工作模式:
| 模式 | 库 | 用途 |
|---|---|---|
Daq(默认) |
Dapper | 数据采集——每个 AddressName 是一条 SQL 语句,每次读取/轮询周期执行一次;查询结果以 JSON 数组返回 |
Default |
SqlSugar(SqlSugarCore) | 增删改查——基于实体类的创建/插入/更新/删除/查询辅助方法(见 CRUD 操作) |
WriteAsync 不支持——该驱动对 DAQ 接口只读。写操作请使用 Default 模式的 CRUD 辅助方法(InsertAsync、UpdateAsync、DeleteAsync 等)。
快速开始
using Snet.DB;
using Snet.Model.data;
using Snet.Model.@enum;
using System.Collections.Concurrent;
// 连接数据库(Daq 模式:数据采集)
var op = await DBOperate.InstanceAsync(new DBData.Basics
{
ConnectStr = "Server=192.168.1.100;Database=MyDB;User Id=sa;Password=xxx;",
DBType = DBData.DBType.SqlServer,
HandlerType = DBData.DBHandlerType.Daq
});
await op.OnAsync();
// 读取——AddressName 就是 SQL 语句(不能为空)
var address = new Address(new List<AddressDetails>
{
new("Temp", "SELECT TOP 10 Temperature, Time FROM SensorData", DataType.String),
new("Avg", "SELECT AVG(Temperature) AS AvgTemp FROM SensorData", DataType.String)
});
var result = await op.ReadAsync(address);
if (result.Status)
{
// ResultData 是以 AddressName 为键的字典;
// 每个 ResultValue 是查询结果(JSON 数组字符串)
var data = result.GetSource<ConcurrentDictionary<string, AddressValue>>();
foreach (var kv in data)
Console.WriteLine($"{kv.Key} = {kv.Value.ResultValue}");
}
else
{
Console.WriteLine($"读取失败: {result.Message}");
}
// 订阅(轮询式)——先绑定事件,再订阅;每 HandleInterval 毫秒执行一次 SQL
op.OnDataEventAsync += async (sender, e) =>
{
if (e.Status)
{
var data = e.GetSource<ConcurrentDictionary<string, AddressValue>>();
foreach (var kv in data)
Console.WriteLine($"{kv.Key} = {kv.Value.ResultValue}");
}
else
{
Console.WriteLine($"{e.Message}");
}
};
var subAddr = new Address(new List<AddressDetails>
{
new("Temp", "SELECT TOP 10 Temperature, Time FROM SensorData", DataType.String)
});
await op.SubscribeAsync(subAddr);
// 后续:取消订阅
// await op.UnSubscribeAsync(subAddr);
await op.DisposeAsync();
安装
dotnet add package Snet.DB
配置参数(DBData.Basics)
继承 SubscribeData.SCData(轮询订阅参数),并 override 了 HandleInterval——DB 驱动默认每 10 秒采集一次(源码:DBData.Basics.HandleInterval = 10000)。
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
SN |
string |
Guid.NewGuid().ToUpperNString() |
唯一标识符 |
ConnectStr |
string |
— | 必填。 数据库连接字符串(格式取决于 DBType,见下) |
DBType |
DBType |
SQLite |
数据库类型:SqlServer / MySql / Oracle / SQLite |
HandlerType |
DBHandlerType |
Daq |
处理模式:Daq(Dapper 采集)/ Default(SqlSugarCore 增删改查) |
ChangeOut |
bool |
true |
(继承)订阅:变化时输出 |
AllOut |
bool |
false |
(继承)当 ChangeOut = true 时:未变化项与变化项一同抛出,保证批次数据完整 |
HandleInterval |
int |
10000 |
(override)轮询间隔——每个订阅周期执行 SQL 的频率(毫秒)。DB 驱动默认每 10 秒采集一次;调小间隔更新更快(数据库负载更高) |
TaskNumber |
int |
5 |
(继承)订阅任务数量 |
各 DBType 的连接字符串示例
DBType |
提供程序 | ConnectStr 示例 |
|---|---|---|
SqlServer |
Microsoft.Data.SqlClient | Server=192.168.1.100;Database=MyDB;User Id=sa;Password=xxx; |
MySql |
MySqlConnector | Server=192.168.1.100;Port=3306;Database=MyDB;User Id=root;Password=xxx; |
Oracle |
Oracle.ManagedDataAccess | User Id=system;Password=xxx;Data Source=192.168.1.100:1521/ORCL |
SQLite |
Microsoft.Data.Sqlite | Data Source=sensor.db |
API 参考
IDaq 生命周期(继承)
| 方法 | 描述 |
|---|---|
Task<OperateResult> OnAsync(CancellationToken token = default) |
Daq 模式:校验配置的数据库类型与 ConnectStr 可创建连接——不持有长连接;实际查询每次操作新建短连接(连接池复用)。Default 模式:创建 SqlSugarClient(SqlSugarCore NuGet,using SqlSugar;)并打开 |
Task<OperateResult> OffAsync(bool hardClose = false, CancellationToken token = default) |
停止轮询订阅(Daq 模式)、关闭并释放连接 / SqlSugarClient |
Task<OperateResult> GetStatusAsync(CancellationToken token = default) |
Daq 模式:OnAsync 置位的状态标志(无实际长连接可检查)。Default 模式:客户端打开状态 |
Task<OperateResult> GetBaseObjectAsync(CancellationToken token = default) |
ResultData:Daq 模式 → null(无持久连接对象);Default 模式 → SqlSugarClient(SqlSugarCore)——用它可以实现内置 CRUD 之外更多的数据库操作 |
Task<OperateResult> ReadAsync(Address address, CancellationToken token = default) |
把每个 AddressName 当作 SQL 查询执行(Dapper)。结果转换为 JSON 数组存入 AddressValue.ResultValue。ResultData:以 AddressName 为键的 ConcurrentDictionary<string, AddressValue>。空的 AddressName 会失败("不能为空,DB操作时此属性为SQL语句") |
Task<OperateResult> WriteAsync(ConcurrentDictionary<string, (object value, EncodingType? encodingType)> values, CancellationToken token = default) |
不支持——始终失败("不支持 Write 操作")。写操作请用 Default 模式 CRUD 辅助方法 |
Task<OperateResult> SubscribeAsync(Address address, CancellationToken token = default) |
基于 SubscribeOperate 的轮询订阅——每 HandleInterval 毫秒(默认 10000,即每 10 秒)经 ReadAsync 执行一次 SQL 并触发 OnDataEventAsync。先绑定事件 |
Task<OperateResult> UnSubscribeAsync(Address address, CancellationToken token = default) |
移除轮询订阅 |
完整的
IDaq语义(事件系统OnDataEventAsync/OnInfoEventAsync、虚拟地址等)参见 IDaq 接口。
采集地址规则
在 Daq 模式下,每个 AddressDetails.AddressName 本身就是 SQL 语句——不能为空。每次读取/轮询周期,驱动执行该语句(Dapper)并将完整结果集以 JSON 数组字符串存储:
[
{ "Temperature": 25.6, "Time": "2026-08-09T10:00:00" },
{ "Temperature": 25.8, "Time": "2026-08-09T10:00:01" }
]
此类点位请使用 DataType.String,需要取个别字段时在应用程序中自行解析 JSON。
CRUD 操作(SqlSugarCore,Default 模式)
当 HandlerType = Default 时可用(使用官方 SqlSugarCore 客户端;需要与表映射的实体类)。同步 + 异步成对提供。
内嵌的 SqlSugar fork 源码(命名空间 Snet.DB.sugar)已移除——自 26.250.1 起包引用官方 SqlSugarCore NuGet(当前 5.1.4.220),代码使用 using SqlSugar; / SqlSugar.DbType.*;dotnet add package Snet.DB 仍然全量引入。
| 操作 | 方法 |
|---|---|
| 执行 SQL | Execute(sql)、Execute(sql, T t)、Execute(sql, List<T> t) + ExecuteAsync(...) |
| 实体查询 | Query<T>(sql) / Query<T>(sql, T t) / QueryList<T>(sql, T t) / Query<T>(sql, int[] ids) + Async 变体;QueryDataTable<T>(sql, T t);QueryMultiple<T>(sql, T t)(返回 List<T>——第一结果集) |
| 创建表 | Create<T>() / CreateAsync<T>()(经 CodeFirst.InitTables<T>) |
| 表是否存在 | Exist<T>() / ExistAsync<T>() |
| 插入 | Insert<T>(T obj) / Insert<T>(List<T> objs) + Async 变体 |
| 修改 | Update<T>(T obj, Expression<Func<T, object>> updateColumns, Expression<Func<T, bool>> condition) + Async |
| 删除 | Delete<T>(Expression<Func<T, bool>> condition) + Async |
CRUD 辅助方法要求 HandlerType = Default。Daq 模式下 SqlSugarCore 客户端不会被创建——调用会抛异常。
获取 SqlSugar 对象(GetBaseObjectAsync)
Default 模式可直接拿到底层 SqlSugarClient——用它实现内置辅助方法之外的更多数据库操作(原生 SQL、存储过程、事务等):
using SqlSugar; // SqlSugarCore NuGet
var op = await DBOperate.InstanceAsync(new DBData.Basics
{
ConnectStr = "Data Source=sensor.db",
DBType = DBData.DBType.SQLite,
HandlerType = DBData.DBHandlerType.Default
});
await op.OnAsync();
var baseObject = (await op.GetBaseObjectAsync()).GetSource<SqlSugarClient>();
int count = await baseObject.Queryable<SensorRecord>().CountAsync(); // 任意 SqlSugarCore API
DBData.DBHandlerType.Default 的源码描述为 “使用 SqlSugarCore;使用 GetBaseObjectAsync 可以获取到 SqlSugarCore 的对象,使用 SqlSugarCore 对象可以自行实现更多的数据库的数据操作”(见 DBData.cs)。
支持的数据类型
AddressName 就是 SQL 语句 — 每个点位都以 JSON 数组字符串返回。
| 数据类型 | 支持 | 说明 |
|---|---|---|
String |
✅ | 查询结果序列化为 JSON 数组字符串(JArray) |
| 其余全部 | — | DataType 被忽略 — 恒返回 JSON 数组字符串 |
DAQ 接口不支持写入 — 写操作请用 HandlerType = Default 的 CRUD 辅助方法。
地址自动组包与解包
每个数采驱动都实现 IPacker。DBOperate 使用符号(SQL 字符串)寻址,未注册在自动组包注册表中,因此 Packer 原样返回地址集合(透传)——见 IPacker 接口。支持字节偏移寻址的协议(如西门子 S7、Modbus、三菱 MC)可组包:地址自动组包。
代码示例
一次调用读取多条 SQL
var address = new Address(new List<AddressDetails>
{
new("Latest", "SELECT TOP 1 * FROM SensorData ORDER BY Time DESC", DataType.String),
new("Count", "SELECT COUNT(*) AS Cnt FROM SensorData", DataType.String)
});
var result = await op.ReadAsync(address);
if (result.Status)
{
var data = result.GetSource<ConcurrentDictionary<string, AddressValue>>();
foreach (var kv in data)
Console.WriteLine($"{kv.Key} = {kv.Value.ResultValue}"); // JSON 数组字符串
}
轮询订阅
var op = await DBOperate.InstanceAsync(new DBData.Basics
{
ConnectStr = "Data Source=sensor.db",
DBType = DBData.DBType.SQLite,
HandlerType = DBData.DBHandlerType.Daq,
HandleInterval = 5000 // 每 5 秒执行一次 SQL
});
await op.OnAsync();
op.OnDataEventAsync += async (sender, e) =>
{
if (e.Status)
{
var data = e.GetSource<ConcurrentDictionary<string, AddressValue>>();
foreach (var kv in data)
Console.WriteLine($"[{DateTime.Now:HH:mm:ss.fff}] {kv.Key} = {kv.Value.ResultValue}");
}
};
await op.SubscribeAsync(new Address(new AddressDetails("Rows", "SELECT * FROM SensorData", DataType.String)));
实体 CRUD(Default 模式)
public class SensorRecord
{
public int Id { get; set; }
public float Temperature { get; set; }
}
var op = await DBOperate.InstanceAsync(new DBData.Basics
{
ConnectStr = "Data Source=sensor.db",
DBType = DBData.DBType.SQLite,
HandlerType = DBData.DBHandlerType.Default
});
await op.OnAsync();
await op.CreateAsync<SensorRecord>(); // 创建表
await op.InsertAsync(new SensorRecord { Temperature = 25.6f });
var rows = (await op.QueryAsync<SensorRecord>(r => r.Temperature > 20)).GetSource<List<SensorRecord>>();
await op.UpdateAsync(new SensorRecord { Id = 1, Temperature = 26.0f },
u => new { u.Temperature }, c => c.Id == 1);
await op.DeleteAsync<SensorRecord>(r => r.Id == 1);
常见问题
Q: ReadAsync 报"不能为空,DB操作时此属性为SQL语句"?
A: DB 模式下 AddressName 就是 SQL 语句——不能为空。Address 中的每个 AddressDetails 必须携带真实的 SQL 查询字符串。
Q: ResultValue 里是什么?
A: 完整查询结果序列化为 JSON 数组字符串(如 [{"Temperature":25.6,...}])。取个别字段请在应用程序中解析 JSON。
Q: WriteAsync 能用吗?
A: 不能——WriteAsync 在 DAQ 接口上不支持(返回"不支持 Write 操作")。写操作请改用 Default 模式的 CRUD 辅助方法(InsertAsync / UpdateAsync / DeleteAsync)配合实体类。
Q: ConnectStr 用什么格式?
A: 取决于 DBType——见 配置参数 中各数据库示例(SqlServer / MySql / Oracle / SQLite)。
Q: Daq 和 Default 模式的区别?
A: Daq(默认)= 数据采集:AddressName 是 SQL 语句,由 Dapper 执行,结果以 JSON 返回。Default = 实体 CRUD(SqlSugarCore 的 InsertAsync/UpdateAsync/DeleteAsync/QueryAsync<T>/CreateAsync<T>),另有 GetBaseObjectAsync 返回 SqlSugarClient 用于自定义数据库操作。按场景选择;每个实例同时只启用一种模式。
Q: 订阅是实时的吗?
A: 是轮询——SubscribeOperate 每 HandleInterval 毫秒经 ReadAsync 执行一次 SQL。调小间隔更新更快(数据库负载更高)。
参见
- DaqAbstract 基类 — 生命周期与同步包装
- IDaq 接口 · 地址模型
- 地址自动组包 — 本驱动为何透传地址
