业务插件二开说明
本文把继承 DataForwardBase 的数据转发目标、协议服务端、数据库写入、云平台对接、同步桥接统一称为“业务插件”。业务插件不负责从现场设备采集原始数据,而是消费采集变量、设备状态、报警或插件事件,把它们写入外部系统,或对外提供服务端协议能力。
如果要扩展 PLC、仪表、OPC、MQTT 采集等设备驱动,请阅读 采集插件二开说明。
源码入口
| 入口 | 作用 |
|---|---|
ThingsGatewayRuntime.Application/Driver/DataForward/DataForwardBase.cs | 所有业务插件的运行时基类。 |
ThingsGatewayRuntime.Application/Driver/DataForward/DataForwardPropertyBase.cs | 目标级插件属性根类。 |
ThingsGatewayRuntime.Application/Driver/DataForward/DataForwardVariablePropertyBase.cs | 目标变量级属性根类。 |
ThingsGatewayRuntime.Application/Driver/DataForward/DataForwardChannelPropertyBase.cs | 需要 TCP、串口、DTU、SSL 的目标属性基类。 |
ThingsGatewayRuntime.Application/Driver/DataForward/Cache/* | 缓存队列、离线文件缓存、Topic 模板、脚本模型、批量上传基类。 |
ThingsGatewayRuntime.Application/Task/DataForward/DataForwardMange/DataForwardMange.cs | 转发组、目标启动、事件分发、组级数据生产、冗余切换的调度入口。 |
ThingsGatewayRuntime.Application/Model/DataForwardGroupRuntime.cs | 转发组范围解析后的运行态变量和设备索引。 |
ThingsGatewayRuntime.Plugin/Plugin/* | MQTT、Kafka、RabbitMQ、Webhook、ModbusSlave、SyncBridge 等实现。 |
ThingsGatewayRuntime.NOAOTPlugin/Plugin/OpcUa/OpcUaServer | OPC UA Server 业务插件实现。 |
ThingsGatewayRuntimePRO/src/ThingsGatewayRuntime.NOAOTPROPlugin/Plugin/* | IEC104、IEC61850 Server、ZeroMQ 等专业版实现。 |
业务插件和采集插件的边界
| 项目 | 采集插件 | 业务插件 |
|---|---|---|
| 基类 | CollectBase | DataForwardBase |
| 配置位置 | 采集设备 | 数据转发目标 |
| 数据来源 | PLC、仪表、外部上报、虚拟计算 | 运行时变量、设备状态、报警、插件事件 |
| 主要任务 | 读写现场点位 | 上传、存储、服务端协议、桥接、反写 |
| 变量范围 | 设备下启用变量 | 转发组解析后的变量范围 |
| 写入入口 | InvokeWriteAsync、InvokeMethodAsync | 目标协议反写后通常再调用采集侧 RPC |
不要把业务插件写成“再去全局扫描所有变量”。变量范围由转发组决定,业务插件应通过 GetVariables()、IdVariableRuntimes、CollectDevices 或组生产入口读取当前目标可见的数据。
插件发现规则
| 规则 | 说明 |
|---|---|
继承 DataForwardBase | PluginService 只把 DataForwardBase 的非抽象派生类识别为数据转发插件。 |
| 公开无参构造 | 运行时通过 Activator.CreateInstance 创建,每个目标一个独立实例。 |
| 插件全名是配置键 | 数据转发目标保存类型 FullName,重命名命名空间或类名会影响旧配置。 |
| 目标属性先注入再初始化 | DataForwardMange.StartTargetAsync 先调用 SetDriverProperties,再 InitTarget、InitAsync、StartAsync。 |
| 目标变量属性不决定变量范围 | 目标变量属性只保存外部映射、权限、数据类型等插件专属配置。变量是否进入目标由转发组范围和组变量关系决定。 |
UsesGroupDataProducer 决定数据入口 | 缓存型目标通常为 true,消费转发组统一生产的数据;服务端协议类通常为 false,直接处理变量变化。 |
基类选择
| 基类 | 适合场景 | 必须重点实现 | 现有参考 |
|---|---|---|---|
DataForwardBase | 协议服务端、同步桥、需要自管内存映射或连接循环的目标。 | TargetProperties、TargetPropertyType、可选 VariablePropertyType、ProtectedInitAsync、AfterVariablesChangedAsync、OnVariableChanged、ProtectedExecuteAsync、IsConnected。 | ModbusSlave、OpcUaServer、IEC104Slave、IEC61850Server、SyncBridge。 |
DataForwardBaseWithCache | 需要离线缓存,但触发模型不完全等同周期变量上传的目标。 | DataForwardPropertyWithCache、启用的模型、Update*Model 或 AcceptProduced*。 | HisAlarmForwardTarget。 |
DataForwardBaseWithCacheInterval | 消费转发组统一生产的数据,支持周期/变化/批处理和离线缓存。 | DataForwardPropertyWithCacheInterval、模型开关、AcceptProduced* 或 Update*Model。 | HisDataForwardTarget、RealDataForwardTarget、ThingsBoardClientProducer。 |
DataForwardBaseWithCacheIntervalScript | 需要 Topic 模板、实体脚本、自定义上传模板,但只想复用转换能力。 | 继承后按模型调用 GetVariableBasicDataTopicArray 等方法。 | 脚本上传类目标的父类。 |
DataForwardBaseWithCacheIntervalScriptAll | MQTT、Kafka、RabbitMQ、Webhook、ZeroMQ 这类“生成 Topic/Payload 后上传”的目标。 | 实现 Upload(TopicArray, CancellationToken),初始化客户端连接,重写 IsConnected。 | MqttClientProducer、KafkaProducer、RabbitMQProducer、Webhook、ZeroMQProducer。 |
DataForwardChannelPropertyBase | 目标需要 TCP 客户端/服务端、串口、DTU、SSL 参数。 | 属性类继承它,并在 ProtectedInitAsync 调用 InitChannelAsync。 | ModbusSlaveProperty、IEC104SlaveProperty。 |
如果目标只是“把变量转成 JSON 发到某个系统”,优先使用 DataForwardBaseWithCacheIntervalScriptAll。如果目标要对外模拟一个协议服务端,让外部系统主动读写内存点表,优先使用 DataForwardBase。
目标属性和变量属性
| 类型 | 作用 |
|---|---|
DataForwardPropertyBase | 目标级属性根类。连接地址、认证、表名、缓存、模板都属于目标属性。 |
DataForwardVariablePropertyBase | 目标变量级属性根类。单个变量在某个目标下的外部地址、数据类型、写入权限等属于这里。 |
DataForwardVariableProperty | 通用预留字段 数据1 到 数据10,适合简单外部映射。 |
DataForwardPropertyWithCache | 离线缓存、内存队列上限、上传分片、过滤离线数据、上传并发。 |
DataForwardPropertyWithCacheInterval | 继承缓存属性,触发模式和周期由转发组统一决定。 |
DataForwardPropertyWithCacheIntervalScript | 增加详细日志、JSON 格式、列表/字典上传、Topic 模板、实体脚本、上传模板。 |
DataForwardChannelPropertyBase | 增加通道类型、远程地址、本地绑定、SSL、串口、心跳、DTU、并发、连接超时等。 |
实现目标属性时要同时保证目标属性实例和类型一致:
private readonly MyTargetProperty _properties = new();
public override DataForwardPropertyBase TargetProperties => _properties;
public override Type TargetPropertyType { get; } = typeof(MyTargetProperty);
如果插件有变量级配置,再提供:
private readonly MyVariableProperty _variableProperties = new();
public override DataForwardVariablePropertyBase VariablePropertys => _variableProperties;
public override Type? VariablePropertyType { get; } = typeof(MyVariableProperty);
DataForwardVariablePropertyBase.Enable 不是动态属性,页面不会出现第二个启用开关。是否转发某个变量由转发组变量范围、组内变量启用状态、目标启用状态和触发条件共同决定。
DynamicProperty 约定
业务插件和采集插件使用同一套动态属性规则。
| 项目 | 说明 |
|---|---|
只有 [DynamicProperty] 会显示和保存 | 运行时连接对象、客户端对象、缓存字段、临时索引不要标记。 |
| 默认值会进入新建表单 | 默认端口、Topic、表名、缓存关闭/开启策略应保守。 |
CertificatePurpose 用于证书下拉 | MQTT、OPC UA、IEC61850 TLS 等证书字段应标明 Client、Server 或 CA。 |
Remark 要写清单位和格式 | 例如毫秒、秒、Topic 模板、SQL 表名、连接字符串格式。 |
| 复杂对象也可以作为属性 | 例如上传模板配置、SOE 配置,但要保证 JSON 序列化和导入导出能恢复。 |
生命周期钩子
| 阶段 | 运行时动作 | 插件开发注意事项 |
|---|---|---|
| 构造函数 | 创建目标插件实例。 | 只初始化轻量字段,不连接外部系统。 |
| 属性注入 | SetDriverProperties(forwarder.TargetProperties, forwarder.TargetPropertyType, target.TargetPropertys)。 | ProtectedInitAsync 可以读取完整目标属性。 |
InitTarget | 挂载 CurrentGroup、CurrentTarget、日志、target.Forwarder,调用 ProtectedInitTarget。 | 可以缓存组/目标基础信息,不要启动连接。 |
InitAsync | 调用 ProtectedInitAsync,再 AfterVariablesChangedAsync。 | 解析属性、初始化缓存、创建客户端配置、构建变量映射。 |
StartAsync | 调用 ProtectedStartAsync,受 StartTimeout 控制,成功后设置 IsStarted。 | 连接外部系统或启动服务端监听,失败要抛异常或返回失败状态。 |
GetTasks | 建立目标调度循环。默认按组周期调用 ProtectedExecuteAsync。 | 大多数缓存型目标复用父类任务;服务端目标可在执行中刷新内存或检查连接。 |
| 初始快照 | 目标启动后,转发管理器会给组生产型目标推一次当前快照。 | 不要假设必须等下一次变量变化才有数据。 |
| 变量/设备/报警/事件变化 | 根据 UsesGroupDataProducer 走 AcceptProduced* 或 On*Changed。 | 缓存型目标不要自己重复订阅全局事件。 |
StopAsync/SafetyDisposeAsync | 停止调度循环,释放通道、日志、缓存、客户端和服务端。 | 必须关闭连接、释放缓存、取消订阅,避免重启后重复发送。 |
转发组是变量范围的唯一入口
DataForwardGroupRuntime.RebuildVariables() 会根据组配置生成两个运行态索引。
| 索引 | 说明 |
|---|---|
IdVariableRuntimes | 当前组范围解析后的变量,key 为变量 Id。 |
CollectDevices | 当前组范围内变量反推出来的采集设备。 |
业务插件应通过这些入口读取数据:
| 入口 | 适合场景 |
|---|---|
GetVariables() | 遍历当前目标可见变量,最常用。 |
IdVariableRuntimes | 按变量 Id 快速查找。 |
CollectDevices | 生成设备快照、协议服务端节点或设备连接状态。 |
GetVariableProperty<TProperty>(variable) | 读取当前目标下某个变量的强类型目标变量属性。 |
GetVariablePropertyValue(variable, propertyName) | 读取目标变量属性单字段,未配置时回退到默认属性实例。 |
TryGetVariablePropertyValue(...) | 需要区分“未配置”和“配置为空值”时使用。 |
目标变量属性不能把变量“拉进”转发范围。现场出现“变量有值但目标没有发送”时,先查转发组范围、组内变量启用、触发模式、目标启用和目标连接。
数据入口:直接事件和组生产
| 模式 | UsesGroupDataProducer | 数据入口 | 适合插件 |
|---|---|---|---|
| 直接事件 | false | OnVariableChanged、OnDeviceChanged、OnAlarmChanged、OnPluginEventChanged | ModbusSlave、IEC104Slave、OPC UA Server、IEC61850Server、SyncBridge。 |
| 组生产 | true | AcceptProducedVariableChange、AcceptProducedVariableSnapshot、AcceptProducedDeviceSnapshot、AcceptProducedAlarm、AcceptProducedPluginEvent | MQTT、Kafka、RabbitMQ、Webhook、ZeroMQ、历史数据、实时数据、历史报警。 |
缓存型基类会把 UsesGroupDataProducer 固定为 true。转发管理器会根据组触发模式、批处理模式和最大批量统一生产数据,再调用 AcceptProduced*。这样多个目标不会各自重复做范围解析和批处理。
缓存和失败语义
DataForwardBaseWithCache 为每种模型维护独立的内存队列和文件缓存。
| 属性 | 说明 |
|---|---|
CacheEnable | 是否启用离线文件缓存。关闭时发送失败的数据按 at-most-once 语义丢弃。 |
CacheFileMaxLength | 单个缓存文件最大行数,超过后删除旧数据。 |
SplitSize | 上传或补发时的批量拆分大小。 |
QueueMaxCount | 内存队列上限,超过后优先落盘,无法落盘时丢弃旧数据。 |
OnlineFilter | 是否过滤离线变量。 |
Concurrency | Topic 上传并发数量,由具体插件发送实现使用。 |
发送方法返回失败时,缓存基类会根据 CacheEnable 决定是否落盘。脚本类上传目标还会把已经生成但未确认的 TopicArray 写入 Topic outbox,连接恢复后 ReplayTopicUploadCache 先补发旧数据,再处理新数据。
Topic、脚本和上传模板
DataForwardBaseWithCacheIntervalScriptAll 已经处理了这些工作:
| 能力 | 说明 |
|---|---|
| Topic 模板 | ThingsGateway/Variable/${DeviceName} 这类模板会按实体属性分组并替换。 |
| 实体脚本 | 变量、设备、报警、插件事件都可以通过动态模型脚本重塑字段。 |
| 上传模板 | 配置内容模板后,用 ${属性名} 生成自定义 payload。 |
| 列表/逐条 | IsVariableList、IsDeviceList 等控制合并发送还是逐条发送。 |
| 字典上传 | 变量/报警可转换为 DeviceName -> Name -> 实体 的字典结构。 |
| JSON 格式 | 缩进和忽略 null 由属性控制。 |
这类插件通常只需要实现 Upload(TopicArray topicArray, CancellationToken cancellationToken)。
脚本上传目标模板
public sealed class MyProducer : DataForwardBaseWithCacheIntervalScriptAll
{
private readonly MyProducerProperty _properties = new();
private readonly DataForwardVariableProperty _variableProperties = new();
private MyClient? _client;
private bool _success = true;
protected override DataForwardPropertyWithCacheIntervalScript DataForwardPropertyWithCacheIntervalScript => _properties;
public override Type TargetPropertyType { get; } = typeof(MyProducerProperty);
public override Type VariablePropertyType { get; } = typeof(DataForwardVariableProperty);
public override DataForwardVariablePropertyBase VariablePropertys => _variableProperties;
public override bool IsConnected() => _success && _client?.Connected == true;
protected override async Task ProtectedInitAsync(CancellationToken cancellationToken)
{
_client = new MyClient(_properties.Endpoint, _properties.Token);
await base.ProtectedInitAsync(cancellationToken).ConfigureAwait(false);
}
protected override async Task ProtectedStartAsync(CancellationToken cancellationToken)
{
await _client!.ConnectAsync(cancellationToken).ConfigureAwait(false);
await base.ProtectedStartAsync(cancellationToken).ConfigureAwait(false);
}
protected override async ValueTask<OperResult> Upload(TopicArray topicArray, CancellationToken cancellationToken)
{
try
{
await _client!.PublishAsync(topicArray.Topic, topicArray.Payload.Memory, cancellationToken).ConfigureAwait(false);
_success = true;
return OperResult.Success;
}
catch (Exception ex)
{
_success = false;
return new OperResult(ex);
}
}
protected override async Task SafetyDisposeAsync(bool disposing)
{
if (_client != null)
await _client.DisposeAsync().ConfigureAwait(false);
await base.SafetyDisposeAsync(disposing).ConfigureAwait(false);
}
}
public sealed class MyProducerProperty : DataForwardPropertyWithCacheIntervalScript
{
[DynamicProperty("服务地址")]
public string Endpoint { get; set; } = "http://127.0.0.1:8080";
[DynamicProperty("访问令牌")]
public string Token { get; set; } = string.Empty;
}
协议服务端目标模板
协议服务端目标通常不使用脚本上传基类,而是维护一份外部协议内存映射。
public sealed class MyServerTarget : DataForwardBase
{
private readonly MyServerProperty _properties = new();
private readonly MyServerVariableProperty _variableProperties = new();
private readonly ConcurrentQueue<VariableRuntime> _changedVariables = new();
private MyServer? _server;
public override DataForwardPropertyBase TargetProperties => _properties;
public override Type TargetPropertyType { get; } = typeof(MyServerProperty);
public override DataForwardVariablePropertyBase VariablePropertys => _variableProperties;
public override Type VariablePropertyType { get; } = typeof(MyServerVariableProperty);
public override bool IsConnected() => _server?.IsRunning == true;
protected override async Task ProtectedInitAsync(CancellationToken cancellationToken)
{
_server = new MyServer(_properties.BindUrl);
await base.ProtectedInitAsync(cancellationToken).ConfigureAwait(false);
}
public override async Task AfterVariablesChangedAsync(CancellationToken cancellationToken)
{
await base.AfterVariablesChangedAsync(cancellationToken).ConfigureAwait(false);
foreach (var variable in GetVariables())
{
var map = GetVariableProperty<MyServerVariableProperty>(variable);
if (map != null)
_server!.Register(map.Address, variable.DataType);
}
}
protected override async Task ProtectedStartAsync(CancellationToken cancellationToken)
{
await _server!.StartAsync(cancellationToken).ConfigureAwait(false);
await base.ProtectedStartAsync(cancellationToken).ConfigureAwait(false);
}
public override void OnVariableChanged(VariableRuntime variableRuntime, VariableBasicData variableData, CancellationToken cancellationToken)
{
if (cancellationToken.IsCancellationRequested)
return;
_changedVariables.Enqueue(variableRuntime);
}
protected override Task ProtectedExecuteAsync(object? state, CancellationToken cancellationToken)
{
while (_changedVariables.TryDequeue(out var variable))
{
var map = GetVariableProperty<MyServerVariableProperty>(variable);
if (map != null)
_server!.SetValue(map.Address, variable.Value);
}
return Task.CompletedTask;
}
}
ModbusSlave、IEC104Slave、OpcUaServer、IEC61850Server 都属于这种思路:变量变化先进入插件维护的映射或队列,再由协议服务端对外提供读写。
当前源码业务插件参考矩阵
下表按当前源码核对了 15 个业务插件。开发新目标前,优先找最接近的类型参考。
| 插件 | 基类 | 目标属性 | 变量属性 | 参考重点 |
|---|---|---|---|---|
HisDataForwardTarget | DataForwardBaseWithCacheInterval | HisDataForwardProperty | HisDataForwardVariableProperty | 历史数据入库、采样策略、条件表达式、分表、保留天数。 |
RealDataForwardTarget | DataForwardBaseWithCacheInterval | RealDataForwardProperty | 无 | 实时数据表覆盖写入、周期/变化触发、数据库写入器。 |
HisAlarmForwardTarget | DataForwardBaseWithCache | HisAlarmForwardProperty | 无 | 历史报警入库、报警等级过滤、缓存补写。 |
ModbusSlave | DataForwardBase | ModbusSlaveProperty | ModbusSlaveVariableProperty | Modbus 从站内存映射、通道属性、变量地址、RPC 写入权限。 |
MqttClientProducer | DataForwardBaseWithCacheIntervalScriptAll | MqttClientProducerProperty | MqttClientProducerVariableProperty | MQTT 客户端上传、Topic/Payload 模板、订阅 RPC 写入。 |
MqttServerProducer | DataForwardBaseWithCacheIntervalScriptAll | MqttServerProducerProperty | MqttServerProducerVariableProperty | MQTT 服务端上传、客户端连接管理、Topic/Payload 模板。 |
ThingsBoardClientProducer | DataForwardBaseWithCacheInterval | ThingsBoardClientProducerProperty | ThingsBoardClientProducerVariableProperty | ThingsBoard Gateway 协议、遥测上传、设备连接/断开、RPC 写入。 |
KafkaProducer | DataForwardBaseWithCacheIntervalScriptAll | KafkaProducerProperty | DataForwardVariableProperty | Kafka Producer 初始化、TopicArray 上传、并发和缓存。 |
RabbitMQProducer | DataForwardBaseWithCacheIntervalScriptAll | RabbitMQProducerProperty | DataForwardVariableProperty | RabbitMQ 连接、Exchange/RoutingKey、TopicArray 上传。 |
Webhook | DataForwardBaseWithCacheIntervalScriptAll | WebhookProperty | DataForwardVariableProperty | HTTP Webhook、请求方法、Header、Payload 模板、失败缓存。 |
SyncBridge | DataForwardBase、IRpcDriver | SyncBridgeProperty | SyncBridgeVariableProperty | 网关间同步桥、变量变化队列、RPC 反写代理。 |
OpcUaServer | DataForwardBase | OpcUaServerProperty | OpcUaServerVariableProperty | OPC UA 服务端节点、证书、安全策略、变量写入权限。 |
IEC104Slave | DataForwardBase | IEC104SlaveProperty | IEC104SlaveVariableProperty | IEC104 从站、遥信/遥测/位数组映射、通道属性。 |
IEC61850Server | DataForwardBase | IEC61850ServerProperty | IEC61850ServerVariableProperty | IEC61850 Server 建模、数据集、写入索引、对象类型转换。 |
ZeroMQProducer | DataForwardBaseWithCacheIntervalScriptAll | ZeroMQProducerProperty | DataForwardVariableProperty | ZeroMQ Push/Pub/Dealer 上传、绑定模式、TopicArray 和缓存。 |
插件类型参考
| 类型 | 推荐参考 |
|---|---|
| 外部消息系统上传 | MqttClientProducer、KafkaProducer、RabbitMQProducer、ZeroMQProducer。 |
| HTTP/REST 推送 | Webhook。 |
| 工业协议服务端 | ModbusSlave、IEC104Slave、OpcUaServer、IEC61850Server。 |
| 数据库存储 | HisDataForwardTarget、RealDataForwardTarget、HisAlarmForwardTarget。 |
| 平台专用协议 | ThingsBoardClientProducer。 |
| 网关间桥接和反写代理 | SyncBridge。 |
常见开发错误
| 错误 | 后果 | 正确做法 |
|---|---|---|
| 在构造函数连接外部系统 | 属性还没注入,日志和取消令牌不可用。 | 连接放到 ProtectedStartAsync,客户端配置放到 ProtectedInitAsync。 |
TargetProperties 和 TargetPropertyType 不一致 | 表单字段和运行时对象错位,导入导出异常。 | 两者始终对应同一个强类型属性。 |
| 把目标变量属性当成范围过滤 | 变量变化不会进入目标或排查方向错误。 | 变量范围只由转发组决定,目标变量属性只做外部映射。 |
| 缓存型目标自己订阅全局变量事件 | 重复发送,批处理和触发模式失效。 | 继承缓存基类后使用 AcceptProduced* 和 Update*Model。 |
| 发送失败仍返回成功 | 离线缓存无法接管,数据丢失且日志误导。 | 发送失败返回失败 OperResult。 |
IsConnected 永远返回 true | 前端状态和缓存补发判断失真。 | 根据客户端、socket、server 或 channel 真实状态返回。 |
忽略 CancellationToken | 停止目标、切换冗余或刷新配置时卡住。 | 所有连接、发送、等待、补发都传入取消令牌。 |
| 释放时不关闭客户端/服务端 | 端口占用、重复连接、文件缓存未释放。 | SafetyDisposeAsync 中关闭连接、释放缓存,再调用 base。 |
上传模板不释放 TopicArray | 高频上传产生内存压力。 | 使用父类 UpdateTopicArrays,不要绕过其释放逻辑。 |
验证清单
| 场景 | 验证点 |
|---|---|
| 插件发现 | 插件管理和数据转发目标下拉能看到插件,类型为数据转发。 |
| 属性表单 | 目标属性、变量属性、默认值、证书、枚举、导入导出都正确。 |
| 范围解析 | Manual、All、CollectDevice、CollectGroup 模式下,GetVariables() 数量符合预期。 |
| 变化触发 | 单点变化能进入目标,组变量 参与组触发、更新模式 生效。 |
| 周期触发 | 周期快照按组间隔、批处理模式和最大批量执行。 |
| 外部连接 | 目标连接失败时 LastErrorMessage 清楚,恢复后能重连或重启成功。 |
| 离线缓存 | 外部系统断开时失败数据进入缓存,恢复后按批次补发。 |
| 反写/RPC | 服务端协议或平台 RPC 反写能落到采集变量,权限关闭时拒绝。 |
| 目标重启 | 修改属性保存、多次启停后没有重复订阅、端口占用或重复发送。 |
| 冗余目标 | 启用冗余时主备切换不丢运行态引用,失败目标能清理干净。 |