- 物联网
- 后端
- 数据可视化
- 消息队列
【免费下载链接】thingsboard
All-in-one IoT Platform - Device management, data collection, processing and visualization.
本文围绕 ThingsBoard 计算字段(Calculated Field)中 TBEL 脚本的时间序列滚动参数(Time Series Rolling Argument)展开,完整解析merge()与mergeAll()两个内置合并函数的输入数据格式、用法、输出结构、时间戳对齐与 NaN 处理规则,并结合仓库源码(TbelCfTsRollingArg)与单元测试揭示底层实现原理。读者读完将能够独立编写多遥测序列对齐、交叉分析的 TBEL 计算字段脚本,例如冰箱除霜状态与温度联合判定、多传感器数据融合等典型场景。
一、背景:计算字段中的三类参数
在 ThingsBoard 计算字段中,用户通过calculate(ctx, arg1, arg2, ...)函数编写 TBEL(ThingsBoard Expression Language)脚本,对遥测与属性数据进行自定义计算。函数接收计算字段配置中按名称传入的参数,分为三种类型:
- 属性与最新遥测参数:单值,类型可为 boolean、int64(long)、double、string 或 JSON;
- 时间序列滚动参数(Time Series Rolling Argument):时间窗口内的时序数据集合;
ctx上下文对象:携带latestTs(参数遥测的最新时间戳,毫秒)以及ctx.args.<argName>形式的参数访问入口。
本文聚焦第 2 类参数——时间序列滚动参数之间的**合并(merge)**操作。相关功能定义可参考 expression_fn.md。
二、输入数据结构:合并操作的前提
merge()与mergeAll()的输入是时间序列滚动参数,其标准结构包含timeWindow(时间窗口起止)与values(采样点数组,每个点包含ts与value)。以下为合并示例的完整输入数据(来自 merge_input.md):
{ "humidity": { "timeWindow": { "startTs": 1741356332086, "endTs": 1741357232086 }, "values": [{ "ts": 1741356882759, "value": 43 }, { "ts": 1741356918779, "value": 46 }] }, "pressure": { "timeWindow": { "startTs": 1741356332086, "endTs": 1741357232086 }, "values": [{ "ts": 1741357047945, "value": 1023 }, { "ts": 1741357056144, "value": 1026 }, { "ts": 1741357147391, "value": 1025 }] }, "temperature": { "timeWindow": { "startTs": 1741356332086, "endTs": 1741357232086 }, "values": [{ "ts": 1741356874943, "value": 76 }, { "ts": 1741357063689, "value": 77 }] } }三个序列共享同一个timeWindow(startTs=1741356332086,endTs=1741357232086,窗口约 15 分钟),但各自采样时刻与采样密度均不同:
humidity:2 个采样点(43、46),时间靠前;pressure:3 个采样点(1023、1026、1025),时间居中;temperature:2 个采样点(76、77),时间跨度最大。
这正是合并函数存在的意义:将采样时间不同步的多条时序数据对齐到同一组时间戳上,从而进行逐点联合分析。
注意:根据文档说明,滚动参数中的数值统一转换为
double类型,当转换失败时以NaN表示,可用isNaN(double): boolean判断有效性。这意味着输入数据中天然可能存在非数值点,合并操作必须正确处理它们。
三、merge():两条时间序列的两两合并
3.1 用法
merge(other, settings)将当前滚动参数与另一个滚动参数合并,对齐时间戳,并以前一个可用值填充缺失值。用法示例(来自 merge_usage.md):
var mergedData = temperature.merge(humidity, { ignoreNaN: false });这里temperature为调用者(结果中的第一列),humidity为被合并方(结果中的第二列),settings传入{ ignoreNaN: false }。
3.2 输出结果
执行上述合并后得到(来自 merge_output.md):
{ "mergedData": { "timeWindow": { "startTs": 1741356332086, "endTs": 1741357232086 }, "values": [{ "ts": 1741356874943, "values": [76.0, "NaN"] }, { "ts": 1741356882759, "values": [76.0, 43.0] }, { "ts": 1741356918779, "values": [76.0, 46.0] }, { "ts": 1741357063689, "values": [77.0, 46.0] }] } }3.3 逐点解读
合并结果逐点分析,可以清晰看出对齐与填充规则:
| 时间戳 | 第一列(temperature) | 第二列(humidity) | 说明 |
|---|---|---|---|
| 1741356874943 | 76.0 | NaN | temperature 自己的采样点;humidity 在此时刻尚无数据,且ignoreNaN=false,故保留该行为 NaN |
| 1741356882759 | 76.0 | 43.0 | humidity 新增采样点;temperature 无新数据,前向填充上一值 76.0 |
| 1741356918779 | 76.0 | 46.0 | humidity 新增采样点 46.0;temperature 继续前向填充 76.0 |
| 1741357063689 | 77.0 | 46.0 | temperature 新增采样点 77.0;humidity 无新数据,前向填充 46.0 |
关键规则:
- 时间戳并集:输出时间戳为所有参与序列采样时刻的并集(升序);
- 前向填充(carry-forward):某个序列在当前时间戳没有新采样时,沿用其最近一次(更早)的值;
ignoreNaN控制行保留:false时,含NaN的行(如首行第二列)也会被保留,便于观察"缺失"状态。
四、mergeAll():多条时间序列的批量合并
4.1 用法
mergeAll(others, settings)将当前滚动参数与一个滚动参数数组批量合并。用法示例(来自 merge_all_usage.md):
var mergedData = temperature.mergeAll([humidity, pressure], { ignoreNaN: true });调用者temperature为第一列,数组[humidity, pressure]依次成为第二、三列。
4.2 输出结果
{ "mergedData": { "timeWindow": { "startTs": 1741356332086, "endTs": 1741357232086 }, "values": [{ "ts": 1741357047945, "values": [76.0, 46.0, 1023.0] }, { "ts": 1741357056144, "values": [76.0, 46.0, 1026.0] }, { "ts": 1741357063689, "values": [77.0, 46.0, 1026.0] }, { "ts": 1741357147391, "values": [77.0, 46.0, 1025.0] }] } }4.3 与 merge() 的差异
对比可见两个重要差异:
ignoreNaN: true过滤行:所有 4 行均为三列有效数值(无 NaN 行)。第一行ts=1741357047945来自pressure的第一个采样点,此时 temperature 前向填充 76.0、humidity 前向填充 46.0。因为ignoreNaN=true,任何一列缺失都会导致整行被跳过,因此早期那些只有 temperature/humidity 数据而 pressure 尚无数据的时刻不会出现在结果中;- 列数扩展:输出
values数组元素个数 = 参与合并的序列总数(1 + others 数量),顺序与传入顺序一致。
这个输出也印证了 mergeAll 的输出行数(4 行)恰好是各序列采样时刻中"所有列都已具备前值"的那些时间戳。
五、settings 参数详解
两个函数都支持可选的settings配置对象(文档定义见 expression_fn.md):
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
ignoreNaN | boolean | true | 控制合并结果中是否跳过含NaN的行。true:任一行任一列存在 NaN 即跳过该行(本示例中 pressure 未就绪的行被跳过);false:保留含 NaN 的行,供观察缺失状态 |
timeWindow | object | 各序列 timeWindow 的并集 | 自定义输出时间窗口,仅保留落在该窗口内的对齐点;支持{startTs, endTs}形式的对象 |
timeWindow自定义示例:
var mergedData = temperature.merge(humidity, { ignoreNaN: true, timeWindow: { startTs: 1741356300000, endTs: 1741357200000 } });从源码实现看,timeWindow设置项同时支持TbTimeWindow对象、{startTs, endTs}Map,以及 JSON 字符串三种传入形态。
六、源码级实现原理
merge()与mergeAll()的底层实现位于 TbelCfTsRollingArg.java。核心设计如下:
6.1 方法委托
merge(other, settings)内部直接委托给mergeAll:
public TbelCfTsRollingData merge(TbelCfTsRollingArg other, Map<String, Object> settings) { return mergeAll(Collections.singletonList(other), settings); }mergeAll(others, settings)则将调用者自身置于序列首位,再追加others:
List<TbelCfTsRollingArg> args = new ArrayList<>(others.size() + 1); args.add(this); args.addAll(others);这解释了输出列顺序:"调用者在前、其余按数组顺序在后"。
6.2 对齐算法
mergeAll的核心算法(源码第 270–337 行)分为四步:
- 解析 settings:读取
ignoreNaN(默认true)与可选timeWindow; - 收集时间戳并集:将所有序列的采样
ts放入TreeSet<Long>(天然升序、去重),并计算各序列timeWindow的最小startTs与最大endTs作为输出窗口; - 前向填充扫描:使用
lastIndex[]游标数组逐序列推进——对每个输出时间戳ts,只要values[lastIndex[i]].getTs() <= ts就持续更新result[i],从而把"最近一个不晚于当前时刻的值"带入该列; - 按规则产出行:若未自定义
timeWindow或当前ts落在窗口内,则依据ignoreNaN决定是否跳过含NaN的行(true时任一列为 NaN 即跳过,false时无条件产出),最终封装为TbelCfTsRollingData。
while (lastIndex[i] < values.size() && values.get(lastIndex[i]).getTs() <= ts) { result[i] = values.get(lastIndex[i]).getValue(); lastIndex[i]++; }这段代码正是"前向填充"的精确实现:对每个时间戳按列推进游标,<=条件保证了取到的是最接近且不晚于当前时刻的值。
6.3 返回值结构
返回类型为TbelCfTsRollingData,它实现了Iterable,因此在 TBEL 中可以直接用foreach(item: merged) { ... }遍历,每个item暴露ts与values(数组,可通过item.v1、item.v2等下标语义访问各列)。
6.4 单元测试佐证
TbelCfTsRollingArgTest.java中针对 merge 覆盖了四个典型场景,与文档行为完全一致:
merge_two_rolling_args_ts_match_test:两个序列时间戳完全一致时,逐行精确对齐;merge_two_rolling_args_with_timewindow_test:传入自定义timeWindow(对象与 Map 两种形式)后仅保留窗口内(本例为前 10000ms)的 2 行;merge_two_rolling_args_ts_mismatch_default_test:时间戳不一致且默认ignoreNaN=true时,只输出 3 行(全部列具备前值的时刻,ts=200处第一列前向填充为 1);merge_two_rolling_args_ts_mismatch_ignore_nan_disabled_test:ignoreNaN=false时输出 4 行,其中ts=100处第二列保持NaN。
var result = arg1.merge(arg2, Collections.singletonMap("ignoreNaN", false)); Assertions.assertEquals(4, result.getSize()); Assertions.assertEquals(Double.NaN, item0.getValues()[1]); // 缺失位置保留 NaN测试结果从工程上验证了本文第三、四节对输出结构的全部解读。
七、实战:冰箱冷冻室温度与除霜状态联合分析
合并函数的典型价值在于"跨序列条件判定"。以下是文档中给出的完整实战示例(来自 expression_fn.md):将temperature(冷冻室温度)与defrost(除霜状态,0/1)合并后,找出未处于除霜模式但温度高于 -5°C的异常时刻:
function calculate(ctx, temperature, defrost) { var merged = temperature.merge(defrost); var result = []; foreach(item: merged) { if (item.v1 > -5.0 && item.v2 == 0) { result.add({ ts: item.ts, values: { issue: { temperature: item.v1, defrostState: false } } }); } } return result; }其中item.v1为合并结果第一列(temperature),item.v2为第二列(defrost)。返回结果为符合条件的时间点列表,可直接用于配置告警规则:
[ { "ts": 1741613833843, "values": { "issue": { "temperature": -3.12, "defrostState": false } } }, { "ts": 1741613923848, "values": { "issue": { "temperature": -4.16, "defrostState": false } } } ]注意:此处merge未显式传入settings,即采用默认ignoreNaN=true——任一行存在 NaN 就会被跳过,从而保证判定条件item.v1 > -5.0始终基于有效数值。
八、常见误区与最佳实践
- 列顺序 = 调用顺序:
a.merge(b)中a永远是第一列,b是第二列;mergeAll中数组顺序决定后续列。条件判断(如item.v1、item.v2)必须与之一一对应; ignoreNaN默认行为:默认true会跳过含 NaN 的行,适合"只关心完整数据点"的统计与告警场景;如需观察缺失状态或显式区分"无数据"与"0 值",应显式传false;- 时间窗口限制:结果窗口默认取各序列
timeWindow的并集;如希望裁剪输出范围,用settings.timeWindow自定义窗口,源码中timeWindow.matches(ts)会在窗口外直接跳过该行; - 空序列保护:滚动参数的内置统计方法(
sum()、max()、mean()等)在values为空时会抛出IllegalArgumentException,合并前应确认各序列至少包含一个有效采样点; - 返回值格式:合并对象作为计算字段的返回值,输出类型默认
Time Series,可返回带ts的对象/数组(推荐使用ctx.latestTs对齐触发遥测时间戳),也可直接返回不带时间戳的 JSON 对象(此时用于 Attribute 输出)。
九、总结
merge()与mergeAll()是 ThingsBoard 计算字段 TBEL 脚本中处理多路时间序列的核心工具,通过"时间戳并集 + 前向填充 + NaN 策略"三个机制,将采样不同步的遥测数据对齐到统一时间轴,为逐点联合分析、异常判定与告警生成提供数据基础。本文从输入结构、用法、输出解读到TbelCfTsRollingArg源码算法与单元测试,完整还原了该功能的运行机制。相关配套文档与源码路径如下:
- 输入示例:merge_input.md
- merge 用法/输出:merge_usage.md、merge_output.md
- mergeAll 用法/输出:merge_all_usage.md、merge_all_output.md
- TBEL 计算字段总览:expression_fn.md
- 核心实现:TbelCfTsRollingArg.java
- 单元测试:TbelCfTsRollingArgTest.java
- 物联网
- 后端
- 数据可视化
- 消息队列
【免费下载链接】thingsboard
All-in-one IoT Platform - Device management, data collection, processing and visualization.
相关推荐
ThingsBoard 计算字段 TBEL 中 merge / mergeAll 时间序列合并函数完整指南
ThingsBoard 计算字段 TBEL 中 merge / mergeAll 时间序列合并函数完整指南 导读 ThingsBoard 的计算字段(Calcu
物联网后端数据可视化消息队列PyPTO 实现 RoPE(旋转位置编码)算子:从 kernel 参考骨架到生产级实现
PyPTO 实现 RoPE(旋转位置编码)算子:从 kernel 参考骨架到生产级实现 导读 本文以 PyPTO Gym 仓库中 RoPE kernel 参考骨
物联网后端数据可视化消息队列Context7 MCP Server 完整配置实战:把最新库文档注入任意 AI 提示词的部署与排错指南
Context7 MCP Server 完整配置实战:把最新库文档注入任意 AI 提示词的部署与排错指南 本文以 Klavis 仓库 mcp_servers/c
物联网后端数据可视化消息队列
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考