跳到主要内容

脚本开发说明

本文面向需要扩展 ThingsGatewayRuntime 脚本能力的开发人员和交付工程师。它不是普通现场配置手册,而是说明脚本如何被编译、加载、注册、绑定和执行,以及每一种脚本类型应该怎么写。

源码入口

源码作用
ThingsGatewayRuntime.Application/Expressions/ExpressionInfo.cs脚本类型枚举、脚本元数据、输入输出参数定义。
ThingsGatewayRuntime.Application/Expressions/DynamicBase.cs数据转换、内存变量、动态模型、动态 SQL、MQTT RPC、完整源码的运行时基类。
ThingsGatewayRuntime.Application/Expressions/CustomExpressionBase.cs自定义节点脚本基类和输出变化回调。
ThingsGatewayScriptCompiler/ExpressionCodeGenerator.cs把页面中填写的脚本包装成 C# 类并编译为 DLL。
ThingsGatewayRuntime.ExpressionsGenerator/ExpressionRegistrationGenerator.cs编译时扫描脚本类,通过 ModuleInitializer 注册到 ExpressionsData
ThingsGatewayRuntime.Application/Expressions/ExpressionsData.cs保存已注册脚本委托,按名称查找脚本实例。
ThingsGatewayRuntime.Application/Expressions/ExpressionsHelper.cs运行时按脚本名称执行数据转换、动态模型、动态 SQL、MQTT RPC。
ThingsGatewayRuntime.Application/Controllers/GatewayScriptController.csWEB 脚本管理接口:创建、保存、编译、批量编译、删除、热加载。
ThingsGatewayRuntime.Application/Controllers/RuleEngine/GatewayCustomNodeController.cs自定义节点创建、保存、编译、热加载。
ThingsGatewayRuntime.Application/Task/RuleEngine/RuleEngineTask.cs规则流程启动、节点实例创建、输入输出传播和防抖触发。
ThingsGatewayRuntime.Application/Expressions/RuntimeDllLoader.cs启动时加载 PluginDllsScriptDllsCustomNodeDlls

总生命周期

  1. 在 WEB 中创建脚本或自定义节点,源码、分类、描述和参数写入数据库。
  2. 点击编译时,运行时调用外部 ThingsGatewayScriptCompiler 进程。
  3. 编译器按脚本类型包装源码:普通脚本写入 ScriptDlls,自定义节点写入 CustomNodeDlls
  4. 源生成器在脚本 DLL 内生成注册代码,模块加载时写入 ExpressionsData
  5. 编译成功后运行时热加载 DLL;服务启动时也会加载 PluginDllsScriptDllsCustomNodeDlls
  6. 功能配置按脚本名称绑定脚本,例如变量读表达式、内存变量表达式、数据转发实体脚本、历史表脚本、MQTT RPC 脚本、规则节点。
  7. 触发条件满足后执行脚本。脚本异常会影响对应功能,例如变量离线、转发批次失败、RPC 无响应或规则节点停止传播。
  8. 删除脚本时会删除数据库记录,并尝试删除对应 DLL;如果 DLL 正被占用,删除逻辑会尝试重命名为 .del

类型总览

类型页面类型值基类用户代码形态主要绑定位置
数据转换DataTransExpressionDatatransExecute 方法体采集变量读表达式、写表达式。
内存变量计算MemoryVariableDatatransMemoryVariableExpressionDatatransExecute 方法体内存变量读表达式。
动态模型-变量DynamicModel_VariableBasicDataDynamicModelBase<VariableBasicData>GetList 方法体数据转发目标的变量实体脚本。
动态模型-设备DynamicModel_DeviceBasicDataDynamicModelBase<DeviceBasicData>GetList 方法体数据转发目标的设备实体脚本。
动态模型-报警DynamicModel_AlarmVariableDynamicModelBase<AlarmVariable>GetList 方法体数据转发目标的报警实体脚本。
动态模型-插件事件DynamicModel_PluginEventDataDynamicModelBase<PluginEventData>GetList 方法体数据转发目标的插件事件实体脚本。
动态 SQL-变量DynamicSQL_VariableBasicDataDynamicSQLBase<VariableBasicData>类成员源码历史数据表脚本、实时数据表脚本。
动态 SQL-设备DynamicSQL_DeviceBasicDataDynamicSQLBase<DeviceBasicData>类成员源码框架支持;标准目标当前未发现直接消费入口,通常用于二开目标。
动态 SQL-报警DynamicSQL_AlarmVariableDynamicSQLBase<AlarmVariable>类成员源码历史报警表脚本。
动态 SQL-插件事件DynamicSQL_PluginEventDataDynamicSQLBase<PluginEventData>类成员源码框架支持;标准目标当前未发现直接消费入口,通常用于二开目标。
MQTT RPCMQTTDynamicRPCMQTTDynamicRPCBase类成员源码MQTT Client/Server 目标的 RPC 脚本。
完整源码FullSourceDynamicBase完整 C# 源码高级扩展,供自定义代码按名称获取实例。
自定义节点CustomNodeCustomExpressionBaseInitAsync/ChangedAsync 方法体或完整节点类规则引擎流程节点。

