跳到主要内容

采集插件二开说明

本文面向需要扩展 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/*专业版采集插件实现。

一句话流程

运行时按下面顺序处理采集插件:

  1. PluginService 扫描所有非抽象 CollectBase 派生类,形成插件列表。
  2. 设备启动时,DeviceThreadManage.CreateDriver 通过插件全名创建实例。
  3. DriverBase.InitDevice 挂载 DeviceRuntime、日志和 device.Driver
  4. PluginServiceUtil.SetDriverProperties 把设备插件属性字典写回强类型属性对象。
  5. DriverBase.InitChannelAsync 挂载通道,调用 AfterVariablesChangedAsync 打包变量。
  6. DriverBase.StartAsync 调用插件的 ProtectedStartAsync,再创建并启动 TaskSchedulerLoop
  7. CollectBase 按变量间隔读取,失败重试,成功后变量上线,失败后变量离线。
  8. 外部写入/RPC 进入 InvokeWriteAsyncInvokeMethodAsync,由插件落到协议写入。
  9. 停止设备时调用 StopAsync,最终进入 SafetyDisposeAsync 释放通道、底层协议对象、日志和锁。

插件发现规则

采集插件必须满足这些条件才会出现在采集设备配置中。

规则说明
继承 CollectBasePluginService 只把 CollectBase 的非抽象派生类识别为采集插件。
公开无参构造运行时通过 Activator.CreateInstance 创建实例;没有无参构造会启动失败。
不要在构造函数连接现场设备构造阶段还没有设备属性、通道、日志和取消令牌。连接应放在 InitChannelAsyncProtectedStartAsync
插件全名是配置键设备实体保存的是类型 FullName,重命名命名空间或类名会影响旧配置。
[DisplayNonePlugin] 会隐藏插件MemoryDriver 使用该标记,作为内部内存设备模板,不作为普通现场插件显示。
[OnlyWindowsSupport] 会限制平台非 Windows 环境下带该属性的插件不会显示。
NOAOT 程序集不支持 AOTPluginService 根据程序集名标记 SupportsAot,OPC/COM 等插件通常位于 NOAOT 程序集。

基类选择

基类适合场景必须重点实现现有参考
CollectFoundationBase协议可以抽象成按地址读写字节,例如 Modbus、S7、DLT645、Omron、Melsec。FoundationDeviceCollectPropertiesDriverPropertyTypeInitChannelAsyncProtectedLoadSourceReadAsync。通常不需要重写 ReadSourceAsyncModbusMasterSiemensS7MasterControlLogixMaster
CollectReceivedFoundationBase数据由对端主动上报或底层设备维护连接状态,普通定时读不成立。FoundationReceivedDevice、属性类型、上报事件到变量的映射,必要时重写 AfterVariablesChangedAsyncHJ212MasterEDPF_NTMasterKELID2008MasterLKSISMaster
CollectBase自定义协议、订阅协议、消息协议、虚拟设备、扫码器、CAN、OPC、IEC61850、ZeroMQ。ProtectedLoadSourceReadAsyncReadSourceAsyncWriteValuesAsyncIsConnected,必要时重写 ProtectedStartAsyncAfterVariablesChangedAsyncOpcUaMasterCanMasterIEC61850MasterMemoryDriverMqttCollectBase
MqttCollectBaseMQTT 客户端/服务端采集,变量值由 Topic 消息更新。派生类处理连接、订阅、消息接收和在线判断;父类负责按变量关系更新值。MqttCollectClientMqttCollectServer

不要为了“省代码”继承过高层的基类。能用 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,最后 AfterVariablesChangedAsyncFoundation 插件通常在这里重建底层 IDevice,把属性赋给底层对象,并调用 InitChannel
AfterVariablesChangedAsync基类重建 VariableSourceReads、脚本变量、方法变量和定时任务。变量新增、删除、属性修改都会触发。重写时注意调用 base,并清理旧订阅或旧映射。
ProtectedStartAsync启动通讯前调用,受 StartTimeout 控制。连接远端、启动服务、订阅消息可放这里。必须尊重取消令牌。
ProtectedGetTasks构建调度任务。CollectBase 已提供设备状态、在线测试、变量读取任务。只有特殊协议才重写。
定时读取调用 ReadSourceAsync,失败按 RetryCount 重试。成功要设置变量值或返回可解析字节;失败要返回失败结果,不要吞异常后假成功。
写入/RPCInvokeWriteAsyncInvokeMethodAsync自定义写入要使用写锁,必要时执行写后回读校验。
StopAsync/SafetyDisposeAsync停止任务循环,释放日志、底层设备、事件订阅、锁。必须取消事件订阅、释放 socket/client/channel,避免重启后重复接收。

变量读取模型

CollectBase.AfterVariablesChangedAsync 会把启用变量分成三类。

类型判断方式运行方式
普通源读取变量地址不是 DeviceStatusScriptScriptRead,且没有 OtherMethod进入 ProtectedLoadSourceReadAsync,生成 VariableSourceRead 后由定时任务调用 ReadSourceAsync
脚本/特殊变量地址为 DeviceStatusScriptScriptRead进入 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首次连接或启动服务的位置。超时后会记录启动失败。
TestOnlineCollectReceivedFoundationBase 默认会尝试重连并把变量置离线;自定义协议可重写做轻量检测。
SetDeviceStatusCollectBase 根据连接状态和变量在线情况更新设备状态。特殊插件可重写,例如 MemoryDriver 永久在线。

现场上“设备在线”不等于“每个变量都在线”。如果某个读包失败,CollectBase 会把该读包内变量置为离线,并记录最后错误。

特殊地址和动态方法

能力写法场景
设备状态变量变量地址写 DeviceStatus把设备在线/离线状态作为普通变量参与展示或转发。
脚本变量地址写 ScriptScriptRead不走现场通讯,通过表达式或脚本计算值。
动态方法方法标记 [DynamicMethod("说明", "备注")]读写日期、调用协议特有命令、执行对象方法。
RPC 写入外部接口写变量。进入 InvokeWriteAsyncInvokeMethodAsync,再由插件写现场设备。

动态方法返回值必须能转成 OperResultIOperResult<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 个采集插件。开发新插件前,先找协议形态最接近的实现参考。

插件基类属性类参考重点
MemoryDriverCollectBaseCollectPropertyNone内存变量、表达式触发、全局变量变化事件、无真实通道。
ModbusMasterCollectFoundationBaseModbusMasterPropertyFoundation 标准模板、DTU、站号、最大打包长度。
SiemensS7MasterCollectFoundationBaseSiemensS7MasterPropertyS7 地址打包、批量写、日期/时间动态方法。
Dlt645_2007MasterCollectFoundationBaseDlt645_2007MasterProperty仪表地址、密码、操作员代码、前导符。
OpcUaMasterCollectBaseOpcUaMasterProperty订阅/轮询混合、证书、安全策略、节点类型加载、订阅刷新。
OpcDaMasterCollectBaseOpcDaMasterPropertyWindows/COM OPC DA、订阅组、重连检查、服务端时间。
MqttCollectClientMqttCollectBaseMqttCollectClientPropertyMQTT 客户端连接、订阅 Topic、消息映射变量。
MqttCollectServerMqttCollectBaseMqttCollectServerPropertyMQTT 服务端监听、客户端校验、消息映射变量。
CanMasterCollectBaseCanMasterPropertyCAN/CAN FD 端点、帧组装、硬件过滤、全帧切片读取。
ControlLogixMasterCollectFoundationBaseControlLogixMasterPropertyAllenBradley CIP,Foundation 读包,自定义写入。
PCCCMasterCollectFoundationBasePCCCMasterPropertyAllenBradley PCCC,Foundation 读写模板。
DCONMasterCollectFoundationBaseDCONMasterPropertyDCON 协议,DTU 和打包模板。
EDPF_NTMasterCollectReceivedFoundationBaseEDPF_NTUdpPropertyUDP 被动接收,变量刷新时重建映射,不实现主动读写。
HJ212MasterCollectReceivedFoundationBaseHJ212MasterProperty环保 HJ212 上报,地址说明和上报映射。
IEC61850MasterCollectBaseIEC61850MasterPropertyMMS 读写、RCB 报告、GOOSE、SOE、TLS、复杂订阅生命周期。
InovanceMasterCollectFoundationBaseInovanceMasterProperty汇川协议,继承 Foundation DTU 打包属性。
KELID2008MasterCollectReceivedFoundationBaseKELID2008MasterProperty上报协议,地址说明扩展,主动读写未实现。
LKSISMasterCollectReceivedFoundationBaseLKSISPropertyUDP 上报协议,变量变更后维护接收映射。
Mc1E_BinaryMasterCollectFoundationBaseMc1E_BinaryMasterProperty三菱 1E 二进制,Foundation 打包模板。
Mc3E_BinaryMasterCollectFoundationBaseMc3E_BinaryMasterProperty三菱 3E 二进制,Foundation 打包模板。
ModbusC1MasterCollectFoundationBaseModbusC1MasterPropertyModbus C1,DTU 打包模板。
ModbusC20MasterCollectFoundationBaseModbusC20MasterPropertyModbus C20,DTU 打包模板。
OmronFinsMasterCollectFoundationBaseOmronFinsMasterPropertyOmron FINS,Foundation 打包模板。
OpcAeMasterCollectBaseOpcAeMasterPropertyOPC AE 事件采集,事件订阅转变量。
SECSMasterCollectReceivedFoundationBaseSECSMasterPropertySECS/GEM,既有接收设备,也实现了主动读写。
TIANXINMasterCollectFoundationBaseTIANXINMasterProperty天信仪表类协议,DTU 打包模板。
TS550MasterCollectFoundationBaseTS550MasterPropertyTS550 协议,Foundation 通用读写属性。
USBScanerCollectBaseUSBScanerProperty本机扫码器 Hook,扫码事件写入变量,不实现主动读写。
VigorMasterCollectFoundationBaseVigorMasterPropertyVigor 协议,DTU 打包模板。
ZeroMQCollectClientCollectBaseZeroMQCollectClientPropertyZeroMQ 订阅/连接、Topic 前缀、高水位、重连和清理。

常见开发错误

错误后果正确做法
构造函数里读取属性或连接设备属性还没注入,日志和取消令牌也不存在。构造函数只初始化字段,连接放到 InitChannelAsyncProtectedStartAsync
CollectPropertiesDriverPropertyType 不一致页面字段、导入导出和运行时读取错位。两者都指向同一个强类型属性类。
重写 InitChannelAsync 后不调用 base变量不会打包,定时任务没有数据源。完成底层对象初始化后调用 base.InitChannelAsync
变量变化后不清理旧订阅重复接收、重复写变量、内存泄漏。AfterVariablesChangedAsync 或释放流程中先移除旧映射。
读失败仍返回成功变量看似在线但值错误。返回失败 OperResult,让基类置离线并记录错误。
忽略 CancellationToken停止设备或重启时卡住。所有网络、串口、等待和重试都传入取消令牌。
写入不加写锁读写并发冲突,现场协议包交叉。自定义 WriteValuesAsync 使用 ReadWriteLock.WriterLockAsync
释放时不取消事件订阅重启后同一数据被处理多次。SafetyDisposeAsync 中取消订阅并释放底层对象。

验证清单

开发完成后至少验证这些场景。

场景验证点
插件发现插件管理能看到插件,采集设备下拉能选择,插件类型为采集。
属性表单所有 [DynamicProperty] 字段显示、默认值、枚举、证书下拉、导入导出都正确。
通道类型SupportedChannelTypes() 与现场连接方式一致;无普通通道时返回 Other
数据类型SupportedDataTypes() 不要暴露协议无法解析的类型。
单点读取使用设备调试先读 3 到 5 个点,覆盖布尔、整数、浮点、字符串或协议特有类型。
批量读取点表导入后确认打包数量、间隔和超时符合预期。
写入/RPC只读点不可写,可写点能写入,开启回读校验时失败能返回错误。
断线重连拔线、关闭服务端或断开串口后变量离线,恢复后能自动上线。
暂停/恢复暂停设备后任务停止,恢复后继续采集。
释放重启多次保存设备或重启服务后没有重复订阅、端口占用、线程泄漏。