TDgpt 算法扩展开发指南:为 TDengine 时序分析平台编写自定义分析与预测算法
【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine
TDgpt 是 TDengine 生态中可扩展的高级时序数据分析 Agent,本文基于官方开发者指南(docs/en/09-ai-and-advanced-analytics/01-tdgpt/06-dev/index.md)并对照仓库源码,系统讲解如何在 anonode 上开发、注册并调用自定义统计、机器学习与基础模型算法。读完本文,你将掌握 TDgpt 的算法加载机制、Python 类开发规范(命名、继承、属性初始化)、SQL 调用方式,以及如何编写一个可直接运行的异常检测与预测算法。
一、TDgpt 的扩展模型:半动态算法加载
TDgpt 的扩展性核心在于anode 的半动态算法加载机制:anode 启动时会扫描指定目录,将符合约定要求的 Python 文件中的算法注册到平台;之后通过 SQL 语句即可直接调用这些算法,无需重启 TDengine 服务端(taosd)。
整个注册链路在源码中非常清晰:
- tools/tdgpt/taosanalytics/app.py 在应用初始化时调用
loader.register_all_services(); - tools/tdgpt/taosanalytics/service_registry.py 中
register_all_services()依次扫描algo/ad、algo/fc、algo/imputat、algo/correl四个内置算法目录,以及algo/custom/ad、algo/custom/fc两个自定义目录; - 同时扫描动态模型目录(
dynamic_model_dir),目录下的 JSON 配置文件可随时增删,模型在服务运行期间即会被加载或卸载(见 service_registry.py 的sync_dynamic_services),这就是"半动态"的含义。
由于 TDgpt 与 TDengine 解耦,在 anode 上新增或升级算法不会影响 taosd 本身;应用侧只需更新 SQL 语句即可使用新算法。
二、添加算法到 TDgpt 的三个步骤
按官方指南,向 TDgpt 添加一个算法只需三步:
- 开发:按照 TDgpt 的规范用 Python 编写分析算法;
- 部署:将源码文件放入 anode 的指定目录,并重启 anode 服务;
- 注册:执行
CREATE ANODE语句,将 anode 加入 TDengine 集群。
之后算法即可通过 SQL 调用。需要说明的是,第 3 步仅在 anode 尚未注册到集群时执行一次;日常算法升级只需完成第 2 步(重启 anode),应用侧更新 SQL 即可。
2.1 准备开发环境
TDgpt 的源码位于本仓库的 tools/tdgpt 目录下,核心 Python 包为 tools/tdgpt/taosanalytics。开发环境要求 Python 3.10 及以上(见 tools/tdgpt/README.md),并依赖 numpy、statsmodels、scikit-learn、Flask 等库;安装完成后 taosanode 会创建虚拟环境venv并注册为系统服务taosanoded。
2.2 anode 目录结构
安装后的 anode 目录结构如下(摘自官方文档):
. ├── bin ├── cfg ├── lib │ └── taosanalytics │ ├── algo │ │ ├── ad │ │ └── fc │ ├── misc │ └── test ├── log -> /var/log/taos/taosanode ├── model -> /var/lib/taos/taosanode/model └── venv -> /var/lib/taos/taosanode/venv| 目录 | 说明 |
|---|---|
taosanalytics | 平台源码。algo子目录存放算法(ad为异常检测、fc为预测),test存放单元与集成测试,misc存放其他辅助文件 |
venv | Python 虚拟环境 |
model | 数据集对应的已训练模型 |
cfg | 配置文件 |
仓库中 tools/tdgpt/taosanalytics/algo/ad 与 tools/tdgpt/taosanalytics/algo/fc 的目录组织与安装布局一一对应,内置算法包括:
- 异常检测(
ad):ksigma、grubbs、iqr、lof、shesd、pyod_stat; - 预测(
fc):holtwinters、arima、prophet、theta、ets、ces、chronos、moirai、timemoe、timesfm、gpt。
三、算法开发硬性规范(决定能否被自动加载)
anode 采用自动扫描注册,因此文件命名、类命名、继承关系必须严格遵守以下约定,否则算法会被静默跳过。
3.1 存放位置限制
- 异常检测算法的 Python 源码必须放在
./taosanalytics/algo/ad; - 预测算法的 Python 源码必须放在
./taosanalytics/algo/fc; - 自定义算法可放入
algo/custom/ad与algo/custom/fc(见 service_registry.py)。
3.2 类命名规则
算法类名必须以下划线开头、以Service结尾。例如_KSigmaService就是 k-sigma 异常检测算法的类名(见 tools/tdgpt/taosanalytics/algo/ad/ksigma.py)。
这条规则在注册源码中得到了验证。在 service_registry.py 中,扫描器遍历模块内的所有类,并做如下过滤:
- 跳过抽象基类(
AbstractAnomalyDetectionService等); - 跳过不以
_开头的类; - 跳过定义在其他模块的类(防止误注册 import 进来的第三方类);
- 跳过
__init__.py、__pycache__及非.py文件。
if class_name in ServiceRegistry._base_class_name or ( not class_name.startswith("_") ): continue3.3 类继承规则
- 所有异常检测算法必须继承
AbstractAnomalyDetectionService,并实现execute方法; - 所有预测算法必须继承
AbstractForecastService,并实现execute方法。
这两个抽象基类定义在 tools/tdgpt/taosanalytics/base.py(异常检测)与 base.py(预测),其公共父类是AbstractAnalyticsService(再上层是AnalyticsService)。
AbstractAnomalyDetectionService的关键约定:
valid_code = 1:结果列表中等于valid_code的值表示正常点,否则视为异常点;type = "anomaly-detection":用于平台按类型索引算法;- 基类已实现
set_input_list(支持一维/二维输入)与set_params(支持valid_code参数)。
AbstractForecastService的关键约定:
- 内置参数
period(周期)、start_ts(起始时间戳)、time_step(步长)、rows(预测行数)、return_conf(是否返回置信区间)、conf(置信度,默认 0.95)、precision(时间精度,默认ms)、tz(时区); set_params强制要求start_ts、time_step、rows三个参数必须存在,且time_step、rows必须大于 0(见 base.py);- 预测结果以字典形式返回,
res中第一行为预测时间戳序列,后续行为预测值及可选的置信区间上下界。
对于基于 StatsForecast 的统计模型,还可继承 base.py 中定义的AbstractStatsForecastService,只需实现_fit_model并复用其execute,即可自动得到预测值、置信区间与 MSE。
3.4 类属性初始化
每个算法类必须初始化以下两个类属性:
name:算法标识符,只能使用小写字母。该标识符会显示在SHOW语句的算法列表中,并在 SQL 中通过'algo=name'指定;desc:算法的基础描述,会展示给用户并用于平台列表接口。
-- 示例:'algo' 键即取类中定义的 'name' 值 SELECT COUNT(*) FROM foo ANOMALY_WINDOW(col_name, 'algo=name')四、源码级示例一:编写异常检测算法(以 k-sigma 为例)
仓库自带的_KSigmaService(tools/tdgpt/taosanalytics/algo/ad/ksigma.py)是规范实现的最小完整范例,可作为模板:
"""ksigma class definition""" import numpy as np from taosanalytics.base import AbstractAnomalyDetectionService class _KSigmaService(AbstractAnomalyDetectionService): """KSigma algorithm is to check the anomaly data in the input list""" name = "ksigma" desc = """the k-sigma algorithm (or 3σ rule) expresses a conventional heuristic that nearly all values are taken to lie within k (usually three) standard deviations of the mean, and thus it is empirically useful to treat 99.7% probability as near certainty""" _builtins = True def __init__(self, k_val=3): super().__init__() self.k_val = k_val def execute(self): def get_k_sigma_range(vals, k_value): avg = np.mean(vals) std = np.std(vals) upper = avg + k_value * std lower = avg - k_value * std return [float(lower), float(upper)] if self.input_is_empty(): return [] threshold = get_k_sigma_range(self.list, self.k_val) return [-1 if k < threshold[0] or k > threshold[1] else 1 for k in self.list] def set_params(self, params): super().set_params(params) if "k" in params: k = int(params["k"]) if k < 1 or k > 3: raise ValueError("k value out of range, valid range [1, 3]") self.k_val = k def get_params(self): return {"k": self.k_val}对照规范逐项检查:
- 命名:文件为
ksigma.py,类名_KSigmaService以_开头、以Service结尾,满足扫描条件; - 继承:继承
AbstractAnomalyDetectionService并实现execute; - 属性:
name = "ksigma"(全小写)、desc描述了 3σ 准则原理,另设_builtins = True标记为内置算法; - 参数:通过重写
set_params支持k参数(取值范围 [1, 3],默认 3),get_params向平台暴露参数列表。
执行语义上,execute返回与输入等长的标签列表:正常点返回1(等于基类valid_code),异常点返回-1。上层 tools/tdgpt/taosanalytics/algo/anomaly.py 的do_ad_check会统计非valid_code的点数,并通过convert_results_to_windows将异常点转换为异常窗口供 SQL 侧使用。
编写同类算法时,仿照该结构即可:继承基类 → 初始化name/desc→ 在execute中读取self.list并返回标签列表 → 可选重写set_params支持自定义参数。
五、源码级示例二:编写预测算法(以 Holt-Winters 为例)
预测算法的参照实现是_HoltWintersService(tools/tdgpt/taosanalytics/algo/fc/holtwinters.py):
"""holt winters definition""" from statsmodels.tsa.holtwinters import ExponentialSmoothing, SimpleExpSmoothing from taosanalytics.algo.forecast import insert_ts_list from taosanalytics.base import AbstractForecastService class _HoltWintersService(AbstractForecastService): """Holt winters algorithm is to do the fc in the input list""" name = "holtwinters" desc = "forecast algorithm by using exponential smoothing" _builtins = True def __init__(self): super().__init__() self.trend_option = None self.seasonal_option = None def set_params(self, params): super().set_params(params) self.trend_option = params["trend"] if "trend" in params else None if self.trend_option is not None: if self.trend_option not in ("add", "mul"): raise ValueError("trend parameter can only be 'mul' or 'add'") self.seasonal_option = params["seasonal"] if "seasonal" in params else None if self.seasonal_option is not None: if self.seasonal_option not in ("add", "mul"): raise ValueError("seasonal parameter can only be 'mul' or 'add'") def get_params(self): p = super().get_params() p.update({"trend": self.trend_option, "seasonal": self.seasonal_option}) return p def execute(self): if self.list is None or len(self.list) < self.period: raise ValueError("number of input data is less than the periods") if self.rows <= 0: raise ValueError("fc rows is not specified yet") res, mse = self.__do_forecast_helper(self.list, self.rows) insert_ts_list(res, self.start_ts, self.time_step, self.rows) return {"mse": mse, "res": res}预测算法与异常检测算法的关键差异在于execute的返回结构:必须返回包含res键的字典,其中res的第一行为预测时间戳序列(由insert_ts_list基于start_ts与time_step生成),后续行为预测值序列及(若return_conf开启)置信区间上下界。上层 tools/tdgpt/taosanalytics/algo/forecast.py 的do_forecast会执行check_forecast_results校验各序列长度一致,并把period与algo写入结果。另外注意,当请求的算法名加载失败时,do_forecast会自动回退到holtwinters,因此 Holt-Winters 是平台默认的兜底预测算法。
AbstractStatsForecastService则提供了更省力的路径:子类只需实现_fit_model,基类的execute会统一完成数据校验、预测、置信区间计算与时间戳生成(base.py),arima、prophet、theta、ets等内置算法均基于此模式。
六、动态模型:免重启加载与卸载
除了内置 Python 算法,TDgpt 还支持通过JSON 配置文件注册动态模型。配置文件放在动态模型目录(默认model_dir/dynamic),格式支持模型名.json(旧格式)或模型名/模型名.json(新格式)两种布局,配置中必须包含algo字段(见 service_registry.py 的register_service_from_file)。
algo字段当前支持以下模型类型(见 service_registry.py):
- 预测(forecast):
arima、prophet、deepar; - 异常检测(anomaly):
iforest、svm; - 回归(regression):
linear_regression、lasso、ridge、elasticnet、svr、polynomial_regression; - 分类(classification):
logistic_regression、decision_tree。
动态模型与内置算法的差异在于:内置算法随代码加载,动态模型则从 JSON 配置 + 序列化的模型文件中加载。平台会为动态模型创建对应的DynamicForecastService/DynamicAnomalyService等包装服务(tools/tdgpt/taosanalytics/handlers/dynamic)。在服务运行期间,get_service每次调用都会触发sync_dynamic_services:新增配置文件会被注册,配置被删除的模型会被从内存移除,孤立的.pkl文件会被清理(service_registry.py)——这就是前文所述"半动态"机制的实际落地。需要训练与部署动态模型时,可参考 tools/tdgpt/taosanalytics/misc/train_ad_model.py 及 tools/tdgpt/taosanalytics/misc/model_downloader.py。
七、在 SQL 中调用自定义算法
算法注册成功后,即可通过 TDengine SQL 直接调用。官方指南给出了异常检测的调用示例:
SELECT COUNT(*) FROM foo ANOMALY_WINDOW(col_name, 'algo=name')其中'algo=name'的name就是算法类中定义的name属性。预测算法同理,例如使用内置的 Holt-Winters 或时间序列基础模型:
SELECT forecast(val, 'algo=holtwinters') FROM foo; SELECT forecast(val, 'algo=tdtsfm_1') FROM forecast.electricity_demand;forecast支持通过 options 传入预测行数、置信度、时间精度等参数;tools/tdgpt/taosanalytics/algo/forecast.py 的add_forecast_params展示了参数映射关系:forecast_rows→rows、start→start_ts、every→time_step、conf→ 置信度、return_conf→ 是否返回置信区间、prec→ 精度、tz→ 时区。
7.1 anode 的集群注册与管理
首次使用时需将 anode 注册进 TDengine 集群,相关管理语句(见 docs/en/09-ai-and-advanced-analytics/01-tdgpt/03-management.md):
CREATE ANODE {node_url}; -- 注册 anode SHOW ANODES; -- 查看已注册 anode SHOW ANODES FULL; -- 查看详情 DROP ANODE {anode_id}; -- 移除 anode UPDATE ANODE {anode_id}; -- 更新 anode 信息 UPDATE ALL ANODES;7.2 通过 RESTful 接口校验算法
anode 以 Flask 提供 RESTful 服务(tools/tdgpt/taosanalytics/app.py),可用以下接口快速验证算法是否注册成功:
GET /:返回 anode 版本号;GET /status:返回服务状态;GET /list:返回全部可用算法清单(含类型、名称、描述、参数与状态,见loader.get_service_list(),service_registry.py);GET /models:返回可用模型列表。
八、服务配置与验证测试
8.1 关键配置项
anode 的配置文件为taosanode.config.py(Linux 下默认位于/etc/taos/,仓库模板见 tools/tdgpt/cfg/taosanode.config.py),也可通过环境变量TDGPT_CONF指定路径(app.py)。常用配置项包括:
| 配置项 | 默认值 | 说明 |
|---|---|---|
bind | 0.0.0.0:6035 | 服务监听地址与端口 |
workers | 2 | worker 进程数(建议 2×CPU 核数 + 1) |
threads | 由 CPU 核数推导 | 每进程线程数(适合模型部署场景) |
timeout/keepalive | 1200 | 超时与 keep-alive(秒) |
log_level | DEBUG | 日志级别 |
model_dir | <install_dir>/model | 模型存储目录 |
dynamic_model_dir | model_dir/dynamic | 动态模型目录,运行期热加载 |
draw_result | False | 是否绘制查询结果图(输出到img_dir) |
models | 见配置文件 | 时间序列基础模型(tdtsfm、timemoe为必需,moirai、chronos、timesfm、moment为可选)及其端口/端点,用于自动推导服务 URL(conf.py) |
默认路径定义可在 tools/tdgpt/taosanalytics/conf.py 中查到(日志/var/log/taos/taosanode/、模型目录/usr/local/taos/taosanode/model/、动态模型目录model/dynamic)。修改配置后重启服务即可生效:systemctl restart taosanoded。
8.2 使用单元测试验证算法
仓库在 tools/tdgpt/tests 与 tools/tdgpt/taosanalytics/test 提供了大量测试用例。例如 tools/tdgpt/tests/unit_test.py 中的test_generate_anomaly_window直接验证了异常检测结果向异常窗口的转换逻辑(convert_results_to_windows),forecast_test.py、anomaly_test.py、restful_api_test.py则分别覆盖预测、异常检测与接口层。开发新算法时,建议参照这些测试为execute的正确性、set_params的参数校验(非法参数应抛出ValueError)补充用例,并运行pytest确保回归。
九、开发建议与常见问题
- 算法未出现在
SHOW列表中:优先检查类名是否以_开头且以Service结尾、类是否定义在扫描目录的模块内、name是否全小写、文件是否为.py且非__init__.py;查看 anode 日志(/var/log/taos/taosanode/taosanode.app.log)中的load algorithm与failed to register service记录。 name冲突:注册表不允许重名,重复注册会抛出RuntimeError(service_registry.py),请确保name全局唯一。- 输入数据校验:异常检测的
execute应先调用input_is_empty()处理空输入;预测算法应在execute中校验数据量是否满足周期要求(如 Holt-Winters 要求数据量不小于period)。 - 参数校验:
set_params中所有非法参数都应抛出ValueError,平台会捕获并返回错误信息(见 tools/tdgpt/taosanalytics/handlers/forecast.py)。 - 动态模型优于内置算法升级:如果算法只是模型权重更新,优先使用动态模型目录热替换 JSON 配置,避免重启 anode 影响正在执行的查询。
- Python 版本兼容:注册器会对特定 Python 版本做兼容性过滤(例如 Python 3.12 下因 pandas 兼容问题跳过
_SHESDService,见 service_registry.py),自研算法应关注运行环境的 Python 版本。
十、总结
TDgpt 通过"约定优于配置"的目录扫描机制,将自定义算法的接入成本降到最低:把遵循命名、继承与属性规范编写的 Python 文件放入algo/ad或algo/fc目录并重启 anode,算法即可通过'algo=name'在 SQL 中调用,而动态模型机制则进一步支持运行期免重启的模型热更新。本文给出的 k-sigma 与 Holt-Winters 两个源码范例覆盖了异常检测与预测两大算法类型的完整开发范式,读者可直接基于 tools/tdgpt/taosanalytics/algo 中的任一内置实现仿写,配合 tools/tdgpt/tests 的测试用例快速验证,即可将自有算法无缝接入 TDengine 的时序分析链路。
【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考