公共规则

脚本名称是运行时查找键。不要在生产环境频繁改名;变量、转发目标、MQTT RPC 或规则流程保存的是脚本名称。

页面脚本不是全部都写完整类。DataTransMemoryVariableDatatrans、四类动态模型写的是方法体;动态 SQL 和 MQTT RPC 写的是类成员;FullSource 写完整源码;自定义节点由节点页面根据输入输出参数生成属性,再把代码插入类中。

脚本默认可使用 SystemSystem.LinqSystem.Collections.GenericThingsGatewayRuntime.ApplicationThingsGatewayRuntimeTUtilityTUtility.ExtensionTUtility.Json.ExtensionTUtility.Log。如果需要额外命名空间,把 using 写在脚本最前面,编译器会把它移到包装类外部。

不要依赖“没有 return 时自动补 return”的包装逻辑。源码中确实有自动补齐,但它只做字符串判断,遇到注释或复杂代码不可靠;正式脚本请显式 return

带输入参数的脚本目前只对 DataTransMemoryVariableDatatrans 有运行时写入入口。参数会变成脚本类属性,变量配置中的参数值会在变量初始化时写入脚本实例。

动态加载依赖 RuntimeFeature.IsDynamicCodeSupported。AOT 或禁止动态代码的发布方式不能热加载脚本 DLL;这种环境应在发布前编译并验证扩展。

脚本实例通常会被复用。不要把单次执行的临时结果放到实例字段中,除非你明确需要跨执行保存状态并已考虑并发、重启和释放。

脚本里可以写日志,但高频变量脚本不要每次正常执行都写 Info 日志。建议只在异常、被过滤、格式不符合预期时写 Warning 或 Debug。

数据实体速查

实体常用字段说明
VariableBasicDataIdDeviceNameNameValueRawValueCollectTimeChangeTimeIsOnlineDataTypeUnitRegisterAddressCollectGroupRemark1Remark5变量当前值或历史采样值。Value 是脚本转换后的工程值,RawValue 是原始值。
DeviceBasicDataIdNameActiveTimeDeviceStatusLastErrorMessagePluginNameDescriptionRemark1Remark5设备状态变化或快照。
AlarmVariableAlarmIdVariableIdDeviceNameNameAlarmLevelAlarmTypeEventTypeAlarmTimeEventTimeFinishTimeConfirmTimeAlarmTextAlarmCodeAlarmLimit报警发生、恢复、确认等事件。
PluginEventDataDeviceNameObjectValue插件自定义事件,事件正文是 JsonElement

数据转换脚本

用途

数据转换脚本把一个变量的原始值转换为工程值,也可在外部写入设备前把工程值转换回设备值。

运行时机

读表达式在 VariableRuntime.SetValue 中执行,顺序是:采集插件读到值、写入 RawValue、执行读表达式、写入 Value。读表达式失败时变量会被置为离线并记录转换失败信息。

写表达式在 CollectBase.InvokeWriteAsyncInvokeMethodAsync 中执行,顺序是:收到外部写入值、执行写表达式、把转换后的值交给采集插件写入设备。

方法形态

public override object Execute(object raw, Logger? logger)
{
// 用户脚本写在这里。
}

页面中只写方法体,可以直接使用 rawlogger

Demo:模拟量比例换算

if (!TryConvertToDouble(raw, out var rawValue))
{
logger?.LogWarning($"AI value is not numeric: {raw}");
return raw;
}

var engineeringValue = rawValue * 0.1;
return Math.Round(engineeringValue, 2);

