跳到主要内容

脚本开发说明

本文面向需要扩展 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 在保留天数清理时执行。

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. 设置 OutputChangedCallback,输出变化时传播到下游节点。
  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 或等价的 OutputChangedCallback 调用传播变化。只改字段不会触发下游节点。

订阅全局事件、创建定时器、打开网络连接的节点必须实现 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 个点验证输入、输出、时间、单位和在线状态,再扩大范围。
回退保留旧脚本内容或导出配置,生产修改前确认可回退。

相关文档