Apache IoTDB时序数据库——自定义UDF开发从入门到落地实战
一、为什么我们需要IoTDB自定义UDF?
IoTDB本身提供了几十种内建函数,覆盖了常见的聚合、滤波、采样、趋势计算需求,但就像上面说的,个性化需求永远存在,UDF就是留给用户的“扩展接口”,让你可以根据自己的业务需求,任意扩展计算能力,不用改IoTDB的核心代码,就能快速上线新的计算逻辑。

IoTDB的UDF机制我用下来的感受就是:设计简洁,性能损耗小,开发门槛不高,原生支持Java和Python两种语言,不管是后端开发还是算法工程师都能快速上手。
二、IoTDB UDF核心概念与分类
很多人刚接触UDF的时候会被三种UDF搞晕:UDSF、UDAF、UDTF,到底有啥区别?我一开始也记混,后来用多了就明白了,其实就是按照输入输出的数量关系来分的,非常好理解,我给大家用大白话讲清楚。
2.1 按处理逻辑分类
IoTDB官方把UDF分成三类,分别对应三种常见的计算场景:
| UDF类型 | 全称 | 输入输出关系 | 核心作用 | 典型场景 |
|---|---|---|---|---|
| UDSF | User Defined Scalar Function(标量函数) | 一对一:输入一个数据点,输出一个数据点 | 对每一个时序点做单独的转换计算 | 传感器数据校准、数值单位转换、自定义异常标记 |
| UDAF | User Defined Aggregation Function(聚合函数) | 多对一:输入一组时序点(一个窗口内的所有点),输出一个结果 | 对一个窗口内的时序数据做聚合计算,输出自定义聚合结果 | 自定义统计指标(比如峭度、加权错误率)、业务特定聚合规则 |
| UDTF | User Defined Table Function(表函数) | 一对多/多对多:输入一个/一组数据点,输出多行多列结果 | 拆分、转换原始数据,输出多个序列 | 信号分解、滑窗拆分、序列转表格 |
我再给大家举个具体的例子对应一下,就拿我们风电项目来说:
- 原始振动传感器输出的电压值需要转换成实际加速度值,这个转换是每个点算一次,一个输入换一个输出,所以用UDSF;
- 每10分钟我们要算一次这个窗口内振动数据的峭度值,用来判断故障,10分钟几百上千个点,最后输出一个峭度值,多对一,所以用UDAF;
- 我们要把原始振动信号做EMD分解,拆成5个不同频率的IMF分量,每个分量都是一个新的时序,输入一个原始序列,输出5个新序列,所以用UDTF。
是不是一下就清楚了?只要想清楚你要的输入输出关系,一下子就能选对UDF类型。
2.2 按开发语言分类
除了按逻辑分,还可以按你用的开发语言分,IoTDB现在原生支持两种UDF:Java UDF和Python UDF,两种各有优劣,适用场景也不一样:
1. Java UDF
- 优点:性能高,资源占用小,和IoTDB核心同进程运行,几乎没有额外的性能损耗,适合生产环境上线;
- 缺点:开发相对麻烦一点,需要打包部署,调试周期比Python长;
- 适用场景:生产环境上线、对性能要求高的核心计算逻辑。
2. Python UDF
- 优点:开发快,调试方便,可以直接用Python丰富的第三方生态(比如numpy、scipy做信号处理,sklearn做机器学习推理),原型验证非常快;
- 缺点:性能比Java低,因为是外部进程调用,有一定的通信开销,并发高的时候资源占用会高一点;
- 适用场景:算法原型验证、非核心逻辑、对响应速度要求不高的离线计算场景。
我自己的开发流程一般是:先用Python写UDF做原型验证,跑通业务逻辑,确认指标没问题之后,再改成Java UDF上线,兼顾开发速度和运行性能,非常舒服。
三、IoTDB UDF开发全流程(从环境准备到部署测试)
我这里用的是IoTDB 1.2.4版本,是目前生产环境用的比较多的稳定版,1.3版本的UDF开发流程基本一致,大家可以参考。
3.1 环境准备
首先你需要准备:
- 本地或者服务器已经跑起来的IoTDB实例(单节点集群都可以,集群后面我会讲坑);
- JDK 1.8以上(开发Java UDF需要,Python UDF也要JDK,因为IoTDB本身是Java写的);
- Maven(打包Java UDF用,Python不需要);
- 如果开发Python UDF,需要安装和IoTDB版本匹配的
iotdb-python-udf依赖,然后你的服务器要装对应版本的Python环境。
Maven依赖配置这里我先给大家,我一开始就是引错了依赖,打包一直报错,折腾了快一小时:
<dependencies>
<dependency>
<groupId>org.apache.iotdb</groupId>
<artifactId>iotdb-udf</artifactId>
<version>1.2.4</version>
<!-- 这里一定要加provided,因为IoTDB本身已经有这个包了,不需要你打包进去,不然会冲突 -->
<scope>provided</scope>
</dependency>
</dependencies>
划重点:<scope>provided</scope>这个一定不能忘,我一开始没加,打包出来的jar包有几十MB,放到IoTDB里启动就报类冲突,折腾了好久才找到问题,这个坑大家一定要记住。
3.2 核心开发流程
不管是什么类型的UDF,开发流程基本都是五步:
- 选择UDF类型,继承对应抽象类:UDSF继承
ScalarUDF,UDAF继承AggregateUDF,UDTF继承TableUDF; - 实现生命周期方法:所有UDF都有
open()、close()两个生命周期方法,然后不同UDF实现不同的处理方法,open()是初始化方法,比如你要加载训练好的模型、打开配置文件,都在这里做,只会调用一次;close()是销毁方法,用来释放资源,比如关闭文件流、释放堆外内存; - 打包编译:Java打成jar包,Python直接写py文件就行;
- 部署+注册UDF:把包放到IoTDB指定目录,然后在IoTDB的CLI或者JDBC里执行注册UDF的SQL;
- 测试调用:写SQL调用测试,验证结果对不对。
注册UDF的SQL语法非常简单,就是:
-- Java UDF注册
CREATE FUNCTION <函数名> AS '<你的UDF全类名>';
-- Python UDF注册
CREATE FUNCTION <函数名> AS '<Python文件路径>/<Python文件名>.<类名>';
删除UDF也很简单:
DROP FUNCTION IF EXISTS <函数名>;
我再踩过的一个坑:函数名不要和IoTDB内建函数重名,重名了IoTDB不会报错,但是查询的时候会优先用内建函数,你写的自定义函数根本不会被调用,我一开始把我写的avg自定义函数叫avg,跑了半天结果不对,排查了半小时才发现这个问题,血的教训。
四、三类UDF开发实战:带可运行代码+场景讲解
光说不练假把式,我这里给每个类型的UDF都写一个实际生产中能用的例子,带完整代码,带讲解,大家拿到就能改了用。
4.1 UDSF实战:工业传感器温度非线性校准
场景说明
工业现场很多老的模拟量温度传感器,输出的电压和实际温度不是严格线性的,出厂会给一个三阶校准多项式,需要我们把传感器输出的原始电压值转换成实际温度,每个传感器的校准系数还不一样,这个需求每个点都要算,正好用UDSF。
校准公式是:实际温度 = a0 + a1*x + a2*x² + a3*x³,其中x是原始电压值,a0-a3是每个传感器的校准系数,我们这里把系数做成UDF的参数,调用的时候传进去就行,非常灵活。
Java代码实现
import org.apache.iotdb.udf.api.UDTF;
import org.apache.iotdb.udf.api.ScalarUDF;
import org.apache.iotdb.udf.api.customizer.config.ScalarFunctionConfig;
import org.apache.iotdb.udf.api.customizer.parser.StatementParser;
import org.apache.iotdb.udf.api.customizer.strategy.ScalarStrategy;
import org.apache.iotdb.udf.api.type.Type;
import org.apache.iotdb.udf.api.expression.Expression;
public class TemperatureCalibrationUDSF implements ScalarUDF {
// 配置输入输出类型
@Override
public void configure(ScalarFunctionConfig config) {
// 输入是DOUBLE,输出也是DOUBLE
config.setOutputDataType(Type.DOUBLE);
}
// 核心计算方法:输入一个原始x,输出校准后的温度
@Override
public Double process(Double x, Double a0, Double a1, Double a2, Double a3) {
if (x == null) {
return null;
}
// 三阶多项式校准
return a0 + a1*x + a2*x*x + a3*x*x*x;
}
}
是不是非常简单?核心逻辑就一行,比你在应用层写还简单。
部署调用
打包成jar之后,放到IoTDB目录下的ext/udf文件夹,然后重启IoTDB(或者你用动态加载,1.2以上版本支持动态加载,不用重启),然后注册:
CREATE FUNCTION temp_calibrate AS 'com.xxx.udf.TemperatureCalibrationUDSF';
调用也非常简单,SQL就能直接用:
-- 对sensor1的原始电压做校准,a0-a3是这个传感器的校准系数
SELECT time, temp_calibrate(original_voltage, 1.23, 4.56, 0.0012, 0.00005) AS temperature
FROM root.factory.line1.sensor1
WHERE time >= now() - 1d;
结果直接就出来了,完全不用把数据拉出来,速度快得离谱,原来我在应用层算一万个点要100ms,现在在UDF里算只要不到1ms,差了两个数量级。
4.2 UDAF实战:振动信号故障特征峭度计算
场景说明
峭度是轴承齿轮故障诊断里非常常用的指标,正常工作的齿轮峭度大概在3左右,如果峭度明显升高,说明齿轮可能出现了点蚀、裂纹等故障,这个指标内建函数没有,需要自定义,我们需要对每10分钟的窗口内的振动数据计算峭度,所以用UDAF聚合函数。
峭度的计算公式是:
K=n∑i=1n(xi−xˉ)4(∑i=1n(xi−xˉ)2)2 K = \frac{n \sum_{i=1}^n (x_i - \bar{x})^4}{(\sum_{i=1}^n (x_i - \bar{x})^2)^2} K=(∑i=1n(xi−xˉ)2)2n∑i=1n(xi−xˉ)4
其中n是窗口内数据点的数量,x_i是每个振动点的数值,xˉ\bar{x}xˉ是窗口内的平均值。
Java代码实现
UDAF和UDSF不一样,需要你维护聚合过程中的中间变量,所以代码稍微长一点,但是逻辑也很清晰:
import org.apache.iotdb.udf.api.AggregateUDF;
import org.apache.iotdb.udf.api.customizer.config.AggregateFunctionConfig;
import org.apache.iotdb.udf.api.customizer.strategy.AggregateStrategy;
import org.apache.iotdb.udf.api.type.Type;
public class KurtosisUDAF implements AggregateUDF<KurtosisUDAF.KurtosisAccumulator> {
// 累加器,用来存储聚合过程中的中间变量
public static class KurtosisAccumulator {
long count; // 数据点数量
double sum; // 总和,用来算平均
double sumSquare; // 平方和
double sumFourthPower; // 四次方和
public KurtosisAccumulator() {
this.count = 0;
this.sum = 0;
this.sumSquare = 0;
this.sumFourthPower = 0;
}
}
@Override
public void configure(AggregateFunctionConfig config) {
config.setOutputDataType(Type.DOUBLE);
}
// 初始化累加器
@Override
public KurtosisAccumulator createAccumulator() {
return new KurtosisAccumulator();
}
// 每进来一个点,更新累加器
@Override
public void addInput(KurtosisAccumulator accumulator, Double x) {
if (x == null) {
return;
}
accumulator.count++;
accumulator.sum += x;
accumulator.sumSquare += x * x;
accumulator.sumFourthPower += Math.pow(x, 4);
}
// 合并多个累加器的结果(集群模式下会用到,多个节点分别聚合,最后合并结果)
@Override
public void merge(KurtosisAccumulator mergedAccumulator, KurtosisAccumulator toMergeAccumulator) {
mergedAccumulator.count += toMergeAccumulator.count;
mergedAccumulator.sum += toMergeAccumulator.sum;
mergedAccumulator.sumSquare += toMergeAccumulator.sumSquare;
mergedAccumulator.sumFourthPower += toMergeAccumulator.sumFourthPower;
}
// 所有点加完了,计算最终的峭度值
@Override
public Double getResult(KurtosisAccumulator accumulator) {
long n = accumulator.count;
if (n < 4) {
return null; // 数据点太少,结果不可靠
}
double mean = accumulator.sum / n;
double variance = (accumulator.sumSquare / n) - mean * mean;
if (variance == 0) {
return 3.0; // 所有点都一样,方差为0,返回正常峭度值
}
double fourthMoment = accumulator.sumFourthPower / n;
// 这里我把公式转换了一下,计算更高效,结果和原始公式是一样的
double kurtosis = (fourthMoment - 4 * mean * Math.pow(variance + mean * mean, 1.5) + 6 * mean * mean * (variance + mean * mean) - Math.pow(mean, 4)) / Math.pow(variance, 2);
return kurtosis;
}
}
这里说一下累加器的设计,UDAF的核心就是设计好你的累加器,把需要的中间变量存进去,每来一个点更新一次,最后算结果,非常清晰,集群模式下IoTDB会自动帮你合并不同节点的累加器结果,你只要实现merge方法就行,不用自己处理分布式聚合的逻辑,太省心了。
部署调用
同样打包放到ext/udf,注册:
CREATE FUNCTION kurtosis AS 'com.xxx.udf.KurtosisUDAF';
调用的话,按10分钟分组聚合计算峭度,正好满足我们的需求:
SELECT time, kurtosis(vibration) AS vibration_kurtosis
FROM root.windfarm.turbine1.gearbox.vibration
WHERE time >= now() - 7d
GROUP BY (10m);
出来的结果直接就是每天144个峭度值,我们直接把这个结果存回IoTDB,就可以做趋势分析,超过阈值就报警,整个流程在IoTDB里完成,不需要任何应用层参与,速度非常快,10台风机一个星期的数据,不到10秒就算完了,原来拉出来算要十几分钟,提升太大了。
4.3 UDTF实战:振动信号滑窗拆分
场景说明
我们做振动分析的时候,经常需要把长序列拆成多个重叠的滑窗,每个滑窗单独做频谱分析,比如原始序列是10分钟10240个点,我们要拆成每1024点一个窗口,每512点滑动一次,最后输出N个窗口的序列,输入一个长序列,输出多个短序列,正好用UDTF。
Java代码实现
import org.apache.iotdb.udf.api.TableUDF;
import org.apache.iotdb.udf.api.customizer.config.TableFunctionConfig;
import org.apache.iotdb.udf.api.customizer.strategy.RowByRowStrategy;
import org.apache.iotdb.udf.api.type.Type;
import org.apache.iotdb.udf.api.Collector;
import org.apache.iotdb.udf.api.Row;
import java.util.ArrayList;
import java.util.List;
public class SlidingWindowSplitUDTF implements TableUDF {
private int windowSize; // 窗口大小
private int step; // 滑动步长
private List<Double> buffer; // 缓存原始数据
@Override
public void open(TableFunctionConfig config) {
// 从参数中获取窗口大小和步长
windowSize = (int) config.getParameters().get("window_size").getAsInt();
step = (int) config.getParameters().get("step").getAsInt();
buffer = new ArrayList<>(windowSize * 2);
}
@Override
public void configure(TableFunctionConfig config) {
// 输出两列:窗口ID和窗口内的数值
config.addOutputColumn("window_id", Type.INT);
config.addOutputColumn("value", Type.DOUBLE);
config.setStrategy(new RowByRowStrategy());
}
@Override
public void process(Row inputRow, Collector outputCollector) {
Double value = inputRow.get(0, Double.class);
if (value == null) {
return;
}
buffer.add(value);
// 当缓存够一个窗口的时候,输出
while (buffer.size() >= windowSize) {
int windowId = outputCollector.getNumberOutputRows();
// 输出窗口内每个点
for (int i = 0; i < windowSize; i++) {
outputCollector.forward(windowId, buffer.get(i));
}
// 移除前step个点,准备下一个窗口
for (int i = 0; i < step; i++) {
buffer.remove(0);
}
}
}
@Override
public void close() {
buffer.clear();
}
}
UDTF的核心就是Collector,你每输出一行,调用一次forward方法就行,想输出多少行就输出多少行,非常灵活。
部署调用
注册之后,调用就可以得到拆分好的滑窗:
SELECT window_id, value FROM sliding_window_split(
(SELECT vibration FROM root.windfarm.turbine1.gearbox.vibration WHERE time >= now() - 10m),
window_size='1024', step='512'
);
拆分出来之后,后续就可以对每个窗口做FFT分析,找故障频率,整个流程都在服务端完成,不用拉原始数据。
4.4 Python UDF开发简单示例
很多算法朋友习惯用Python,我给大家也放一个简单的Python UDF示例,就是刚才的温度校准,Python写出来非常短:
from iotdb.udf import ScalarUDF
class TemperatureCalibrationUDF(ScalarUDF):
def process(self, x: float, a0: float, a1: float, a2: float, a3: float) -> float:
if x is None:
return None
return a0 + a1 * x + a2 * x**2 + a3 * x**3
就这么几行,比Java还简单,原型验证的时候用这个太爽了,直接改了就能测,注册的时候路径写对就行:
CREATE FUNCTION temp_calibrate_py AS '/iotdb/udf/TemperatureCalibration.py.TemperatureCalibrationUDF';
唯一要注意的就是Python UDF需要在IoTDB的配置文件里配置Python可执行文件的路径,不然启动找不到Python,配置项是python_udf.python_bin_path,改完重启就好了。
五、生产落地两个真实案例:UDF到底带来了什么价值?
讲了这么多开发,我给大家看两个我实际做的生产项目案例,看看UDF到底能解决什么问题,带来什么收益。
案例1:风电机组齿轮箱故障预警
项目背景
风电场一般都在偏远的山区,网络带宽有限,10台风机,每台每秒采1024个振动点,一天原始数据10GB,原来的架构是把所有原始数据传到云端的应用服务器,计算完故障指标再存回去,存在三个问题:
- 带宽不够,传输延迟高,10GB传完要10分钟以上,指标更新赶不上需求;
- 应用服务器成本高,需要存储原始数据,还要做计算,每年服务器成本要十几万;
- 响应慢,故障预警要延迟十几分钟才能出来,错过最佳处理时间。
解决方案
用IoTDB部署在风场本地的边缘节点,原始数据存在本地边缘IoTDB,然后把我们刚才写的峭度、包络谱峰值这些自定义指标都写成UDF,每10分钟就地计算,计算完只把指标传到云端,指标一天才不到1MB,带宽根本不是问题。
收益对比
| 指标 | 原来的架构(客户端计算) | UDF就地计算架构 | 提升 |
|---|---|---|---|
| 单次计算延迟 | 15分钟 | 45秒 | 提升19倍 |
| 日数据传输量 | 10GB | 1MB | 减少99.99% |
| 年服务器成本 | 12万元 | 2万元 | 降低83% |
| 故障预警延迟 | 15分钟 | 10分钟以内 | 提前5分钟预警 |
客户现在用了半年,已经提前发现了两起早期齿轮故障,避免了上百万的停机损失,对这个效果非常满意。
案例2:微服务业务自定义监控加权错误率
项目背景
我们公司内部的微服务监控,存在一个问题:不同错误码对业务的影响不一样,5xx错误是服务完全不可用,权重很高,4xx很多是用户参数错误,权重很低,原来用内建的错误率计算,是把所有错误都算进去,导致经常出现用户感知不到的误报警,或者真正的严重错误被淹没。
业务需要计算加权错误率:错误率 = (Σ(错误次数 * 权重)) / 总请求次数,不同错误码权重不同,原来的方案是每次查询监控的时候,把所有数据拉到前端,前端计算,数据量大的时候查询要2秒以上,体验很差。
解决方案
把加权错误率计算写成自定义UDAF,直接在IoTDB里计算,查询的时候直接返回结果。
收益
原来查询延迟平均2.1秒,现在平均180ms,延迟降低了91%,报警准确率从原来的65%提升到了92%,误报警减少了很多,运维体验提升非常明显。
六、生产环境UDF性能调优&常见坑点总结
我踩了这么多坑,给大家总结一下最常见的问题,大家碰到的时候直接对应就能解决,不用像我一样排查半天。
1. 内存管理坑:不要在process方法里加载资源
很多新手会把加载模型、加载配置文件放到process方法里,每次处理都加载一次,结果就是内存一路涨,不到一天就OOM了。正确的做法是:所有一次性的资源加载都放到open()方法里,open()只会调用一次,process只是处理数据,这样就不会内存泄漏。我上次就是犯了这个错,排查了一天才找到问题,血的教训。
2. 类型匹配坑:注意输入输出的数据类型
IoTDB对数据类型检查很严格,如果你的UDF输入要求是DOUBLE,你传了INT,IoTDB会帮你隐式转换,一般没问题,但如果你输入要求是INT,传了DOUBLE,就会丢精度,我上次算峭度,把DOUBLE写成FLOAT,结果小的故障特征被精度丢了,一直出不对的结果,排查了一天才发现是类型错了。开发的时候一定要对应好数据类型,不确定就都用DOUBLE。
3. 线程安全坑:不要用非线程安全的全局变量
IoTDB是多线程查询,同一个UDF实例可能会被多个线程同时调用,如果你的UDF有全局变量,而且是非线程安全的(比如SimpleDateFormat),就会出现并发问题,结果错,甚至抛异常。正确的做法是:不要用全局可变变量,需要的变量要么放到方法里作为局部变量,要么如果是不可变的配置,作为常量也没问题。我上次用了全局SimpleDateFormat,并发查询的时候时间戳错乱,结果一直不对,查了好久才发现是并发问题。
4. 打包依赖坑:第三方依赖要用shade插件打包
如果你开发UDF用到了第三方依赖(比如我做FFT用到了JTransforms),不要直接把依赖放到lib里,要用Maven的shade插件把第三方依赖打包到你的jar包里,不然IoTDB核心里没有你的第三方依赖,运行就会抛ClassNotFound异常。给大家放一个我常用的shade配置:
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>3.2.4</version>
<executions>
<execution>
<phase>package</phase>
<goals>
<goal>shade</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>
这样打包出来的jar包就包含了你所有的第三方依赖,放到IoTDB里直接就能用,不会缺类。
5. 集群部署坑:每个数据节点都要放UDF包
IoTDB集群模式下,数据是分布在多个数据节点上的,查询会路由到对应的数据节点执行,所以你必须把你的UDF jar包放到所有数据节点的ext/udf目录下,不然查询路由到没有包的节点,就会报“找不到UDF”的错误。我一开始只放到了第一个节点,查了一半就报错,排查了好久才想到这个问题,集群部署一定要记住这点。
性能调优小技巧
- Java UDF性能肯定比Python好,生产上线能用Java就用Java,Python只用来做原型;
- 能提前初始化的资源都在
open()里做,不要在process里做,减少每次处理的开销; - 大的UDF不要注册太多,用完不用的UDF及时DROP,减少资源占用;
- Python UDF可以配置
python_udf.pool_size,根据你的并发需求调整进程池大小,避免不够用或者资源浪费。
七、UDF vs 其他扩展方案:什么时候该用UDF?
很多朋友会问,我有自定义计算需求,用UDF好还是用其他方式好?我给大家对比一下,大家可以根据自己的场景选:
| 扩展方式 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
| 自定义UDF | 性能高,SQL调用灵活,部署简单,就地计算减少传输 | 适合单查询内的计算,复杂多步骤逻辑不太方便 | 动态查询计算、自定义指标、个性化转换聚合,大部分场景都能用 |
| 触发器预计算 | 写入的时候就算好,查询的时候直接拿,查询速度快 | 规则改了要重新计算所有历史数据,不灵活 | 固定不变的指标,查询频繁 |
| 应用层客户端计算 | 灵活,开发自由度高 | 需要拉大量原始数据,传输开销大,延迟高 | 非常复杂的多步骤计算,需要和其他业务系统交互 |
| 存储过程 | 支持复杂多步骤逻辑 | 开发复杂,耦合度高 | 复杂的批处理任务 |
总结一下:只要是单查询内能完成的自定义计算,优先用UDF,性能好,灵活,开发简单,大部分需求都能满足。
八、总结
IoTDB的UDF扩展机制真的是解决个性化时序计算需求的杀器,设计非常简洁,开发门槛低,性能几乎没有损耗,完美解决了内建函数覆盖不了的个性化需求。我从最开始踩了一堆坑,到现在项目稳定运行,最深的感受就是:UDF把扩展的权利交给了用户,你不用等IoTDB官方给你加函数,自己就能快速实现,大大提升了开发效率,降低了系统成本。
更多推荐
所有评论(0)