Demo:带参数的线性换算

在脚本输入参数中新增:

参数名类型初始值说明
ScaleDouble0.1比例系数。
OffsetDouble0偏移量。

脚本内容:

if (!TryConvertToDouble(raw, out var rawValue))
{
return raw;
}

return Math.Round(rawValue * Scale + Offset, 3);

注意事项

返回值类型要和变量数据类型匹配。例如变量配置为 Double,不要返回无法转换的字符串。

读表达式不要直接访问外部网络或数据库。它在采集热路径执行,阻塞会拖慢变量刷新。

写表达式用于反写前转换,失败会阻止本次写入。反写脚本里要优先做输入校验,错误信息写清楚。

内存变量计算脚本

用途

内存变量脚本用于计算网关内部变量,不直接访问 PLC。它可以读取其它变量,生成派生量、状态量、汇总量或联锁判断结果。

运行时机

内存变量由 MemoryDriver 管理。变量没有读表达式时按写入值保存;有读表达式时,运行时会按内存变量触发方式执行计算。触发方式可以按周期执行,也可以尝试根据 Tag(device, variable) 记录的依赖变量变化触发。

方法形态

public override object Execute(object raw, Logger? logger)
{
// 用户脚本写在这里。
}

可用 Tag("设备名", "变量名") 获取其它变量运行时对象。Tag 同时会把依赖加入 Tags 集合,运行时据此建立变化触发关系。

Demo:计算两路温度平均值

var left = Tag("PLC_1", "Temp_Left").Value;
var right = Tag("PLC_1", "Temp_Right").Value;

if (!double.TryParse(Convert.ToString(left), out var leftValue) ||
!double.TryParse(Convert.ToString(right), out var rightValue))
{
logger?.LogWarning("Temperature source value is not numeric.");
return null;
}

return Math.Round((leftValue + rightValue) / 2.0, 1);

Demo:运行允许状态

var autoMode = Convert.ToBoolean(Tag("PLC_1", "AutoMode").Value);
var emergencyStop = Convert.ToBoolean(Tag("PLC_1", "EmergencyStop").Value);
var pressureOk = Convert.ToBoolean(Tag("PLC_1", "PressureOk").Value);

return autoMode && !emergencyStop && pressureOk;

注意事项

Tag 的设备名和变量名必须是运行时名称,不是描述、地址或外部映射名。

不要把 Tag 放在不一定执行的分支里,否则依赖变量集合可能不完整,变化触发关系也可能不完整。

依赖变量不存在时 Tag 会抛异常,脚本计算失败。上线前先用少量变量验证依赖名称。

如果脚本依赖时间、累计值或外部状态,即使有 Tag,也建议保留周期触发,避免只有变量变化才计算。

动态模型脚本

用途

动态模型脚本在数据转发前重塑实体。它不负责连接外部系统,只负责把 VariableBasicDataDeviceBasicDataAlarmVariablePluginEventData 变成目标系统需要的对象、匿名对象或字典。

转发基类会继续使用脚本输出做 Topic 分组、默认 JSON 序列化或内容模板替换。Topic 模板中的 ${属性名} 会从脚本输出对象上读取属性。

方法形态

public override IEnumerable<object> GetList(IEnumerable<T> datas, Logger? logger)
{
// 用户脚本写在这里。
}

页面中只写方法体,必须返回 IEnumerable<object> 或可赋给它的集合。

Demo:变量动态模型

类型选择 DynamicModel_VariableBasicData

return datas.Select(data => new
{
device = data.DeviceName,
tag = data.Name,
value = data.Value,
raw = data.RawValue,
unit = data.Unit,
quality = data.IsOnline ? "Good" : "Bad",
ts = data.CollectTime
});

如果转发目标 Topic 写成 factory/${device}/${tag},运行时会按脚本输出的 devicetag 分组。

Demo:设备动态模型

类型选择 DynamicModel_DeviceBasicData

return datas.Select(device => new
{
device = device.Name,
status = device.DeviceStatus.ToString(),
online = device.DeviceStatus.ToString() == "OnLine",
plugin = device.PluginName,
activeTime = device.ActiveTime,
message = device.LastErrorMessage
});

Demo:报警动态模型

类型选择 DynamicModel_AlarmVariable

