采集插件二开说明
本文面向需要扩展 ThingsGatewayRuntime 采集协议的开发人员。采集插件负责把 PLC、仪表、传感器、上位系统或虚拟数据源读入运行时变量,并在允许时把外部写入/RPC 反写到现场设备。
如果要扩展数据转发、协议服务端、数据库写入或云平台对接,请阅读 业务插件二开说明。
源码入口
| 入口 | 作用 |
|---|---|
ThingsGatewayRuntime.Application/Driver/IDriver.cs | 所有设备驱动运行时接口。 |
ThingsGatewayRuntime.Application/Driver/DriverBase.cs | 设备插件生命周期、日志、任务调度、通道挂载、释放逻辑。 |
ThingsGatewayRuntime.Application/Driver/Collect/CollectBase.cs | 采集插件核心基类,负责变量打包、定时读、脚本变量、方法变量、写入/RPC 流程。 |
ThingsGatewayRuntime.Application/Driver/Collect/CollectFoundationBase.cs | 基于 Foundation IDevice 的主站类协议模板。 |
ThingsGatewayRuntime.Application/Driver/Collect/CollectReceivedFoundationBase.cs | 基于 Foundation IReceivedDevice 的被动接收类协议模板。 |
ThingsGatewayRuntime.Application/Task/Collect/DeviceManage/DeviceThreadManage.cs | 设备插件创建、属性注入、通道初始化、启动任务循环的调度入口。 |
ThingsGatewayRuntime.Application/Service/Plugin/PluginService.cs | 扫描 CollectBase 派生类并生成插件清单。 |
ThingsGatewayRuntime.Plugin/Plugin/* | 开源版采集插件实现。 |
ThingsGatewayRuntime.NOAOTPlugin/Plugin/* | 非 AOT 采集插件实现,例如 OPC UA、OPC DA。 |
ThingsGatewayRuntimePRO/src/ThingsGatewayRuntime.NOAOTPROPlugin/Plugin/* | 专业版采集插件实现。 |
一句话流程
运行时按下面顺序处理采集插件:
PluginService扫描所有非抽象CollectBase派生类,形成插件列表。- 设备启动时,
DeviceThreadManage.CreateDriver通过插件全名创建实例。 DriverBase.InitDevice挂载DeviceRuntime、日志和device.Driver。PluginServiceUtil.SetDriverProperties把设备插件属性字典写回强类型属性对象。DriverBase.InitChannelAsync挂载通道,调用AfterVariablesChangedAsync打包变量。DriverBase.StartAsync调用插件的ProtectedStartAsync,再创建并启动TaskSchedulerLoop。CollectBase按变量间隔读取,失败重试,成功后变量上线,失败后变量离线。- 外部写入/RPC 进入
InvokeWriteAsync或InvokeMethodAsync,由插件落到协议写入。 - 停止设备时调用
StopAsync,最终进入SafetyDisposeAsync释放通道、底层协议对象、日志和锁。
插件发现规则
采集插件必须满足这些条件才会出现在采集设备配置中。
| 规则 | 说明 |
|---|---|
继承 CollectBase | PluginService 只把 CollectBase 的非抽象派生类识别为采集插件。 |
| 公开无参构造 | 运行时通过 Activator.CreateInstance 创建实例;没有无参构造会启动失败。 |
| 不要在构造函数连接现场设备 | 构造阶段还没有设备属性、通道、日志和取消令牌。连接应放在 InitChannelAsync 或 ProtectedStartAsync。 |
| 插件全名是配置键 | 设备实体保存的是类型 FullName,重命名命名空间或类名会影响旧配置。 |
[DisplayNonePlugin] 会隐藏插件 | MemoryDriver 使用该标记,作为内部内存设备模板,不作为普通现场插件显示。 |
[OnlyWindowsSupport] 会限制平台 | 非 Windows 环境下带该属性的插件不会显示。 |
NOAOT 程序集不支持 AOT | PluginService 根据程序集名标记 SupportsAot,OPC/COM 等插件通常位于 NOAOT 程序集。 |
基类选择
| 基类 | 适合场景 | 必须重点实现 | 现有参考 |
|---|---|---|---|
CollectFoundationBase | 协议可以抽象成按地址读写字节,例如 Modbus、S7、DLT645、Omron、Melsec。 | FoundationDevice、CollectProperties、DriverPropertyType、InitChannelAsync、ProtectedLoadSourceReadAsync。通常不需要重写 ReadSourceAsync。 | ModbusMaster、SiemensS7Master、ControlLogixMaster。 |
CollectReceivedFoundationBase | 数据由对端主动上报或底层设备维护连接状态,普通定时读不成立。 | FoundationReceivedDevice、属性类型、上报事件到变量的映射,必要时重写 AfterVariablesChangedAsync。 | HJ212Master、EDPF_NTMaster、KELID2008Master、LKSISMaster。 |
CollectBase | 自定义协议、订阅协议、消息协议、虚拟设备、扫码器、CAN、OPC、IEC61850、ZeroMQ。 | ProtectedLoadSourceReadAsync、ReadSourceAsync、WriteValuesAsync、IsConnected,必要时重写 ProtectedStartAsync、AfterVariablesChangedAsync。 | OpcUaMaster、CanMaster、IEC61850Master、MemoryDriver、MqttCollectBase。 |
MqttCollectBase | MQTT 客户端/服务端采集,变量值由 Topic 消息更新。 | 派生类处理连接、订阅、消息接收和在线判断;父类负责按变量关系更新值。 | MqttCollectClient、MqttCollectServer。 |
不要为了“省代码”继承过高层的基类。能用 CollectFoundationBase 就不要手写读写锁、重试和字节解析;必须自定义连接和订阅时再用 CollectBase。
属性模型
插件属性分两层:基类属性负责公共行为,具体属性类负责协议参数。
| 类型 | 作用 |
|---|---|
DriverPropertyBase | 所有插件属性根类型。只有标记 [DynamicProperty] 的属性会暴露给前端、导入导出和运行时属性注入。 |
CollectPropertyBase | 采集属性根类型,包含并发、离线恢复间隔、重试、读写占空比、写优先等内部字段。 |
CollectPropertyRetryBase | 暴露 失败重试次数、读写占空比、写优先。 |
CollectFoundationPropertyBase | 增加 读写超时时间、帧前时间、字符串反转字节、数据解析顺序。 |
CollectFoundationPackPropertyBase | 增加 最大打包长度。 |
CollectFoundationDtuPropertyBase | 增加 DTU ID。 |
CollectFoundationDtuPackPropertyBase | 同时包含打包长度和 DTU ID。 |
CollectPropertyNone | 无页面属性的内部插件属性,例如 MemoryDriver。 |
实现属性时要同时保证下面两个成员指向同一个属性类型实例:
private readonly MyDriverProperty _driverProperties = new();
public override CollectPropertyBase CollectProperties => _driverProperties;
public override Type DriverPropertyType { get; } = typeof(MyDriverProperty);
DriverPropertyType 用于前端动态表单、Excel 导入导出和属性反序列化;CollectProperties 是运行时真正读取的对象。两者不一致时,页面看到的字段和插件实际使用的字段会错位。
DynamicProperty 约定
| 项目 | 说明 |
|---|---|
Description | 页面显示名称,也是未配置本地化资源时的默认中文名。 |
Remark | 页面提示说明,可写地址格式、单位、注意事项。 |
GroupName | 前端分组名,属性较多时建议使用。 |
ExpressionType | 标识脚本输入类型,常用于动态模型、表达式编辑器。 |
CertificatePurpose | 标识证书选择用途,前端会从证书管理中提供下拉选择。 |
| 枚举属性 | 建议加 JsonStringEnumConverter<T> 或确保枚举项已进入源生成缓存,避免导入导出和 JSON 显示不一致。 |
| 默认值 | 前端新建配置会读取属性实例默认值,默认值必须是可直接用于测试环境的保守值。 |
没有 [DynamicProperty] 的属性不会被页面保存,也不会从设备属性字典注入。不要把运行时缓存、连接对象、锁、队列标成动态属性。
生命周期钩子
| 阶段 | 运行时动作 | 插件开发注意事项 |
|---|---|---|
| 构造函数 | Activator.CreateInstance 创建实例。 | 只初始化轻量字段,不读配置、不连设备、不启动线程。 |
InitDevice | 设置 CurrentDevice、日志、device.Driver,再调用 ProtectedInitDevice。 | 如需读取设备运行态,可在 ProtectedInitDevice 后使用;属性值此时还未注入完成。 |
| 属性注入 | PluginServiceUtil.SetDriverProperties 写入 [DynamicProperty] 值。 | 具体协议参数应在后续阶段读取,不要在构造函数缓存。 |
InitChannelAsync | 设置 ChannelObject,必要时 Channel.SetupAsync,最后 AfterVariablesChangedAsync。 | Foundation 插件通常在这里重建底层 IDevice,把属性赋给底层对象,并调用 InitChannel。 |
AfterVariablesChangedAsync | 基类重建 VariableSourceReads、脚本变量、方法变量和定时任务。 | 变量新增、删除、属性修改都会触发。重写时注意调用 base,并清理旧订阅或旧映射。 |
ProtectedStartAsync | 启动通讯前调用,受 StartTimeout 控制。 | 连接远端、启动服务、订阅消息可放这里。必须尊重取消令牌。 |
ProtectedGetTasks | 构建调度任务。 | CollectBase 已提供设备状态、在线测试、变量读取任务。只有特殊协议才重写。 |
| 定时读取 | 调用 ReadSourceAsync,失败按 RetryCount 重试。 | 成功要设置变量值或返回可解析字节;失败要返回失败结果,不要吞异常后假成功。 |
| 写入/RPC | InvokeWriteAsync、InvokeMethodAsync。 | 自定义写入要使用写锁,必要时执行写后回读校验。 |
StopAsync/SafetyDisposeAsync | 停止任务循环,释放日志、底层设备、事件订阅、锁。 | 必须取消事件订阅、释放 socket/client/channel,避免重启后重复接收。 |
变量读取模型
CollectBase.AfterVariablesChangedAsync 会把启用变量分成三类。
| 类型 | 判断方式 | 运行方式 |
|---|---|---|
| 普通源读取变量 | 地址不是 DeviceStatus、Script、ScriptRead,且没有 OtherMethod。 | 进入 ProtectedLoadSourceReadAsync,生成 VariableSourceRead 后由定时任务调用 ReadSourceAsync。 |
| 脚本/特殊变量 | 地址为 DeviceStatus、Script、ScriptRead。 | 进入 VariableScriptReads,由脚本任务或特殊地址逻辑更新。 |
| 方法变量 | 配置了 OtherMethod。 | 根据 [DynamicMethod] 找方法,读时调用方法,写时把参数传给方法。 |
ProtectedLoadSourceReadAsync 的职责不是读数据,而是把变量整理成“读包”。Foundation 插件通常这样写:
protected override Task<List<VariableSourceRead>> ProtectedLoadSourceReadAsync(List<VariableRuntime> deviceVariables)
{
List<VariableSourceRead> reads = new();
foreach (var group in deviceVariables.GroupBy(a => a.CollectGroup))
{
reads.AddRange(_plc.LoadSourceRead<VariableSourceRead, VariableRuntime>(
group,
_driverProperties.MaxPack,
CurrentDevice.IntervalTime));
}
return Task.FromResult(reads);
}
这样做有两个好处:同一采集组可独立打包,MaxPack 能限制单包长度,变量自己的 IntervalTime 或设备默认间隔能继续生效。
读写锁和写优先
CollectBase 使用 AsyncReadWriteLock 管理读写冲突。
| 机制 | 说明 |
|---|---|
DutyCycle | 写入较多时,按占空比穿插读取,避免长时间只写不读。 |
WritePriority | 写入时可取消正在等待的读取,让控制命令优先下发。 |
ReadSourceAsync | 定时读进入读锁。自定义读不要再长期阻塞线程。 |
WriteValuesAsync | 写入实现应进入写锁,Foundation 基类已经处理;直接继承 CollectBase 时要自己处理。 |
Check | 写入成功后,开启 RpcWriteCheck 的变量会回读校验。自定义写入实现可复用该方法。 |
不要在写锁中做无关的长耗时工作,例如等待外部 HTTP、写大文件或启动线程。锁只保护协议读写临界区。
连接和在线状态
| 方法 | 说明 |
|---|---|
IsConnected() | 前端状态、在线测试、设备状态判断都会用。必须反映真实连接状态,不能永远返回 true,除非像 MemoryDriver 这类永久在线。 |
ProtectedStartAsync | 首次连接或启动服务的位置。超时后会记录启动失败。 |
TestOnline | CollectReceivedFoundationBase 默认会尝试重连并把变量置离线;自定义协议可重写做轻量检测。 |
SetDeviceStatus | CollectBase 根据连接状态和变量在线情况更新设备状态。特殊插件可重写,例如 MemoryDriver 永久在线。 |
现场上“设备在线”不等于“每个变量都在线”。如果某个读包失败,CollectBase 会把该读包内变量置为离线,并记录最后错误。
特殊地址和动态方法
| 能力 | 写法 | 场景 |
|---|---|---|
| 设备状态变量 | 变量地址写 DeviceStatus。 | 把设备在线/离线状态作为普通变量参与展示或转发。 |
| 脚本变量 | 地址写 Script 或 ScriptRead。 | 不走现场通讯,通过表达式或脚本计算值。 |
| 动态方法 | 方法标记 [DynamicMethod("说明", "备注")]。 | 读写日期、调用协议特有命令、执行对象方法。 |
| RPC 写入 | 外部接口写变量。 | 进入 InvokeWriteAsync 或 InvokeMethodAsync,再由插件写现场设备。 |
动态方法返回值必须能转成 OperResult 或 IOperResult<T>。变量上配置 OtherMethod 后,基类会按方法名找到对应方法。
Foundation 插件模板
CollectFoundationBase 已经实现了大部分固定流程。
| 你需要写 | 说明 |
|---|---|
底层 _plc 字段 | 使用 Foundation 中的协议设备类。 |
FoundationDevice | 返回当前 _plc。 |
InitChannelAsync | 重建 _plc,从属性对象写入协议参数,调用 _plc.InitChannel(channelObject, LogMessage),最后调用 base.InitChannelAsync。 |
ProtectedLoadSourceReadAsync | 用底层设备的地址解析/打包能力生成 VariableSourceRead。 |
可选 WriteValuesAsync | 协议有批量写、特殊数据类型或字符串写法时重写。 |
可选 [DynamicMethod] | 暴露协议特有方法给变量或 RPC。 |
ModbusMaster 是最小 Foundation 模板;SiemensS7Master 展示了批量写和动态方法;ControlLogixMaster 展示了 Foundation 写入不够用时的自定义写入。
直接继承 CollectBase 的模板
直接继承 CollectBase 时,插件要自己完成读包、读、写和连接状态。
[DynamicallyAccessedMembers(DynamicallyAccessedMemberTypes.All)]
public sealed class MyCollectDriver : CollectBase
{
private readonly MyCollectProperty _properties = new();
private MyClient? _client;
public override CollectPropertyBase CollectProperties => _properties;
public override Type DriverPropertyType { get; } = typeof(MyCollectProperty);
public override bool IsConnected() => _client?.Connected == true;
public override ChannelTypeEnum[] SupportedChannelTypes() => [ChannelTypeEnum.Other];
public override DataTypeEnum[] SupportedDataTypes() =>
[DataTypeEnum.Boolean, DataTypeEnum.Int16, DataTypeEnum.Int32, DataTypeEnum.Float, DataTypeEnum.String];
protected override async Task ProtectedStartAsync(CancellationToken cancellationToken)
{
_client = new MyClient(_properties.Endpoint);
await _client.ConnectAsync(cancellationToken).ConfigureAwait(false);
}
protected override Task<List<VariableSourceRead>> ProtectedLoadSourceReadAsync(List<VariableRuntime> deviceVariables)
{
var reads = deviceVariables
.GroupBy(a => string.IsNullOrWhiteSpace(a.IntervalTime) ? CurrentDevice.IntervalTime : a.IntervalTime)
.Select(group =>
{
var read = new VariableSourceRead { IntervalTime = new(group.Key) };
read.AddVariableRange(group);
return read;
})
.ToList();
return Task.FromResult(reads);
}
protected override async ValueTask<OperResult<ReadOnlyMemory<byte>>> ReadSourceAsync(
VariableSourceRead sourceRead,
CancellationToken cancellationToken)
{
foreach (var variable in sourceRead.Variables)
{
var value = await _client!.ReadAsync(variable.RegisterAddress!, cancellationToken).ConfigureAwait(false);
variable.SetValue(value, DateTime.UtcNow);
}
return new OperResult<ReadOnlyMemory<byte>>();
}
protected override async ValueTask<Dictionary<string, OperResult>> WriteValuesAsync(
Dictionary<VariableRuntime, JsonElement> writeInfoLists,
CancellationToken cancellationToken)
{
using var writeLock = await ReadWriteLock.WriterLockAsync().ConfigureAwait(false);
var results = new Dictionary<string, OperResult>();
foreach (var (variable, value) in writeInfoLists)
{
results[variable.Name] = await _client!.WriteAsync(
variable.RegisterAddress!,
value.GetObjectFromJsonElement(),
cancellationToken).ConfigureAwait(false);
}
return results;
}
protected override async Task SafetyDisposeAsync(bool disposing)
{
if (_client != null)
await _client.DisposeAsync().ConfigureAwait(false);
await base.SafetyDisposeAsync(disposing).ConfigureAwait(false);
}
}
public sealed class MyCollectProperty : CollectPropertyRetryBase
{
[DynamicProperty("连接地址", Remark = "例如 127.0.0.1:12345")]
public string Endpoint { get; set; } = "127.0.0.1:12345";
}
当前源码采集插件参考矩阵
下表按当前源码核对了 30 个采集插件。开发新插件前,先找协议形态最接近的实现参考。
| 插件 | 基类 | 属性类 | 参考重点 |
|---|---|---|---|
MemoryDriver | CollectBase | CollectPropertyNone | 内存变量、表达式触发、全局变量变化事件、无真实通道。 |
ModbusMaster | CollectFoundationBase | ModbusMasterProperty | Foundation 标准模板、DTU、站号、最大打包长度。 |
SiemensS7Master | CollectFoundationBase | SiemensS7MasterProperty | S7 地址打包、批量写、日期/时间动态方法。 |
Dlt645_2007Master | CollectFoundationBase | Dlt645_2007MasterProperty | 仪表地址、密码、操作员代码、前导符。 |
OpcUaMaster | CollectBase | OpcUaMasterProperty | 订阅/轮询混合、证书、安全策略、节点类型加载、订阅刷新。 |
OpcDaMaster | CollectBase | OpcDaMasterProperty | Windows/COM OPC DA、订阅组、重连检查、服务端时间。 |
MqttCollectClient | MqttCollectBase | MqttCollectClientProperty | MQTT 客户端连接、订阅 Topic、消息映射变量。 |
MqttCollectServer | MqttCollectBase | MqttCollectServerProperty | MQTT 服务端监听、客户端校验、消息映射变量。 |
CanMaster | CollectBase | CanMasterProperty | CAN/CAN FD 端点、帧组装、硬件过滤、全帧切片读取。 |
ControlLogixMaster | CollectFoundationBase | ControlLogixMasterProperty | AllenBradley CIP,Foundation 读包,自定义写入。 |
PCCCMaster | CollectFoundationBase | PCCCMasterProperty | AllenBradley PCCC,Foundation 读写模板。 |
DCONMaster | CollectFoundationBase | DCONMasterProperty | DCON 协议,DTU 和打包模板。 |
EDPF_NTMaster | CollectReceivedFoundationBase | EDPF_NTUdpProperty | UDP 被动接收,变量刷新时重建映射,不实现主动读写。 |
HJ212Master | CollectReceivedFoundationBase | HJ212MasterProperty | 环保 HJ212 上报,地址说明和上报映射。 |
IEC61850Master | CollectBase | IEC61850MasterProperty | MMS 读写、RCB 报告、GOOSE、SOE、TLS、复杂订阅生命周期。 |
InovanceMaster | CollectFoundationBase | InovanceMasterProperty | 汇川协议,继承 Foundation DTU 打包属性。 |
KELID2008Master | CollectReceivedFoundationBase | KELID2008MasterProperty | 上报协议,地址说明扩展,主动读写未实现。 |
LKSISMaster | CollectReceivedFoundationBase | LKSISProperty | UDP 上报协议,变量变更后维护接收映射。 |
Mc1E_BinaryMaster | CollectFoundationBase | Mc1E_BinaryMasterProperty | 三菱 1E 二进制,Foundation 打包模板。 |
Mc3E_BinaryMaster | CollectFoundationBase | Mc3E_BinaryMasterProperty | 三菱 3E 二进制,Foundation 打包模板。 |
ModbusC1Master | CollectFoundationBase | ModbusC1MasterProperty | Modbus C1,DTU 打包模板。 |
ModbusC20Master | CollectFoundationBase | ModbusC20MasterProperty | Modbus C20,DTU 打包模板。 |
OmronFinsMaster | CollectFoundationBase | OmronFinsMasterProperty | Omron FINS,Foundation 打包模板。 |
OpcAeMaster | CollectBase | OpcAeMasterProperty | OPC AE 事件采集,事件订阅转变量。 |
SECSMaster | CollectReceivedFoundationBase | SECSMasterProperty | SECS/GEM,既有接收设备,也实现了主动读写。 |
TIANXINMaster | CollectFoundationBase | TIANXINMasterProperty | 天信仪表类协议,DTU 打包模板。 |
TS550Master | CollectFoundationBase | TS550MasterProperty | TS550 协议,Foundation 通用读写属性。 |
USBScaner | CollectBase | USBScanerProperty | 本机扫码器 Hook,扫码事件写入变量,不实现主动读写。 |
VigorMaster | CollectFoundationBase | VigorMasterProperty | Vigor 协议,DTU 打包模板。 |
ZeroMQCollectClient | CollectBase | ZeroMQCollectClientProperty | ZeroMQ 订阅/连接、Topic 前缀、高水位、重连和清理。 |
常见开发错误
| 错误 | 后果 | 正确做法 |
|---|---|---|
| 构造函数里读取属性或连接设备 | 属性还没注入,日志和取消令牌也不存在。 | 构造函数只初始化字段,连接放到 InitChannelAsync 或 ProtectedStartAsync。 |
CollectProperties 和 DriverPropertyType 不一致 | 页面字段、导入导出和运行时读取错位。 | 两者都指向同一个强类型属性类。 |
重写 InitChannelAsync 后不调用 base | 变量不会打包,定时任务没有数据源。 | 完成底层对象初始化后调用 base.InitChannelAsync。 |
| 变量变化后不清理旧订阅 | 重复接收、重复写变量、内存泄漏。 | 在 AfterVariablesChangedAsync 或释放流程中先移除旧映射。 |
| 读失败仍返回成功 | 变量看似在线但值错误。 | 返回失败 OperResult,让基类置离线并记录错误。 |
忽略 CancellationToken | 停止设备或重启时卡住。 | 所有网络、串口、等待和重试都传入取消令牌。 |
| 写入不加写锁 | 读写并发冲突,现场协议包交叉。 | 自定义 WriteValuesAsync 使用 ReadWriteLock.WriterLockAsync。 |
| 释放时不取消事件订阅 | 重启后同一数据被处理多次。 | SafetyDisposeAsync 中取消订阅并释放底层对象。 |
验证清单
开发完成后至少验证这些场景。
| 场景 | 验证点 |
|---|---|
| 插件发现 | 插件管理能看到插件,采集设备下拉能选择,插件类型为采集。 |
| 属性表单 | 所有 [DynamicProperty] 字段显示、默认值、枚举、证书下拉、导入导出都正确。 |
| 通道类型 | SupportedChannelTypes() 与现场连接方式一致;无普通通道时返回 Other。 |
| 数据类型 | SupportedDataTypes() 不要暴露协议无法解析的类型。 |
| 单点读取 | 使用设备调试先读 3 到 5 个点,覆盖布尔、整数、浮点、字符串或协议特有类型。 |
| 批量读取 | 点表导入后确认打包数量、间隔和超时符合预期。 |
| 写入/RPC | 只读点不可写,可写点能写入,开启回读校验时失败能返回错误。 |
| 断线重连 | 拔线、关闭服务端或断开串口后变量离线,恢复后能自动上线。 |
| 暂停/恢复 | 暂停设备后任务停止,恢复后继续采集。 |
| 释放重启 | 多次保存设备或重启服务后没有重复订阅、端口占用、线程泄漏。 |