数据库采集 - Snet Docs

📘 数据库协议驱动

命名空间: 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 辅助方法(InsertAsyncUpdateAsyncDeleteAsync 等)。

快速开始

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 = DefaultDaq 模式下 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: DaqDefault 模式的区别? A: Daq(默认)= 数据采集:AddressName 是 SQL 语句,由 Dapper 执行,结果以 JSON 返回。Default = 实体 CRUD(SqlSugarCore 的 InsertAsync/UpdateAsync/DeleteAsync/QueryAsync<T>/CreateAsync<T>),另有 GetBaseObjectAsync 返回 SqlSugarClient 用于自定义数据库操作。按场景选择;每个实例同时只启用一种模式。

Q: 订阅是实时的吗? A: 是轮询——SubscribeOperateHandleInterval 毫秒经 ReadAsync 执行一次 SQL。调小间隔更新更快(数据库负载更高)。

参见