return datas.Select(alarm => new
{
id = alarm.AlarmId,
device = alarm.DeviceName,
tag = alarm.Name,
level = alarm.AlarmLevel,
alarmType = alarm.AlarmType?.ToString(),
eventType = alarm.EventType.ToString(),
text = alarm.AlarmText ?? alarm.AlarmCode,
eventTime = alarm.EventTime
});

Demo:插件事件动态模型

类型选择 DynamicModel_PluginEventData

return datas.Select(item => new
{
device = item.DeviceName,
eventBody = item.ObjectValue,
receivedAt = DateTime.UtcNow
});

注意事项

返回对象的属性名会影响 Topic 模板和内容模板。改字段名前先检查转发目标中的 ${属性名}

不要在动态模型里过滤掉大量数据,除非目标协议明确需要。范围过滤应该优先放在数据转发组中。

返回 Dictionary<string, object?> 也可以,但字段名要和模板保持一致。

动态模型脚本只改变上传内容,不改变运行时变量本身。

动态 SQL 脚本

用途

动态 SQL 脚本把数据库建表、删除和保存逻辑交给脚本。它用于项目要求固定表结构、宽表字段、特殊字段命名或特殊清理策略的场景。

方法形态

动态 SQL 不是方法体脚本,而是写入生成类中的成员源码。必须实现三个方法:

public override Task DBInit(OrmClient db, Logger? logger, CancellationToken cancellationToken);

public override Task<int> DBDeletable(OrmClient db, int days, Logger? logger, CancellationToken cancellationToken);

public override Task DBSavable(OrmClient db, IEnumerable<T> datas, Logger? logger, CancellationToken cancellationToken);

DBInit 在目标初始化表结构时执行;DBSavable 在批量写入时执行;DBDeletable 在保留天数清理时执行。

注意:QuestDB 和 TDengine 的历史数据目标使用各自数据库的 TTL/KEEP 保留机制;当目标配置了“历史表脚本”时,清理策略由脚本或数据库自行负责,不要只依赖目标上的“保留天数”设置。

Demo:历史变量宽表

类型选择 DynamicSQL_VariableBasicData,可绑定到历史数据转发目标的“历史表脚本”。

[OrmTable("tg_his_variable_wide")]
private sealed class HisVariableWideRow
{
[OrmColumn(IsPrimaryKey = true)]
public string Id { get; set; } = string.Empty;
public string DeviceName { get; set; } = string.Empty;
public string VariableName { get; set; } = string.Empty;
public string Value { get; set; } = string.Empty;
public string? Unit { get; set; }
public bool IsOnline { get; set; }
public DateTime CollectTime { get; set; }
}

public override async Task DBInit(OrmClient db, Logger? logger, CancellationToken cancellationToken)
{
await db.CodeFirst.InitTableAsync<HisVariableWideRow>(cancellationToken).ConfigureAwait(false);
}

public override async Task<int> DBDeletable(OrmClient db, int days, Logger? logger, CancellationToken cancellationToken)
{
var before = DateTime.UtcNow.AddDays(-days);
return await db.Deletable<HisVariableWideRow>()
.Where(row => row.CollectTime < before)
.ExecuteAsync(cancellationToken)
.ConfigureAwait(false);
}

public override async Task DBSavable(OrmClient db, IEnumerable<VariableBasicData> datas, Logger? logger, CancellationToken cancellationToken)
{
var rows = datas.Select(data => new HisVariableWideRow
{
Id = $"{data.Id}:{data.CollectTime:O}",
DeviceName = data.DeviceName,
VariableName = data.Name,
Value = data.Value?.ToString() ?? string.Empty,
Unit = data.Unit,
IsOnline = data.IsOnline,
CollectTime = data.CollectTime
}).ToList();

if (rows.Count == 0)
{
return;
}

await db.BulkCopy<HisVariableWideRow>()
.BulkInsertAsync(rows, cancellationToken)
.ConfigureAwait(false);
}

Demo:实时变量表

类型选择 DynamicSQL_VariableBasicData,可绑定到实时数据转发目标的“实时表脚本”。

[OrmTable("tg_real_variable")]
private sealed class RealVariableRow
{
[OrmColumn(IsPrimaryKey = true)]
public long VariableId { get; set; }
public string DeviceName { get; set; } = string.Empty;
public string VariableName { get; set; } = string.Empty;
public string Value { get; set; } = string.Empty;
public bool IsOnline { get; set; }
public DateTime CollectTime { get; set; }
public DateTime UpdateTime { get; set; }
}

public override async Task DBInit(OrmClient db, Logger? logger, CancellationToken cancellationToken)
{
await db.CodeFirst.InitTableAsync<RealVariableRow>(cancellationToken).ConfigureAwait(false);
}

public override Task<int> DBDeletable(OrmClient db, int days, Logger? logger, CancellationToken cancellationToken)
{
return Task.FromResult(0);
}

public override async Task DBSavable(OrmClient db, IEnumerable<VariableBasicData> datas, Logger? logger, CancellationToken cancellationToken)
{
var rows = datas.Select(data => new RealVariableRow
{
VariableId = data.Id,
DeviceName = data.DeviceName,
VariableName = data.Name,
Value = data.Value?.ToString() ?? string.Empty,
IsOnline = data.IsOnline,
CollectTime = data.CollectTime,
UpdateTime = DateTime.UtcNow
}).ToList();

if (rows.Count > 0)
{
await db.BulkCopy<RealVariableRow>()
.BulkMergeAsync(rows, cancellationToken)
.ConfigureAwait(false);
}
}

Demo:设备状态表

类型选择 DynamicSQL_DeviceBasicData。当前标准目标主要使用变量和报警动态 SQL;设备动态 SQL 是框架支持能力,通常由二开数据转发目标调用。

[OrmTable("tg_device_state")]
private sealed class DeviceStateRow
{
[OrmColumn(IsPrimaryKey = true)]
public long DeviceId { get; set; }
public string DeviceName { get; set; } = string.Empty;
public string Status { get; set; } = string.Empty;
public string PluginName { get; set; } = string.Empty;
public string? LastErrorMessage { get; set; }
public DateTime ActiveTime { get; set; }
public DateTime UpdateTime { get; set; }
}

public override async Task DBInit(OrmClient db, Logger? logger, CancellationToken cancellationToken)
{
await db.CodeFirst.InitTableAsync<DeviceStateRow>(cancellationToken).ConfigureAwait(false);
}

public override Task<int> DBDeletable(OrmClient db, int days, Logger? logger, CancellationToken cancellationToken)
{
return Task.FromResult(0);
}

public override async Task DBSavable(OrmClient db, IEnumerable<DeviceBasicData> datas, Logger? logger, CancellationToken cancellationToken)
{
var rows = datas.Select(data => new DeviceStateRow
{
DeviceId = data.Id,
DeviceName = data.Name,
Status = data.DeviceStatus.ToString(),
PluginName = data.PluginName,
LastErrorMessage = data.LastErrorMessage,
ActiveTime = data.ActiveTime,
UpdateTime = DateTime.UtcNow
}).ToList();

if (rows.Count > 0)
{
await db.BulkCopy<DeviceStateRow>()
.BulkMergeAsync(rows, cancellationToken)
.ConfigureAwait(false);
}
}

Demo:历史报警表

类型选择 DynamicSQL_AlarmVariable,可绑定到历史报警转发目标的“历史报警表脚本”。

[OrmTable("tg_alarm_history_wide")]
private sealed class AlarmHistoryRow
{
[OrmColumn(IsPrimaryKey = true)]
public string AlarmId { get; set; } = string.Empty;
public string DeviceName { get; set; } = string.Empty;
public string VariableName { get; set; } = string.Empty;
public int AlarmLevel { get; set; }
public string EventType { get; set; } = string.Empty;
public string? AlarmText { get; set; }
public DateTime EventTime { get; set; }
}

public override async Task DBInit(OrmClient db, Logger? logger, CancellationToken cancellationToken)
{
await db.CodeFirst.InitTableAsync<AlarmHistoryRow>(cancellationToken).ConfigureAwait(false);
}

public override async Task<int> DBDeletable(OrmClient db, int days, Logger? logger, CancellationToken cancellationToken)
{
var before = DateTime.UtcNow.AddDays(-days);
return await db.Deletable<AlarmHistoryRow>()
.Where(row => row.EventTime < before)
.ExecuteAsync(cancellationToken)
.ConfigureAwait(false);
}

public override async Task DBSavable(OrmClient db, IEnumerable<AlarmVariable> datas, Logger? logger, CancellationToken cancellationToken)
{
var rows = datas.Select(alarm => new AlarmHistoryRow
{
AlarmId = $"{alarm.AlarmId}:{alarm.EventType}:{alarm.EventTime:O}",
DeviceName = alarm.DeviceName,
VariableName = alarm.Name,
AlarmLevel = alarm.AlarmLevel,
EventType = alarm.EventType.ToString(),
AlarmText = alarm.AlarmText,
EventTime = alarm.EventTime
}).ToList();

if (rows.Count > 0)
{
await db.BulkCopy<AlarmHistoryRow>()
.BulkInsertAsync(rows, cancellationToken)
.ConfigureAwait(false);
}
}

Demo:插件事件表

类型选择 DynamicSQL_PluginEventData。当前标准目标未发现直接消费入口,通常由二开目标调用。

[OrmTable("tg_plugin_event")]
private sealed class PluginEventRow
{
[OrmColumn(IsPrimaryKey = true)]
public string Id { get; set; } = string.Empty;
public string DeviceName { get; set; } = string.Empty;
public string Payload { get; set; } = string.Empty;
public DateTime CreateTime { get; set; }
}

public override async Task DBInit(OrmClient db, Logger? logger, CancellationToken cancellationToken)
{
await db.CodeFirst.InitTableAsync<PluginEventRow>(cancellationToken).ConfigureAwait(false);
}

public override async Task<int> DBDeletable(OrmClient db, int days, Logger? logger, CancellationToken cancellationToken)
{
var before = DateTime.UtcNow.AddDays(-days);
return await db.Deletable<PluginEventRow>()
.Where(row => row.CreateTime < before)
.ExecuteAsync(cancellationToken)
.ConfigureAwait(false);
}

public override async Task DBSavable(OrmClient db, IEnumerable<PluginEventData> datas, Logger? logger, CancellationToken cancellationToken)
{
var rows = datas.Select(item => new PluginEventRow
{
Id = Guid.NewGuid().ToString("N"),
DeviceName = item.DeviceName,
Payload = item.ObjectValue.GetRawText(),
CreateTime = DateTime.UtcNow
}).ToList();

if (rows.Count > 0)
{
await db.BulkCopy<PluginEventRow>()
.BulkInsertAsync(rows, cancellationToken)
.ConfigureAwait(false);
}
}

注意事项

动态 SQL 方法签名必须和 DynamicSQLBase<T> 完全一致。旧模板或手写模板如果缺少 OrmClient dbLogger? loggerint days,会编译失败。

不要在 DBSavable 中逐条写数据库。优先批量写入或合并,否则高频历史数据会拖垮数据库。

DBDeletable 必须尊重保留天数。分表、宽表、实时表都要明确是否需要清理。

如果目标使用离线缓存,DBSavable 抛异常会让本批次失败并保留缓存;不要把数据库失败吞掉当成功。

MQTT RPC 脚本

用途

MQTT RPC 脚本用于自定义 MQTT 写入、查询或控制消息的解析和响应。标准 MQTT Client/Server 目标在没有脚本时默认解析:

{
"PLC_1": {
"StartCommand": true,
"SpeedSet": 1200
}
}

脚本可以把项目自定义 payload 转成这个结构,再调用 getRpcResult 执行网关内部写入。

方法形态

MQTT RPC 写的是完整方法,必须实现:

public override Task RPCInvokeAsync(
MqttArrivedMessage message,
Func<TopicArray, CancellationToken, Task> publish,
Func<Dictionary<string, Dictionary<string, JsonElement>>, ValueTask<Dictionary<string, Dictionary<string, IOperResult>>>> getRpcResult,
Logger? logger = null,
CancellationToken cancellationToken = default);

message 是收到的 MQTT 消息;publish 用于回复消息;getRpcResult 用于把“设备名 -> 变量名 -> 值”的写入请求交给网关执行。

Demo:兼容自定义写入格式

输入 payload:

{
"device": "PLC_1",
"values": {
"StartCommand": true,
"SpeedSet": 1200
}
}

脚本内容:

using System.Text;

private sealed class RpcRequest
{
public string Device { get; set; } = string.Empty;
public Dictionary<string, JsonElement> Values { get; set; } = new();
}

public override async Task RPCInvokeAsync(
MqttArrivedMessage message,
Func<TopicArray, CancellationToken, Task> publish,
Func<Dictionary<string, Dictionary<string, JsonElement>>, ValueTask<Dictionary<string, Dictionary<string, IOperResult>>>> getRpcResult,
Logger? logger = null,
CancellationToken cancellationToken = default)
{
var json = Encoding.UTF8.GetString(message.Payload);
var request = json.FromSystemTextJsonString<RpcRequest>();
if (request == null || string.IsNullOrWhiteSpace(request.Device))
{
logger?.LogWarning("Invalid MQTT RPC payload.");
return;
}

var rpcData = new Dictionary<string, Dictionary<string, JsonElement>>
{
[request.Device] = request.Values
};

var result = await getRpcResult(rpcData).ConfigureAwait(false);

using var response = new TopicArray
{
Topic = $"{message.TopicName}/Response",
Payload = new
{
success = result.Values.SelectMany(device => device.Values).All(item => item.IsSuccess),
detail = result
}.ToSystemTextJsonUtf8Bytes()
};

await publish(response, cancellationToken).ConfigureAwait(false);
}

注意事项

必须回复到项目约定的 Topic。默认逻辑使用 ${原Topic}/Response,自定义脚本可以改,但外部系统也要同步修改。

不要直接操作变量运行时绕过 getRpcResult,否则会绕开变量写权限、写表达式、采集插件写入流程和 RPC 日志。

脚本应捕获并记录无法解析的 payload。解析失败直接抛异常时,外部系统通常只能看到超时。

ThingsBoard 客户端插件有自己的 RPC 解析流程,当前源码未发现它使用 BigTextScriptRpc

完整源码脚本

用途

完整源码脚本适合高级扩展:把完整 C# 类型编译进脚本 DLL,运行时通过 GetFullSource() 获取实例。它不等同于数据转换或动态模型,普通变量和转发目标不会自动调用它。

源码形态

FullSource 不做方法包装,页面内容必须是完整 C# 源码。源生成器通过继承 DynamicBase 识别完整源码脚本。

Demo:维护窗口判断工具

using System;
using System.Diagnostics.CodeAnalysis;
using ThingsGatewayRuntime.Application;

#nullable enable

[DynamicallyAccessedMembers(DynamicallyAccessedMemberTypes.All)]
public sealed class MaintenanceWindowHelper : DynamicBase
{
public bool IsInWindow(DateTime utcNow, int startHourUtc, int endHourUtc)
{
var hour = utcNow.Hour;
if (startHourUtc <= endHourUtc)
{
return hour >= startHourUtc && hour < endHourUtc;
}

return hour >= startHourUtc || hour < endHourUtc;
}
}

调用侧示意:

var helper = "MaintenanceWindowHelper".GetFullSource() as MaintenanceWindowHelper;
var allow = helper?.IsInWindow(DateTime.UtcNow, 16, 18) == true;

注意事项

完整源码脚本必须自己保证类名、命名空间、依赖和访问级别。建议类名和脚本名称保持一致。

DynamicBase 当前不定义 Name 抽象属性;不要照搬带 public override string Name 的旧模板,否则会编译失败。

完整源码脚本更适合作为二开扩展,不建议现场交付人员直接维护。

自定义节点脚本

用途

自定义节点用于规则引擎。节点有输入、输出、输入输出参数,运行时把节点连接成有向图:上游输出变化后写入下游输入,再触发下游 ChangedAsync

生命周期

  1. 规则流程启动时读取布局数据。
  2. 每个节点按 NodeTypeNameExpressionsData 创建实例。
  3. 运行时把页面配置的输入参数和输入输出参数写入实例。
  4. 设置 OutputsCommitted,接收节点提交的输出批次;同批状态输入先一起应用,再触发下游执行。
  5. 调用每个节点的 InitAsync
  6. 没有上游的起始节点会执行一次 ChangedAsync
  7. 上游输出变化时,目标节点通过 SmartChangedTriggerScheduler 触发 ChangedAsync。默认 10 ms 防抖,流程配置 NoDebounce 后取消防抖。
  8. 流程停止时释放节点实例;实现了 IDisposable 的节点会被 TryDispose 调用。

页面生成属性的代码形态

如果在节点页面配置输入 Input、输入 Scale、输出 Output,脚本内容只需要写方法:

public override Task InitAsync()
{
Output = 0;
return Task.CompletedTask;
}

public override Task ChangedAsync()
{
Output = Input * Scale;
return Task.CompletedTask;
}

Demo:泵运行小时累计节点

节点参数建议:

参数名方向类型初始值说明
Running输入Booleanfalse泵运行状态。
Reset输入Booleanfalse复位累计值。
Hours输出Double累计运行小时。

脚本内容:

private DateTime _lastTime = DateTime.UtcNow;
private double _hours;

public override Task InitAsync()
{
_lastTime = DateTime.UtcNow;
_hours = 0;
Hours = 0;
return Task.CompletedTask;
}

public override Task ChangedAsync()
{
var now = DateTime.UtcNow;
if (Reset)
{
_hours = 0;
}
else if (Running)
{
_hours += (now - _lastTime).TotalHours;
}

_lastTime = now;
Hours = Math.Round(_hours, 3);
return Task.CompletedTask;
}

Demo:完整自定义节点类

嵌入式节点和外部 DLL 可直接继承 CustomExpressionBase

using System.ComponentModel;
using ThingsGatewayRuntime.Application;

[Category("计算")]
public sealed class HighLimitNode : CustomExpressionBase
{
public override string Name => "高限判断";

[ExpressionInput]
public double Input { get; set; }

[ExpressionInput]
public double Limit { get; set; } = 100;

private bool _alarm;

[ExpressionOutput]
public bool Alarm
{
get => _alarm;
private set => SetOutput(ref _alarm, value);
}

public override Task InitAsync()
{
Alarm = false;
return Task.CompletedTask;
}

public override Task ChangedAsync()
{
Alarm = Input > Limit;
return Task.CompletedTask;
}
}

注意事项

输出属性必须通过 SetOutput 写入。调度执行成功后自动提交输出批次,也可用 CommitOutputs 显式提交中间批次;失败时舍弃未提交部分。不要直接调用宿主的 OutputsCommitted,只改字段也不会触发下游节点。

订阅全局事件、创建定时器、打开网络连接的节点必须实现 IDisposable,在 Dispose 中取消订阅并释放资源。

ChangedAsync 可能被频繁触发。长耗时动作要考虑防抖、超时、取消和重复触发。

规则引擎会检测循环传播,同一传播链重复访问节点会停止传播。不要用节点环路实现高速循环控制。

排障

现象检查项
编译失败先看第一条红色错误;确认脚本类型和代码形态匹配,动态 SQL 和 MQTT RPC 是否写了 override 方法。
编译成功但列表找不到脚本确认 DLL 已热加载;刷新“已加载脚本”;检查脚本名称、类型和分类。
变量表达式不生效确认变量绑定的是脚本名称;读表达式类型必须是 DataTrans,内存变量读表达式必须是 MemoryVariableDatatrans
写入没有到设备检查写表达式是否抛异常;检查采集插件是否支持写入或 RPC;检查变量保护类型和权限。
内存变量不变化检查触发方式、依赖变量名称、Tag 是否执行、依赖变量是否在线、脚本是否返回了可转换类型。
动态模型 Topic 报属性不存在Topic 模板 ${字段} 与脚本输出对象属性名不一致。
动态 SQL 未执行确认目标属性中绑定了对应表脚本;历史变量/实时变量使用变量动态 SQL,历史报警使用报警动态 SQL。
MQTT RPC 无响应检查 RPC Topic 是否匹配;payload 是否能解析;脚本是否调用 publish;外部系统是否监听正确响应 Topic。
自定义节点不触发检查节点输出是否通过 SetOutput 设置;连线端口是否对应参数;流程是否启用;防抖是否影响观察。
AOT 环境脚本不可用动态代码不支持时不会热加载脚本 DLL,需要改为非 AOT 部署或在发布流程中预编译并验证。

上线检查清单

检查项要求
类型匹配脚本类型、绑定位置、输入实体类型一致。
编译单个编译成功,批量编译成功,重启后已加载脚本仍存在。
参数输入参数名称、类型、初始值和变量配置中的参数值一致。
性能高频脚本不访问慢速外部资源,不逐条写数据库,不正常路径刷 Info 日志。
异常解析失败、变量不存在、数据库失败、RPC 失败都有明确日志。
数据先用 3 到 5 个点验证输入、输出、时间、单位和在线状态,再扩大范围。
回退保留旧脚本内容或导出配置,生产修改前确认可回退。

相关文档