新用户入门指南(工业篇)

注:

教程目标:通过一个工业传感器数据处理案例,快速了解如何用 DolphinDB 完成物联网数据的建模、接入、分析、实时计算和可视化监控。

学习收获:

  • 了解 物联网时序数据的核心特点与常见需求。

  • 掌握 DolphinDB 的建模思路,包括分区规划、存储引擎选择、表结构和字段设计。

  • 学会 接入实时数据、历史数据和离线文件。

  • 掌握 常见查询分析方法,如设备查询、窗口聚合、状态统计和异常分析。

  • 理解 流计算框架,构建实时清洗、指标计算和告警链路。

  • 了解 DolphinDB 与 Grafana 等工具的可视化集成方式。

1. 概述与环境准备

本教程面向物联网入门用户,以工业传感器产生的时序数据为主线,围绕数据建模、数据接入、查询分析、实时计算与可视化监控等核心环节,介绍如何使用 DolphinDB 搭建物联网数据处理与监控的基础流程。

1.1 DolphinDB 在物联网场景下的应用

物联网场景需要对海量工业传感器产生的时序数据进行实时接入、存储、分析与计算,具有数据规模大、采集频率高、实时性要求强等典型特征。DolphinDB 依托高性能时序数据库引擎、分布式计算能力与流批一体架构,已在能源电力、高端制造、科学研究、公用事业及快消品等行业实现广泛应用。

  • 电力行业:发电设备监测、电网运行分析、智能电表管理、电力交易分析等,支撑设备状态监控、故障预警与负荷分析。

  • 高端制造:工业设备数据采集、产线监控、质量追溯、设备运维及测试数据分析,助力企业提升生产效率与设备利用率。

  • 科学研究:大型实验装置与科研设备的高频数据管理,实现实验数据存储、趋势分析、异常回溯与科研计算。

  • 公用事业:可应用于水务、燃气等行业,实现管网监测、设备运行分析、告警管理和运营优化。

  • 快消品:可用于生产设备监控、冷链物流管理、仓储运营分析和供应链优化,保障产品质量和运营效率。

总体来看,DolphinDB 在物联网场景中的典型应用包括设备数据统一接入、海量时序数据存储、实时监控与告警、设备状态分析、质量追溯、预测性维护以及可视化展示等,为企业构建统一的数据分析平台提供支撑。

1.2 环境准备与教程说明

开始前,请先完成 DolphinDB 的下载与部署。安装包可在 DolphinDB 产品下载页 获取,具体部署步骤请参考 部署。本文实验环境基于 DolphinDB server 3.00.5 2026.02.06 LINUX x86_64,使用社区版 license,并采用单节点方式启动。通过 Web 管理界面、VS Code 插件、DolphinDB GUI 或 Shell 执行脚本与 DolphinDB 数据库进行交互。在生产环境中,可根据数据规模、写入压力和查询并发情况,进一步扩展为分布式集群部署。

1.3 DolphinDB 架构简介

物联网场景中,随着设备数量、采样频率和数据保留周期不断增加,系统通常需要具备高吞吐写入、低延迟查询、并行计算和水平扩展能力。DolphinDB 采用 Shared-nothing 分布式架构,各节点拥有独立的计算与存储资源,可通过增加节点数量提升整体处理能力。

1. 图 1-1 DolphinDB 系统架构图

DolphinDB 主要有单节点和集群两种运行模式:单节点适合本地学习、功能验证和小规模试用;集群模式适合生产环境,可支撑更大的数据规模、写入压力和查询并发。

集群模式下,DolphinDB 通常由 Controller、Agent、DataNode 和 ComputeNode 组成。Controller 负责集群状态、分布式文件系统元数据和事务日志管理;Agent 负责接收 Controller 指令并管理节点启停;DataNode 负责数据存储,并承担主要查询与计算任务;ComputeNode 只负责计算和查询处理,不存储数据,适用于计算与存储分离的场景。

对于入门用户,可以简单理解为:Controller 管理整个集群,Agent 管理节点启停,DataNode 负责存储和计算,ComputeNode 负责计算但不存储数据。

2. 数据建模

本章以一个典型的工业传感器监控场景为例。假设某工厂有 10,000 个传感器,部署在不同产线和设备上。每个传感器每 30 秒上报一次数据,上报的指标包括温度、压力、湿度、电压、电流和设备状态。

数据量估算:

  • 每个传感器每天上报:24 × 60 × 2 = 2,880 条

  • 10,000 个传感器每天产生:2,880 × 10,000 = 2,880 万条记录

  • 单条记录约 50~60 字节,每天原始数据约 1.5 GB

查询场景方面,物联网数据有两个最典型的查询特征:一是按设备和时间范围查询历史数据(例如“查 device_0001 昨天一整天的平均温度”),二是获取所有设备的最新状态(例如“当前所有设备的实时温度是多少”)。这两个特征会直接影响后续的分区策略和存储引擎选择。

为什么先算数据量? 数据规模决定了分区方案、存储引擎和硬件资源配置。

2.1 表结构设计

明确了数据规模后,接下来定义每条上报数据包含哪些字段。

字段 类型 说明 设计原因
device_id SYMBOL 设备 ID 适合作为 HASH 分区列和排序列,便于按设备查询
ts TIMESTAMP 数据时间戳 适合作为时间分区列,支持时间范围查询和窗口计算
temperature DOUBLE 温度 数值型测点,便于聚合分析
pressure DOUBLE 压力 数值型测点,便于聚合分析
humidity DOUBLE 湿度 数值型测点,便于聚合分析
voltage DOUBLE 电压 数值型测点,便于聚合分析
current DOUBLE 电流 数值型测点,便于聚合分析
status INT 设备状态 适合开停机、告警状态等离散状态统计

字段类型选择说明:

  • 时间字段 ts:使用 TIMESTAMP(毫秒精度),字段名统一为 ts,便于按时间范围查询和窗口计算。

  • 设备编号 device_id:使用 SYMBOL 类型。SYMBOL 是 DolphinDB 的字符串压缩存储类型,适合存储重复度高的字符串(如设备 ID),可大幅节省存储空间。

  • 数值型测点:温度、压力等使用 DOUBLE,兼顾精度和存储效率。若指标精度要求不高(如整数型状态值),可使用 INTFLOAT 进一步节省空间。

  • 状态字段 status:使用 INT,便于分组统计和条件判断。

设计原则:字段设计应围绕“查询条件、写入频率、字段类型”展开,将时间、设备 ID、测点值、状态字段统一规划,避免后续频繁改表。

2.2 宽表 vs 窄表 选型

在物联网场景中,同一设备通常会采集多个传感器指标,例如温度、电压、湿度、压力、电流、转速、状态量等。建模时需要判断这些指标应存储在一张“宽表”中,还是拆分为“窄表”或多张主题表。

判断维度 更适合宽表 更适合窄表
上传时间是否一致 同一设备多个指标总是在同一时间点一起上报 不同指标上报频率不同,时间戳不完全一致
指标数量是否稳定 指标集合相对固定,很少新增或下线 指标经常新增、下线或字段变化频繁
查询方式 经常按设备和时间查询多个指标,例如查看某设备一段时间内的温度、电压、电流 经常按指标查询,例如查询所有设备的某一个指标
数据稀疏程度 大部分指标每次都会上报 部分指标大量为空,存在明显稀疏数据
分析方式 多指标联合分析较多,例如温度、电流、压力之间的关联分析 单指标分析较多,例如单独分析温度曲线或电压曲线

在本教程的示例场景中,假设同一设备的温度、压力、湿度、电压、电流、状态等指标在同一时间点统一上报,且指标集合相对稳定,因此推荐采用宽表模型。宽表可以减少关联查询,提升按设备、按时间范围查询时的读取效率,也更便于后续进行多指标聚合分析和状态监控。

如果实际项目中存在大量异构设备、指标频繁变化、不同指标采样频率差异较大等情况,可以采用“按设备类型建宽表”或“核心指标宽表 + 扩展指标窄表”的混合建模方式,兼顾查询性能和模型灵活性。

2.3 存储引擎选型

在 DolphinDB 中,存储引擎的选择直接影响到写入性能、查询效率以及数据压缩比。在物联网场景下,OLAP、TSDB、IOT 是三种常用的存储引擎,分别用于覆盖从传统数仓离线分析到时序数据存储查询的不同需求。关于 OLAP、TSDB、IOT 存储结构的描述性说明,请参照教程:OLAP 存储引擎TSDB 存储引擎 以及物联网点位管理引擎

存储引擎 索引机制 查询特点 适用场景
TSDB sortKey 索引 + zonemap 按设备+时间查询极快,聚合高效 通用时序监控、历史数据查询
IOT sortKey + 最新值缓存 最新值查询毫秒级响应 海量点位管理、高频最新值查询
OLAP 无独立索引,依赖分区剪枝 全表扫描友好,点查慢 传统数仓、离线分析

TSDB 通过排序列和分区裁剪提升时序查询效率,适合本教程中“按设备 + 时间范围”查询、窗口聚合和状态统计等需求。IOT 更偏向海量测点和最新值查询优化,适合后续进阶学习;OLAP 更适合传统离线分析,不适合作为本教程宽表时序数据的首选。

2.4 分区设计

上一节确定了使用 TSDB 引擎,但 TSDB 的高效查询还依赖于合理的分区设计。如果分区方案不合理,即使选对了引擎,查询性能也可能大打折扣。本节将围绕“数据怎么分,查询才能快”这个目标,从分区机制讲起,逐步推导出本教程的分区方案。

2.4.1 分区机制与作用

在 DolphinDB 的分布式架构中,海量数据的存储与并行计算主要通过分布式表来实现,而分区机制是分布式表组织数据的核心方式。对用户而言,分布式表通常指存储在以 dfs:// 开头的分布式数据库中的表。

如果所有数据都放在一张大表中,查询时可能需要扫描大量无关数据。DolphinDB 通过分区机制将一张逻辑表拆分为多个物理分区,查询时只访问命中的分区,从而减少扫描范围,提高查询性能。

对数据库进行分区主要有以下好处:

  • 缩小查询范围:查询条件命中分区列时,只扫描相关分区,减少 I/O 开销。

  • 提高计算性能:多个分区可以分布在不同节点上,并行执行查询和计算任务。

  • 支撑高可用:分区数据可以配置副本,节点故障时可从副本读取数据。

  • 支持水平扩展:数据规模增长后,可通过增加节点承载更多分区。

2. 图 1-2 时间分区示意图

图 1-2 以时间分区为例,展示数据在逻辑分区、副本和数据节点之间的对应关系。

如上图所示,逻辑表 sensor_data 按时间字段进行分区。以 2025年1月1日的数据为例,系统会将其写入分区 2025.01.01,并为该分区维护多个副本,分别存储在不同的数据节点上。当某个数据节点发生故障时,由于其他节点上仍存在该分区的数据副本,系统可以继续从副本读取数据,从而保障数据访问的连续性和高可用性。

2.4.2 分区类型与选型

DolphinDB 支持多种分区类型:范围分区(RANGE)、哈希分区(HASH)、值分区(VALUE)、列表分区(LIST)与复合分区(COMPO),详情可见:数据分区 。物联网场景中最常见的是以下三种:

分区类型 说明 适用场景
VALUE(值分区) 每个唯一值作为一个分区 适用于基数较小且固定的字段,如日期
HASH(哈希分区) 通过哈希函数将数据分配到指定的分区 适用于基数很大且无明显分布规律的字段,如设备 ID
COMPO(复合分区) 多层级分区(两层或三层),每层可选用不同的分区类型 数据量大且查询条件经常涉及多列时使用。

复合分区是 DolphinDB 处理海量时序数据的首选方案。它允许用户结合时间维度和非时间维度进行分层分区,例如“第一层按日期 VALUE 分区 + 第二层按设备 ID HASH 分区”。这种设计既能利用时间分区实现快速的日期范围裁剪,又能利用设备 HASH 分区实现并行写入和查询。

2.4.3 分区粒度与数据量估算

根据分布式架构的特性,分区粒度(每个分区包含的数据量)对性能影响显著,需在"分区过大"和"分区过小"之间找到平衡点。

分区粒度过大带来的问题:

  • 查询时 IO 开销高。即使只需要几条记录,也可能需要读取整个分区的元数据

  • 内存压力大。DolphinDB 在查询时需要加载分区元数据,过大的分区可能触发 license 内存限制

  • 数据维护困难。压缩、备份等操作的最小单位是分区,过大的分区意味着更粗的维护粒度

调整分区粒度的方法:

  • 降低分区粒度:(1)采用 COMPO 分区;(2)增加分区个数;(3)将 RANGE 分区改为 VALUE 分区。

  • 增加分区粒度:(1)采用 RANGE 分区取代 VALUE 分区;(2)HASH 分区减小哈希映射数。

分区粒度过小带来的问题

  • 查询涉及分区过多,节点间通信开销增加。极端情况下,一次查询可能需要访问数万个分区

  • 控制节点元数据膨胀。每个分区都需要在控制节点注册元信息,过多的小分区会导致控制节点内存不足

  • 涉及的分区多,导致系统读写频繁,从而造成很多低效的磁盘访问(小文件读写),造成系统负荷过重。

推荐的分区大小(压缩前)如下表所示:

存储引擎 推荐分区大小 说明
OLAP 100MB - 300MB 分析型场景,适合批量读取
TSDB 400MB - 1GB 时序优化引擎,需要更大的块来发挥压缩优势
IOT 400MB - 1GB 物联网专用引擎,与 TSDB 类似

为了估算整体的数据量大小,可以先估算单行数据的大小。

根据表字段的存储的字节数估算。可通过 schema 函数获取表字段的类型信息,再对照 数据类型中对应类型的字节数,将各字段字节数相加即可估算单行数据的大小,这里的估算方式为:单行数据大小 = Σ字段类型对应字节数。

// 调用 schema 函数获取表字段信息(字段名,类型)
t.schema().colDefs

/* output:
name        typeString typeInt extra comment  sensitive
----------- ---------- ------- ----- -------- ---------
device_id   SYMBOL     17            设备id 0        
ts          TIMESTAMP  12            日期   0        
temperature DOUBLE     16            温度   0        
pressure    DOUBLE     16            压力   0        
humidity    DOUBLE     16            湿度   0        
voltage     DOUBLE     16            电压   0        
current     DOUBLE     16            电流   0        
status      INT        4                   0   
*/
单行大小 ≈ 4 (device_id) + 8 (ts) + 8×5 (5 个 DOUBLE 指标) + 4 (status) = 56 字节
28800000 行 * 56 bytes ≈ 1,612,800,000 bytes ≈ 约 1.5 GB

如果仅按天分区,每天约 1.5 GB,接近或超过推荐分区大小。因此,需要采用组合分区来降低单个分区的数据量:采用“日期 VALUE 分区 + 设备 ID HASH 分区”的组合分区,降低单分区数据量并提升并行处理能力。

2.4.4 HASH 分区规划

HASH 分区数需在建库时确定,后续调整通常需要迁移或重写历史数据。因此,建议在项目初期结合未来 1~3 年的数据增长预期一次性规划。

本教程的规划计算

参数 数值
每天数据量 1.5 GB
HASH 分区数 5 个
每个子分区数据量 1.5 ÷ 5 ≈ 300 MB

300 MB 略低于 TSDB 推荐的下限(400MB),但在可接受范围内。如果选择 4 个 HASH 分区:

参数 数值
HASH 分区数 4 个
每个子分区数据量 1.5 ÷ 4 = 0.375 GB ≈ 384 MB

4 个 HASH 分区时,接近推荐范围的下限,但还需要考虑未来增长:

  • 当前 10,000 个设备,未来 1~3 年可能增长到 20,000 个

  • 数据量翻倍时,每个子分区也会翻倍:384 MB × 2 = 768 MB,仍在推荐范围内

  • 如果预留 5 个分区,未来翻倍后:1.5 × 2 ÷ 5 = 600 MB,也在推荐范围内

综合建议

分区数 当前单分区大小 未来翻倍后 评价
4 个 0.375 GB(384 MB) 0.75 GB(768 MB) 当前略低但可接受,未来表现良好
5 个 0.30 GB(300 MB) 0.60 GB(600 MB) 当前偏小,未来更宽松

两个方案均可。本教程建议预留 4~5 个 HASH 分区,具体数值可根据对未来数据增长的信心来定:

  • 如果认为设备数量增长较快(1~2 年内翻倍),建议选 5 个

  • 如果增长较为平缓,选 4 个 即可

重要提醒:HASH 分区数一旦建库确定,后续无法直接修改。如果要调整,需要新建库表并重新导入历史数据。因此请结合业务规划一次性做好决策。

本教程最终分区方案汇总

层级 分区类型 分区列 分区数/范围
第一层 VALUE(值分区) ts(日期) 按天,每天一个分区
第二层 HASH(哈希分区) device_id(设备 ID) 4~5 个(本教程示例采用 5 个)

2.5 建库建表实战

DolphinDB SQL 建表语法和传统 SQL 类似,除了指定的表名、字段名、字段类型外,额外增加了一些数据库的配置参数。配置参数的用法详见create 参数说明。

create table dbPath.tableName (
    schema[columnDescription] 
    // 字段名 字段类型 
    // [comment = 字段注解, compress = 压缩方案, index = 向量索引]
)
[partitioned by partitionColumns], // 分区列
[sortColumns], // 排序列
[keepDuplicates=ALL], // 去重机制
[sortKeyMappingFunction] // 排序列映射函数
[softDelete=false] // 软删除配置
[comment] // 表注释

根据物联网传感器数据最佳实践方案,创建数据库 "dfs://iot_sensor" 下的分布式表 "sensor_data",设置分区列为 ts(对应按日 VALUE 分区)和 device_id(对应 SYMBOL HASH 分区)。后续数据写入时,将依据这分区字段的值,将数据划分存储到不同的分区中。

// 步骤1:创建复合分区数据库
create database "dfs://iot_sensor"
partitioned by VALUE(2025.01.01..2025.01.03), HASH([SYMBOL, 5])
engine='TSDB'
// 步骤2:创建分区表
// 表结构包含:设备标识、时间戳、温度、压力、湿度、电压、电流、设备状态
// 按 ts + device_id 复合分区,按 device_id + ts 排序以优化查询
create table "dfs://iot_sensor"."sensor_data" (
    device_id SYMBOL[comment="设备id"],
    ts TIMESTAMP[comment="日期", compress="delta"],
    temperature DOUBLE[comment="温度"],
    pressure DOUBLE[comment="压力"],
    humidity DOUBLE[comment="湿度"],
    voltage DOUBLE[comment="电压"],
    current DOUBLE[comment="电流"],
    status INT
)
partitioned by ts, device_id
sortColumns=["device_id","ts"]

如果希望通过图形化方式完成库表设计,也可以使用 DolphinDB Web 界面库表创建指南。该工具适合在不熟悉脚本的情况下配置分区方案、字段类型和排序列,可作为本节脚本建模方式的补充。

2.6 常见问题

Q1:是否建议按年分区?

通常不建议按年分区。按年分区容易导致单分区过大,影响并行查询、数据迁移和维护效率。即使业务按年归档,也建议数据库层面按分区;查询跨多个分区时,系统会自动合并结果。

Q2:分区建好后可以调整吗?

由于数据存储是依赖于分区的,因此不支持直接修改分区方案。最佳做法是重新创建一个新分区方案的库表,然后将数据迁移过去,最后删除原来的库表。

  • 可以追加分区:VALUE 分区可通过配置 newValuePartitionPolicy='add' 自动创建新值分区,也可使用 addValuePartitions 手动增加;RANGE 分区可使用 addRangePartitions 向后追加范围。

  • 不支持直接修改分区规则:例如将 HASH 分区数从 10 改为 20,通常需要新建库表并将历史数据重新写入。

Q3:查询慢时优先检查什么?

查询或者写入慢的原因多样,可以按照配置情况,查询/写入情况,库表设计,系统负载这几个角度去定位。具体可见教程:查询/写入慢

  • 是否命中分区剪枝WHERE 条件中应包含分区列,例如按天分区时明确指定时间范围,避免扫描无关分区。

  • 单分区是否过大:TSDB/IOTDB 场景下,建议单分区控制在 400 MB~1 GB。若超过范围,可增加二级分区,如按设备 ID 做 HASH 分区。

  • 是否存在全表扫描:避免无过滤条件的 SELECT *,建议只查询必要列,并结合时间、设备 ID 等高频条件缩小扫描范围。

3. 数据接入

在物联网项目中,数据通常来自网关、传感器、业务系统、消息队列或历史文件。不同来源的数据链路不同,因此接入方式也不完全一样。下表整理列出了数据导入 DolphinDB 数据库的不同方式、对应方法以及相关参考文档的链接:数据导入方法

1. 表 3-1 DolphinDB 物联网数据接入方法
接入方式 典型来源 方法 适用特点
API 写入 网关程序、边缘服务、业务应用 通过 DolphinDB C++ API,Python API,Java API 等 API 所提供的接口导入数据 应用程序直接写入,开发灵活
消息中间件接入 Kafka、MQTT Broker、物联网平台等 通过 Kafka、MQTT 等消息中间件的插件接入 削峰填谷,解耦生产与消费
数据库插件同步 MySQL、Oracle、SQL Server 等 通过 MySQL 和 ODBC 插件同步数据 存量数据迁移,异构数据库同步
文件导入 CSV、TXT、Parquet 等 通过内置文本文件加载函数 loadText、ploadText、loadTextEx、textChunkDS 离线批处理,历史数据迁移

以下以一个简单的传感器数据表为例,贯穿介绍多种接入方式。

3.1 API 接入

通过编程 API 写入数据是最通用的接入方式之一,适用于网关程序、边缘服务和业务应用等场景。应用侧可先完成数据采集、清洗和格式转换,再调用 DolphinDB API 将数据写入数据库。DolphinDB 提供 Python、Java、C++、Go 等多种语言 API,便于与现有系统集成。具体的 API 相关教程可见:连接器 和 API

DolphinDB 提供多种 API 写入方式,不同方式适用于不同的数据规模和实时性要求。

  • 单条写入(insert into):适用于测试验证或低频数据场景,实现简单且便于调试,但吞吐量低,不适合生产环境。

  • 批量写入(tableInsert:适用于周期性上传或设备数据攒批写入,网络开销小、性能较好,但需要客户端自行缓存数据。

  • TableAppender:适用于低频批量写入场景,接口简单、易于上手,支持批量追加且便于与现有脚本集成,但吞吐量有限,需自行控制批大小。

  • PartitionedTableAppender:适用于高频大批量数据写入,数据自动路由到分区,可充分利用多分区并行写入优势,但需要正确设置分区键且单批数据量需足够大并涵盖多分区才能充分发挥性能。

  • 异步写入 MultithreadedTableWriter (MTW):适用于高并发、高吞吐的 IoT 实时流以及稳定低延迟的单行持续写入,通过多线程并行、客户端缓冲与批量发送显著提升吞吐量,支持失败重试与异步回调,降低单条写入开销,但需合理配置 batchSize 和线程数。

下面示例演示如何用 Python 生成一批模拟传感器数据,并通过 tableInsert 批量写入 DolphinDB 分布式表。

import dolphindb as ddb
import pandas as pd
# 1. 连接 DolphinDB
s = ddb.session()
s.connect("127.0.0.1", 8848, "admin", "123456")
# 2. 构造一批模拟数据
data = pd.DataFrame({
    "device_id": ["dev001", "dev002", "dev003"],
    "ts": pd.to_datetime([
        "2025-01-01 10:00:00.001",
        "2025-01-01 10:00:00.002",
        "2025-01-01 10:00:00.003"
    ]),
    "temperature": [25.6, 26.1, None],
    "pressure": [101.3, None, 102.1],
    "humidity": [60.5, 58.2, None],
    "voltage": [None, None, 220.5],
    "current": [0.5, None, None],
    "status": [0, 1, 0]
})
# 3. 上传 DataFrame 到 DolphinDB 会话
s.upload({"data": data})
# 4. 批量写入分布式表
s.run("""
pt = loadTable("dfs://iot_sensor", "sensor_data");
tableInsert(pt, data)
""")
print("write success")

入门学习或功能验证阶段,可先使用小批量写入,便于理解流程和调试脚本;生产环境建议根据数据量和实时性要求,采用批量攒批、多线程并发或异步写入方式提升吞吐量。

3.2 消息中间件接入

消息中间件是物联网实时接入中常用的数据通道,可将设备端、网关和数据平台解耦,提升系统的扩展性、可靠性和削峰能力。DolphinDB 可通过 Kafka、MQTT 等插件订阅设备消息,并将数据实时写入流数据表或分布式表,形成从设备采集、消息消费到实时分析的完整链路。详细用法可参考 Kafkamqtt

本教程以 Kafka 为例,演示如何使用 DolphinDB Kafka 插件订阅 Kafka Topic,解析设备消息,并将数据写入流数据表或分布式表,完成实时接入链路的基础搭建。

Kafka 接入流程可概括为以下几个步骤:

  1. 加载 Kafka 插件,并完成基础连接配置;

  2. 创建 Kafka Consumer,指定 Broker、消费组和 offset 策略;

  3. 创建流数据表,作为 Kafka 消息进入 DolphinDB 的接收通道;

  4. 定义消息解析函数,将 Kafka 消息转换为目标表结构;

  5. 订阅 Kafka Topic,启动后台消费任务;

  6. 将解析后的数据写入流表,并通过订阅任务异步持久化到分布式表。

Kafka 接入脚本示例如下,具体代码可见附录 createKafka.dos

loadPlugin("kafka")

consumer=kafka::consumer(...)

kafka::subscribe(...)

kafka::createSubJob(...)

3.3 数据库插件同步

数据库插件主要用于将 MySQL、Oracle、SQL Server、PostgreSQL 等关系型数据库中的存量数据迁移或同步到 DolphinDB,实现统一存储与分析。典型场景包括:迁移设备台账、历史采样数据、告警记录、业务明细,以及定期同步增量数据。

DolphinDB 提供 mysql(对接 MySQL)和 ODBC(对接 Oracle、SQL Server、PostgreSQL 等)。

以 ODBC 插件为例,使用前需安装对应数据库驱动并配置 ODBC 数据源(DSN)。

// 1. 加载 ODBC 插件
loadPlugin("odbc")
// 2. 建立 ODBC 连接
conn = odbc::connect("DSN=OracleDB;UID=iot_user;PWD=iot_password")
t = odbc::query(
    conn, 
    "select device_id, ts, temperature, pressure, humidity, 
    voltage, current, status from sensor_history 
    where 
    ts >= timestamp '2025-01-01 00:00:00' 
    and 
    ts < timestamp '2025-01-04 00:00:00'"
    ) 
t = select 
    device_id$SYMBOL as device_id, 
    timestamp(ts) as ts, 
    temperature$DOUBLE as temperature, 
    pressure$DOUBLE as pressure, 
    humidity$DOUBLE as humidity, 
    voltage$DOUBLE as voltage, 
    current$DOUBLE as current, 
    status$INT as status 
    from 
    t
// 4. 写入 DolphinDB 分布式表
pt = loadTable("dfs://iot_sensor", "sensor_data")
tableInsert(pt, t)

ODBC 插件能够统一访问多种关系型数据库,是异构数据库迁移和数据整合的常用方案。实际使用过程中,需要重点关注字段类型映射问题。例如 Oracle 的 NUMBER、DATE、TIMESTAMP 等类型,在导入 DolphinDB 时应映射为对应的数据类型。建议先使用少量样例数据验证字段映射和数据质量,再执行大规模数据导入。

3.4 文件导入

文件导入适用于历史数据迁移、离线批处理、数据验证和大规模回灌等场景。DolphinDB 提供多种文本文件处理函数:loadText 可将小文件加载为内存表,便于预览和清洗;loadTextEx 可将大文件直接导入分布式表,避免全量加载到内存;ploadText 支持并行加载;textChunkDS 适合按块处理超大文件。

本教程以 loadText 为例进行演示。该函数使用简单,导入后生成普通内存表,适合在写入分布式表前预览数据、检查字段类型、补充字段或完成简单清洗,因此更适用于学习、测试和小规模数据验证阶段。

// 将小型 CSV 文件加载到内存表
data = loadText("/path/to/small_data.csv")
由于数据已经加载到内存中,因此还可以在写入前进行检查:
// 查看数据
select top 10 * from data
// 查看字段类型 
schema(data)
// 将内存表追加到分布式表
pt = loadTable("dfs://iot_sensor", "sensor_data")
tableInsert(pt, data)

3.5 常见问题

Q1:写入时报类型不匹配怎么办?

优先检查字段顺序和字段类型。比如设备编号应为字符串或 SYMBOL,采样值应为 DOUBLE,采样时间应为 TIMESTAMP。如果 Python DataFrame 中的时间字段没有正确转换为 datetime 类型,写入时也可能报错,详细可见 PROTOCOL_DDB

Q2:API 写入性能不高怎么办?

不要使用单条 insert into 高频写入。建议改为批量 tableInsert、PartitionedTableAppender 或 MultithreadedTableWriter。同时检查 batchSize、线程数和分区字段是否合理。

Q3:Kafka 或 MQTT 连接不上怎么办?

先检查 DolphinDB 服务器是否能访问 Kafka Broker 或 MQTT Broker 的 IP 和端口。再检查用户名密码、Topic 名称、Consumer Group 和网络防火墙配置。

Q4:文件导入失败怎么办?

先检查文件编码、分隔符、表头字段名和时间格式。对于大文件,可以先截取前 100 行进行验证,再导入完整文件。

4. 数据查询与分析

完成数据建模和接入后,如何高效查询和分析这些数据是工程实践中的核心问题。本章介绍 DolphinDB 在物联网场景下的常用查询模式和分析方法,涵盖时序查询、聚合分析、状态监控、关联查询和数据分析等典型场景。

为便于验证本章查询示例,使用数据建模章节的物联网案例。该数据集包含 10,000 个测点,时间范围为 2025.01.01 至 2025.01.03,每 30 秒上报一次多指标数据,包括温度、压力、湿度、电压、电流和设备状态等字段,则每天的数据量约为 2,880 万条,适合 DolphinDB 社区版运行验证。请参考附录 genData.dos 模拟数据。

以下章节依次讲解常见的物联网统计分析需求实现,帮助读者快速上手 DolphinDB,并实现典型物联网分析需求。

4.1 基础查询

4.1.1 按设备与时间范围查询

查询设备 "device_0001" 在 2025-01-01 当天的 100 条数据记录。该查询利用分区剪枝,仅扫描涉及的分区。

select
    *
from
    loadTable("dfs://iot_sensor", "sensor_data")
where
    device_id = 'device_0001'
    and 
    ts 
    between 2025.01.01T00:00:00.000
    and 2025.01.01T23:59:59.999
limit
    100

上述语句也可以等价写为:

select
    *
from
    loadTable("dfs://iot_sensor", "sensor_data")
where
    device_id = 'device_0001'
    and date(ts) = 2025.01.01
limit
    100

在 TIMESTAMP、DATETIME 等细粒度时间列上使用 datemonth 等时间函数作为过滤条件时,DolphinDB 会自动进行解析优化,当分区列来源于时间字段时,可以触发分区剪枝。因此,上述查询只会扫描 2025.01.01 对应的日期分区;同时结合 device_id 的 HASH 分区,还可以进一步定位到该日期下对应的设备分区(如 2025.01.01/Key2)。这种写法更简洁、直观,推荐在日常查询中使用。当然,使用方式 1 中显式指定时间范围的写法同样正确,可根据个人编程习惯选择。

4.1.2 最新状态查询

获取设备最新状态是实时监控中的常见需求。DolphinDB 提供 context by 子句实现窗口计算,在窗口内的执行排序、应用函数计算、取 Top N 等语义。

select *
from loadTable("dfs://iot_sensor", "sensor_data")
context by device_id
csort ts desc
limit 1

上述查询按 device_id 分组,每组内按时间戳降序排列,取每组的第一条记录,即各设备的最新状态。context by 是 DolphinDB 的扩展语法,相比传统 SQL 的窗口函数更为简洁。

4.1.3 截面查询

查询某一时间范围内多个设备的运行指标截面序列,形成全局视图。DolphinDB 的 pivot by 子句支持将长格式数据转换为宽格式,便于横向对比。

select
    temperature,
    pressure,
    humidity
from
    loadTable("dfs://iot_sensor", "sensor_data")
where
    ts between 2025.01.01T00:00:00.000
    and 2025.01.01T00:01:00.000
    and device_id between "device_0001"
    and "device_0100" 
pivot by ts,device_id

pivot by 会按 ts 保留时间维度,并将 device_id 展开为列,把长表转换为适合横向对比的宽表。该结果便于在 Grafana 等可视化工具中展示多设备同一时间段内的截面状态,也能简化时序数据的透视分析过程。当 pivot byexec 配合使用时,查询结果会以矩阵形式返回,适合用于批量计算多设备之间的相关性或相似度等分析指标。

注意:设备截面分析可能生成列数较多、数据量较大的宽表。建议在 where 条件中同时限制时间范围和设备范围,只查询关注的设备集合,避免一次性返回过多截面数据,导致 VS Code、Web 等前端工具渲染耗时过长。

4.2 聚合分析

物联网原始数据通常以秒级或亚秒级频率持续产生,直接基于明细数据进行查询和分析往往会带来较高的存储扫描与计算开销。聚合分析通过降采样和统计计算,在保留整体趋势与关键特征的同时显著减少数据量,从而更高效地支撑后续应用场景,例如分时曲线监控、设备数据上报统计以及机器学习模型训练等。

4.2.1 时间窗口聚合

降采样将高频率原始数据聚合为低频率的统计值,常用于生成趋势图或报表。DolphinDB 的 bar 函数将时间戳按指定间隔对齐,配合聚合函数实现降采样。

以下示例计算各设备每 10 分钟的平均温度和最大湿度:

select
    device_id,
    ts_bar,
    avg(temperature) as temp_avg,
    max(humidity) as humidity_max
from
    loadTable("dfs://iot_sensor", "sensor_data")
where
    ts between 2025.01.01T00:00:00.000
    and 2025.01.01T00:59:59.999
group by
    device_id,
    bar(ts, 10m) as ts_bar

bar(ts, 10m) 将时间戳按 10 分钟间隔取整,生成时间桶。由于示例数据为 1 小时长度,使用 10 分钟窗口可生成 6 个数据点便于观察。该模式适用于生成分钟级、小时级、日级的趋势报表。

bar 函数外,DolphinDB 还提供 interval 函数用于时间窗口聚合。interval 函数在物联网场景中更为实用,支持对空值进行插值填充,可处理设备上报间隔不稳定或数据缺失的情况。例如,当设备因网络中断导致某几分钟无数据时,interval 可自动填充空值或进行线性插值,保证聚合结果的连续性。

select device_id,
       ts,
       avg(temperature) as temp_avg,
       max(humidity) as humidity_max
from 
      loadTable("dfs://iot_sensor", "sensor_data")
where 
      date(ts) = 2025.01.01
group by 
      device_id, 
      interval(ts, 1m, "prev") as ts

关于 interval 函数的详细用法(包括填充模式、边界处理等),可参考 DolphinDB 官方文档:interval 函数

4.2.2 全量聚合统计

全量聚合用于统计分析整体数据的统计特征。以下示例统计指定时间范围内所有设备的平均温度、最大压力和最小电压:

select
    nunique(device_id) as device_count,
    avg(temperature) as avg_temperature,
    max(pressure) as max_pressure,
    min(voltage) as min_voltage,
    avg(current) as avg_current
from
    loadTable("dfs://iot_sensor", "sensor_data")
where
    date(ts) = 2025.01.01

nunique 是 count(distinct …) 的等价替代语法。

4.2.3 滑动窗口计算

滑动窗口计算(Moving Window Calculation)用于分析指标的局部趋势,如计算移动平均值、移动标准差等。

传统 SQL 使用 OVER 子句实现滑动窗口:

select
    device_id,
    ts,
    temperature,
    avg(temperature) over (
        partition by device_id
        order by
            ts rows between 10 preceding
            and current row
    ) as moving_avg
from
    loadTable("dfs://iot_sensor", "sensor_data")
where
    date(ts) = 2025.01.01
limit
    100

DolphinDB 提供内置的移动窗口函数,语法更为简洁:

select
    device_id,
    ts,
    temperature,
    mavg(temperature, 10) as moving_avg_10
from
    loadTable("dfs://iot_sensor", "sensor_data")
where
    date(ts) = 2025.01.01 
context by 
    device_id csort ts
limit
    100

mavg(temperature, 10) 用于计算当前记录及其前序记录构成的 10 点移动平均值,帮助平滑短期波动、观察温度变化趋势。DolphinDB 内置了丰富的移动窗口函数,如 mavgmsummmaxmminmstd 等,可直接在 context by 分组后的有序序列上使用,并支持不同指标配置不同窗口大小,无需编写冗长的 OVER 子句,语法更简洁,适合物联网时序指标的趋势分析与异常波动识别。

对窗口计算感兴趣的读者,可以进一步阅读 窗口计算

4.3 状态分析

工业场景中,设备状态监控和安全阈值告警是运维系统的核心功能。本节介绍如何使用 SQL 实现设备开停机统计、越限告警和空载检测等常见需求。

4.3.1 累计开停机时长统计

设备累计开停机时长是设备全寿命周期管理的关键指标。假设设备状态字段 status 的值为 1 表示开机,0 表示关机。

计算逻辑如下:使用 bar(ts, 10m) 按 10 分钟分组;按设备状态进一步分组;对同一时间状态内的差值“按设备分组”求和得到累计时长。

select
    device_id,
    status,
    ts,
    sum(duration_ms) / 60000.0 as duration_minutes
from
    (
        select
            device_id,
            ts,
            status,
            iif(isNull(next(ts)), 0, next(ts) - ts) as duration_ms
        from
            loadTable("dfs://iot_sensor", "sensor_data")
        where
            device_id = "device_0001"
        and 
            date(ts) = 2025.01.01 
        context by 
            device_id 
        csort 
            ts
    )
group by
    device_id,
    status,
    bar(ts, 10m) as ts

4.3.2 开停机次数统计

统计指定时间段内设备的开关机次数,用于评估设备使用频率。

计算逻辑如下:使用 deltas 函数计算相邻状态值的差值;差值为 1 表示由关到开,差值为 -1 表示由开到关;统计差值为 1 和 -1 的记录数即为开关机次数。

select
    device_id,
    sum(iif(deltas(status) == 1, 1, 0)) as start_count,
    sum(iif(deltas(status) == -1, 1, 0)) as stop_count
from
    loadTable("dfs://iot_sensor", "sensor_data")
where
    device_id = 'device_0001'
and 
    date(ts) = 2025.01.01
group by
    bar(ts, 10m),
    device_id

iif 函数实现条件判断,仅对状态变化的记录进行计数。

4.3.3 最大持续运行时间

最大持续运行时间反映设备的稳定运行能力,是预测性维护的重要参考指标。

计算逻辑如下:使用 deltas 计算时间戳差值;使用 segment 函数按状态分组,将连续相同状态的记录分到同一组;筛选状态为开机(1)的组,计算各组时长;取最大值即为最大持续运行时间。

select
    device_id,
    max(running_duration) / 1000 as max_running_sec
from
    (
        select
            device_id,
            sum(iif(status == 1, deltas(ts), 0)) as running_duration
        from
            loadTable("dfs://iot_sensor", "sensor_data")
        where
            date(ts) = 2025.01.01
        group by
            device_id,
            segment(status)
    )
group by
    device_id

segment(status) 函数在状态值变化时产生新的分组标识,将连续开机或关机的记录归为一组。

4.3.4 越限告警查询

越限告警用于检测传感器数据是否超出安全阈值。以下示例查询温度超过 80 度, 或小于 20 度的异常记录:

select
    device_id,
    ts,
    temperature
from
    loadTable("dfs://iot_sensor", "sensor_data")
where
    date(ts) = 2025.01.01
    and (
        temperature > 80
        or temperature < 20
    )

对于变化率告警(如电流突变检测),可使用 ratios 计算相邻值的比例:

select
    *
from
    (
        select
            ts,
            device_id,
            current,
            ratios(current) - 1 as changePct
        from
            loadTable("dfs://iot_sensor", "sensor_data")
        where
            device_id = 'device_0001'
            and date(ts) = 2025.01.01
    ) t
where
    abs(t.changePct) > 1

上述查询检测变化率超过 100% 的异常数据。

4.3.5 空载时间统计

设备空载指设备运行但无工作负载或负载极低的状态,通常通过电流阈值判断(如电流 < 20 表示空载)。

计算逻辑如下:筛选电流小于阈值的记录;使用 deltas 计算空载状态的持续时长;按时间窗口汇总。

select
    device_id,
    sum(deltas(ts)) / 60000 as no_load_minutes
from
    loadTable("dfs://iot_sensor", "sensor_data")
where
    device_id = 'device_0001'
    and current < 20
    and date(ts) = 2025.01.01
group by
    bar(ts, 10m),
    device_id

4.4 多表关联

物联网数据分析常需要关联设备元数据或多传感器数据,以获取更完整的上下文信息。

设备元数据(如设备名称、位置、型号)通常存储在单独的表中,需通过关联查询与传感器数据结合。假设存在设备元数据表 device_metadata 包含设备位置信息:

/ / 创建设备元数据表 ( 内存表示例 )
device_list = "device_" + string(1..10000).lpad(4, "0")
device_metadata = table(
    device_list as id,
    take(["Factory_A", "Factory_B", "Factory_C"], 10000) as location,
    take(["Type_A", "Type_B"], 10000) as device_type
) 
/ / 关联查询: 获取位置为 Factory_A 的相关设备数据
select
    d.location,
    s.device_id,
    s.temperature,
    s.ts
from
    loadTable("dfs://iot_sensor", "sensor_data") s
    inner join device_metadata d on s.device_id = d.id
where
    d.location = 'Factory_A'
    and date(s.ts) = 2025.01.01

除了 inner join 之外,DolphinDB 也支持其常见的左/右/外/半/全连接,以及扩展的 window join、asof join、prefix join 关联语法。

4.5 常见问题

Q1:时间类型不匹配怎么办?

DolphinDB 支持多种时间类型(DATE、DATETIME、TIMESTAMP、MINUTE 等),不同类型之间不可直接比较。当 sql 执行提示“Temporal data comparison should have the same data type” 错误时, 需要进行类型转换。

-- 方式1:使用时间转换函数
select
    *
from
    loadTable("dfs://iot_sensor", "sensor_data")
where
    date(ts) = 2025.01.01
limit
    10
-- 方式2:使用对应类型完整时间戳字符串
select
    *
from
    loadTable("dfs://iot_sensor", "sensor_data")
where
    ts 
    between 2025.01.01T00:00:00.000
    and 2025.01.01T23:59:59.999
limit
    10

Q2:怎么优化查询性能?

SQL 的查询性能受两个核心因素影响

分区剪枝:确保查询条件包含分区列,使系统能过滤不相关分区。对于按时间和设备哈希分区的表,查询应同时指定时间范围和设备标识。

索引利用:TSDB 引擎的查询性能依赖于排序列索引。查询条件应包含排序列的前缀,以利用索引定位数据。

  • 分区剪枝

在 SQL 中加上 [HINT_EXPLAIN] 可以获得执行计划,从其中的属性"partitions" 可以了解到查询引擎所扫描的分区数。

select
    [HINT_EXPLAIN] *
from
    loadTable("dfs://iot_sensor", "sensor_data")
where
    date(ts) = 2025.01.01

从上述 SQL 执行计划中的 "partitions" 可以看到是访问了 5 个分区,即 2025.01.01 下的 5 个 hash 分区,20250101/Key2/z。where date(ts) = 2025.01.01 条件使得查询引擎只检索了 2025.01.01 分区,而非全部分区( 2025.01.01 至 2025.01.03 )。

{
    "measurement": "microsecond",
    "explain": {
        "from": {
            "cost": 16
        },
        "map": {
            "partitions": {
                "local": 5,
                "remote": 0
            },
            "cost": 90424,
            "partitionRoute": {
                "local8848": 5
            },
            "detail": {
                "most": {
                    "sql": "select [253967] device_id,ts,temperature,pressure,humidity,voltage,current,status from sensor_data [partition = /iot_sensor/20250101/Key9/3cj]",
                    "explain": {
                        "rows": 35700,
                        "cost": 78764
                    }
                },
          //skip other information
    }}}
}
  • 索引利用

当 where 条件中带上索引列(device_id)时,查询引擎就可以根据 TSDB 引擎的索引进行快速查找,显著减少扫描 block 数量。

select [HINT_EXPLAIN] * from loadTable("dfs://iot_sensor", "sensor_data")
where date(ts) = 2025.01.01 and device_id = "device_0001"

执行计划中可以看到"TSDBIndexPrefiltering"属性,matchedWhereConditions=1 代表引擎根据 where 条件进行索引查找。

"where": {
            "TSDBIndexPrefiltering": {
                "blocksToBeScanned": 7,
                "matchedWhereConditions": 1
            },
            "rows": 300,
            "cost": 4105
}

当执行如下的 sql 时,则无法有效利用 sortKey, 需要避免。

select [HINT_EXPLAIN] * from loadTable("dfs://iot_sensor", "sensor_data")
where date(ts) = 2025.01.01 and substr(device_id, 7, 4) = "0001"

其执行计划将变为:

"where": {
    "TSDBIndexPrefiltering": {
        "blocksToBeScanned": 651,
        "matchedWhereConditions": 0
    },
    "rows": 0,
    "cost": 12621
}

其执行时间将会由 5ms 增加至 50ms,在实际的生产场景中,性能差距(时延、IO、内存消耗)会更大。

4.6 阅读指引

本章系统介绍了物联网场景下常用的数据查询与分析方法,涵盖时序查询、聚合统计、状态监控、关联分析、趋势分析以及查询性能优化等典型内容。如需进一步学习 SQL 编写技巧和性能分析方法,可参考以下文档:

5. 实时计算

针对工业物联网中设备指标秒级聚合、异常即时发现及控制指令触发等实时计算需求,传统轮询数据库或Kafka+Flink 方案存在延迟高、压力大或技术栈复杂等问题。DolphinDB 内置统一流计算框架,以流表管理数据,提供时序聚合、异常检测、规则及会话窗口等多种即开即用的声明式引擎,支持流批一体与引擎级联,无需编写复杂代码即可实现高效实时处理。

本章将介绍如何基于 DolphinDB 流计算引擎构建一套完整的实时监控体系,涵盖数据质量处理、实时告警以及在线推理,最终通过一个综合案例串联全流程。

5.1 数据清洗

原始传感器数据往往存在噪声、缺失、单位不一致或采样频率过高的问题,直接使用会降低监控的准确性。因此,在进入业务计算之前,需要对流数据进行实时清洗与规整。

本示例以电压、电流数据为例,演示 DolphinDB 如何在流数据订阅过程中完成数据质量控制:将 voltage <= 122current 为空的数据过滤掉,仅将有效数据写入时序聚合引擎进行计算。

整体流程如下:

  1. 创建原始流表 electricity 和聚合结果表 outputTable

  2. 自定义过滤函数 append_after_filtering,对订阅数据进行清洗;

  3. 通过部分应用将过滤函数与时序聚合引擎绑定;

  4. 订阅原始流表,实时接收数据并触发过滤逻辑;

  5. 模拟写入传感器数据,查看原始数据和聚合结果。

其中,select * from electricity 可查看原始流数据,select * from outputTable 可查看过滤后参与计算得到的聚合结果。

share streamTable(
    1000:0,
    `timev`voltage`current,
    [TIMESTAMP, DOUBLE, DOUBLE]
) as electricity

outputTable = table(
    10000:0,
    `timev`avgVoltage`avgCurrent,
    [TIMESTAMP, DOUBLE, DOUBLE]
)

//自定义数据处理过程,过滤 voltage<=122 或 current=NULL 的无效数据。
def append_after_filtering(inputTable, msg){
	t = select * from msg where voltage>122, isValid(current)
	if(size(t)>0){
		insert into inputTable values(
                    t.timev,
                    t.voltage,
                    t.current)		
	}
}
electricityAggregator = createTimeSeriesEngine(
    name = "electricityAggregator",
    windowSize = 6,
    step = 3,
    metrics = < [avg(voltage), avg(current)] >,
    dummyTable = electricity,
    outputTable = outputTable,
    timeColumn = `timev, 
    garbageSize = 2000
)
subscribeTable(
    tableName = "electricity",
    actionName = "avgElectricity",
    offset = 0,
    handler = append_after_filtering{electricityAggregator},
    msgAsTable = true
)

//模拟产生数据
def writeData(t, n){
    timev = 2018.10.08T01:01:01.001 + timestamp(1..n)
    voltage = 120+1..n * 1.0
    current = take([1, NULL, 2]*0.1, n)
    insert into t values(timev, voltage, current);
}
writeData(electricity, 10)

5.2 实时告警

实时告警是物联网监控中的核心功能,常见需求包括阈值越限、状态跳变、变化率异常和设备离线等。DolphinDB 可以结合流数据表、异常检测引擎和自定义消息处理函数,实现实时告警规则的在线计算。

本示例演示温度超限告警逻辑:通过 createAnomalyDetectionEngine() 创建异常检测引擎。规则为:每 30 秒计算一次最近 3 分钟窗口内的数据,如果同一传感器温度超过 40℃ 的次数大于 2 次,且超过 30℃ 的次数大于 3 次,则触发温度异常告警。

// 定义输入输出流数据表
st = streamTable(
    1000000:0,
    ["device_id", "ts", "temperature"],
    [INT, DATETIME, FLOAT]
)
enableTableShareAndPersistence(
    table = st,
    tableName = `sensor,
    asynWrite = false,
    compress = true,
    cacheSize = 1000000
)
share streamTable(
    1000:0,
    ["time", "device_id", "anomalyType", "anomalyString"],
    [DATETIME, INT, INT, SYMBOL]
) as warningTable
// 创建异常检测引擎,实现传感器温度异常报警的功能
engine = createAnomalyDetectionEngine(
    name = "engine1",
    metrics = <[
        sum(temperature > 40) > 2 
        && 
        sum(temperature > 30) > 3
    ]>,
    dummyTable = sensor,
    outputTable = warningTable,
    timeColumn = `ts,
    keyColumn = `device_id,
    windowSize = 180,
    step = 30
)
subscribeTable(
    tableName = "sensor",
    actionName = "sensorAnomalyDetection",
    offset = 0,
    handler = append!{engine},
    msgAsTable = true
)

上述代码中,windowSize = 180 表示窗口长度为 180 秒,即 3 分钟;step = 30 表示每 30 秒计算一次。异常检测引擎会按 device_id 分组,对每个传感器分别计算温度超限规则,并将满足条件的结果写入 warningTable

createAnomalyDetectionEngine() 支持在异常指标中使用聚合函数,也可以组合聚合结果、常量、列和非聚合函数来构建复杂告警规则。当异常指标中包含聚合函数时,必须显式指定窗口长度 windowSize 和计算步长 step。因此,该引擎适合处理阈值越限、窗口统计异常、变化率异常等实时告警场景。更多示例请参考:传感器数据异常检测

5.3 实时推理

本节为进阶内容,入门阶段可先跳过。

在工业物联网预测性维护场景中,DolphinDB 流计算框架支持集成 KNN 回归模型实现实时推理:以风力发电机组发电量异常预警为例,系统通过流表接收高频监测数据,利用时序聚合引擎每 10 秒生成设备级特征均值;每次新聚合数据到达时,读取最近 60 秒内的历史聚合结果(约 600 条记录)训练全局 KNN 模型,并对当前数据进行发电量预测;将预测值与真实值比对,若相对偏差超过 0.1% 则写入告警表。该方案在 100 台风机的规模下,可应对 10 万条/秒的高吞吐原始数据,通过流内持续更新模型和低延迟处理链路,满足实时异常预警需求。

首先创建原始数据流表 dataTable ,聚合结果流表 aggrTable,预测缓存表 predCache,告警输出表 alertTable。

// 定义列名与类型
orgColNames = [
    "time", 
    "deviceCode", 
    "wind", 
    "humidity", 
    "air_pressure", 
    "temperature", 
    "life", 
    "propertyValue"
]
    
orgColTypes = [
    TIMESTAMP,
    INT,
    DOUBLE,
    DOUBLE,
    DOUBLE,
    DOUBLE,
    INT,
    INT
]

// 创建共享流表:原始数据表
share streamTable(
    1000000:0,
    orgColNames,
    orgColTypes
) as dataTable

// 创建聚合结果表(含窗口结束时间)
aggrColNames = [
    "timeWindowEnd", 
    "deviceCode", 
    "avgWind", 
    "avgHumidity", 
    "avgPressure", 
    "avgTemp", 
    "avgLife", 
    "avgProperty"
]
aggrColTypes = [
    TIMESTAMP,
    INT,
    DOUBLE,
    DOUBLE,
    DOUBLE,
    DOUBLE,
    DOUBLE,
    DOUBLE
]

share streamTable(
    1000000:0,
    aggrColNames,
    aggrColTypes
) as aggrTable

// 创建预测缓存表(存储预测值与真实值)
share streamTable(
    1000000:0,
    ["timeWindowEnd", "deviceCode", "predictedValue", "avgProperty"],
    [TIMESTAMP, INT, DOUBLE, DOUBLE]
) as predCache

// 创建告警输出表
share streamTable(
    1000000:0,
    ["timeWindowEnd", "deviceCode", "anomalyRate"],
    [TIMESTAMP, INT, DOUBLE]
) as alertTable

接下来创建时序聚合引擎,按设备分组,每 10 秒计算一次特征均值。由于输入时间列 time 的类型为 TIMESTAMP,时间精度为毫秒,因此 windowSize = 10000step = 10000 表示窗口长度和滑动步长均为 10 秒。

// 创建时序聚合引擎:每 10 秒滑动窗口计算各设备特征均值
aggEngine = createTimeSeriesEngine(
    name = "aggEngine",
    windowSize = 10000,
    step = 10000,
    metrics = <[
        avg(wind),
        avg(humidity),
        avg(air_pressure),
        avg(temperature),
        avg(life),
        avg(propertyValue)
    ]>,
    dummyTable = dataTable,
    outputTable = aggrTable,
    timeColumn = `time,
    keyColumn = `deviceCode,
    useSystemTime = false
)
// 订阅原始数据,输入聚合引擎
subscribeTable(
    tableName = "dataTable",
    actionName = "writeToAggEngine",
    offset = 0,
    handler = append!{aggEngine},
    msgAsTable = true
)

上述代码中,keyColumn = `deviceCode 表示按设备分别进行窗口聚合;useSystemTime = false 表示按照数据本身的时间戳进行窗口切分,而不是按照系统接收时间进行聚合。

为了确保实时数据到达后有足够的历史聚合数据用于训练,需要先导入一段历史数据进行模型预热。

生产环境中,这部分历史数据通常来自历史库或已有数据文件。为便于演示,本示例通过函数生成过去 60 秒的模拟数据,并写入原始流表。写入后,前面创建的时序聚合引擎会自动生成历史聚合结果。

// ----- 导入历史数据用于初始训练 -----
def generateHistoricalData(
        startTime, 
        seconds, 
        devices_number, 
        rate) {
    // 生成从 startTime 开始的连续 seconds 秒数据,每秒 rate 条/设备
    n_total = devices_number * rate * seconds
    // 生成时间戳:每秒内均匀分布
    timeStamps = array(TIMESTAMP, 0, n_total)
    deviceCodes = array(INT, 0, n_total)
    for(i in 0..(seconds - 1)) {
        base = startTime + i*1000
        tmp_ts = take(base + (0..999), devices_number * rate)
        timeStamps.append!(tmp_ts)
        tmp_code = take(1..devices_number, devices_number * rate)
        deviceCodes.append!(tmp_code)
    }
    n = timeStamps.size()
    // 生成特征(正态分布模拟)
    x1 = randNormal(25.0, 2.0, n)      
    x2 = randNormal(55.0, 5.0, n)      
    x3 = randNormal(1.01325, 0.00001, n) 
    x4 = randNormal(75.0, 5.0, n)      
    x5 = int(randNormal(10.0, 3.0, n))  
    // 生成线性系数(加入随机扰动)
    b1 = randNormal(0.4, 0.05, n)
    b2 = randNormal(0.3, 0.05, n)
    b3 = randNormal(0.2, 0.05, n)
    b4 = randNormal(0.09, 0.05, n)
    b5 = randNormal(0.01, 0.001, n)
    bias = randNormal(5.0, 1.0, n)
    propertyValue = int(
      b1 * x1 * 10 +
      b2 * x2 * 2 +
      b3 * x3 * 1000 +
      b4 * x4 +
      b5 * x5 +
      bias
  )
    // 构造表并返回
    return table(timeStamps as time, 
                deviceCodes as deviceCode, 
                x1 as wind, 
                x2 as humidity, 
                x3 as air_pressure, 
                x4 as temperature, 
                x5 as life, 
                propertyValue as propertyValue
                )
}
// 为便于本地演示,示例默认使用 10 台设备
// 生产场景可将 demoDeviceNum 调整为 100 或更高
demoDeviceNum   = 10
sampleRate      = 1000
historySeconds  = 60
// 以当前时间为基准生成过去 60 秒的历史数据
baseTime        = timestamp(now())
historicalStart = baseTime - historySeconds * 1000
historicalData = generateHistoricalData(
    historicalStart,
    historySeconds,
    demoDeviceNum,
    sampleRate
)
// 将历史数据写入原始流表,聚合引擎会自动生成聚合结果
dataTable.append!(historicalData)
print("Historical data appended, waiting for aggregation...")
sleep(3000)

本地演示中,为降低计算压力,默认使用 10 台设备。若模拟 100 台设备、每台设备每秒 1000 条数据,则原始写入速率为 10 万条/秒,对机器资源要求较高。

接着,定义在线训练与预测函数predictStream

def predictStream(msg) {
    endTime = max(msg.timeWindowEnd)
    startTime = endTime - 60 * 1000
    hist_data = select *
        from aggrTable
        where timeWindowEnd between startTime: endTime
    features = select
            avgWind,
            avgHumidity,
            avgPressure,
            avgTemp,
            avgLife
        from hist_data
    labels = hist_data.avgProperty
    // 训练 KNN 回归模型(k = 200)
    model = knn(
        labels,
        features,
        "regressor",
        200
    )
    predict_data = select *
        from msg
    predict_features = select
            avgWind,
            avgHumidity,
            avgPressure,
            avgTemp,
            avgLife
        from predict_data
    predicted = predict(
        model,
        predict_features
    )
    cache_data = table(
        predict_data.timeWindowEnd as timeWindowEnd,
        predict_data.deviceCode as deviceCode,
        predicted as predictedValue,
        predict_data.avgProperty as avgProperty
    )
    predCache.append!(cache_data)
}

注意:上述函数中的 endTime 固定为历史时间点,实际部署时应根据数据动态调整,例如使用 now()max(timeWindowEnd)。此处为简化示例,采用固定时间戳。

定义检测函数 detectAnomaly 函数,订阅预测缓存表 predCache。函数对每条预测结果计算预测值与真实值之间的相对偏差,若偏差超过阈值,则写入告警表 alertTable

相对偏差计算公式为:

相对偏差 = abs(真实值 - 预测值) / abs(真实值)

示例中阈值设置为 0.001,即 0.1%。

// 异常检测函数:将真实值与缓存预测值比对
def detectAnomaly(msg) {
    tmp = select * from msg
    timeWindowEnd_li = array(TIMESTAMP, 0)
    deviceCode_li = array(INT, 0)
    anomalyRate_li = array(DOUBLE, 0)
    for (i in 0:tmp.size()) {
        predicted = tmp.predictedValue[i]
        real = tmp.avgProperty[i]
        if (real != 0) {
            diffRatio = abs(real - predicted) / real
            if (diffRatio > 0.001) {
                timeWindowEnd_li.append!(tmp.timeWindowEnd[i])
                deviceCode_li.append!(tmp.deviceCode[i])
                anomalyRate_li.append!(diffRatio)
            }
        }
    }
    if (size(timeWindowEnd_li) > 0) {
        alertTable.append!(
            table(
                timeWindowEnd_li as timeWindowEnd,
                deviceCode_li as deviceCode,
                anomalyRate_li as anomalyRate
            )
        )
    }
}

然后启动流订阅,设置两个订阅:

  • 订阅 aggrTable(从最新数据开始,offset=-1),触发在线训练与预测;

  • 订阅 predCache,触发异常检测。

// 订阅聚合表
// 从最新数据开始,offset = -1 表示只消费之后的新数据
subscribeTable(
    tableName = "aggrTable",
    actionName = "predict",
    offset = -1,
    handler = predictStream,
    msgAsTable = true
)
subscribeTable(
    tableName = "predCache",
    actionName = "detectAnomaly",
    offset = -1,
    handler = detectAnomaly,
    msgAsTable = true
)

最后,启动一个后台作业持续生成实时数据,写入 dataTable。本地演示中,默认模拟 10 台设备、每台设备每秒 1000 条数据,持续写入 360 秒。

// ----- 6. 启动实时数据模拟(后台持续运行)-----
def send_realtime_data(
    begintime,
    durationSeconds,
    devices_number,
    rate,
    mutable dest
) {
    btime = timestamp(begintime)
    endtime = btime + durationSeconds * 1000
    do {
        n = devices_number * rate
        time = sort(take(btime + (0..999), n))
        deviceCode = take(1..devices_number, n)
        x1 = randNormal(25.0, 2.0, n)
        x2 = randNormal(55.0, 5.0, n)
        x3 = randNormal(1.01325, 0.00001, n)
        x4 = randNormal(75.0, 5.0, n)
        x5 = int(randNormal(10.0, 3.0, n))
        b1 = randNormal(0.4, 0.05, n)
        b2 = randNormal(0.3, 0.05, n)
        b3 = randNormal(0.2, 0.05, n)
        b4 = randNormal(0.09, 0.05, n)
        b5 = randNormal(0.01, 0.001, n)
        bias = randNormal(5.0, 1.0, n)
        propertyValue = int(
            b1 * x1 * 10 +
            b2 * x2 * 2 +
            b3 * x3 * 1000 +
            b4 * x4 +
            b5 * x5 +
            bias
        )
        table_ps = table(
            time,
            deviceCode,
            x1,
            x2,
            x3,
            x4,
            x5,
            propertyValue
        )
        dest.append!(table_ps)
        btime = btime + 1000
        etime = timestamp(now())
        timediff = btime - etime
        if (timediff > 0) {
            sleep(timediff)
        }
    } while (btime < endtime)
}
// 模拟从现在开始,持续写入 360 秒
submitJob(
    "simulate_realtime",
    "Simulate real-time data",
    send_realtime_data,
    now(),
    360,
    10,
    1000,
    dataTable
)

5.4 常见问题

Q1:如何监控流计算消费及时延情况?

通过 getStreamingStat() 函数可以监控消息的消费情况,包括是否有错误消息,消息堆积等。在引擎中可以通过设置参数 outputElapsedMicroseconds= true 来记录引擎的计算耗时。

Q2:为什么时序聚合引擎没有按预期窗口输出结果?

引擎的 windowSize 参数的单位与 timeColumn 字段相匹配,例如当 timeColumn 字段是 DATETIME 类型时,windowSize = 60000 代表着 60000 秒 而不是 60 秒。请根据数据情况设置 windowSize、step。

Q3:如何处理乱序数据?

在默认的情况下,乱序数据是丢弃的。也可以通过 acceptedDelay 参数来设置延迟容忍,在延迟容忍的范围内乱序数据达到也可以被计算。在流计算中,准确性和时延始终是需要权衡的,请根据实际场景设置一个合理的值。

5.5 阅读指引

6. 第三方可视化平台对接实时数据监控

实时计算的价值不仅在于快速产出结果,更在于将结果以直观、可交互的方式呈现,帮助运维人员实时掌握设备运行状态并快速定位异常。Grafana 是常用的开源可视化与监控平台,配合 DolphinDB Grafana 数据源插件,可以直接查询流表、实时计算结果表或物化视图中的数据,构建支持自动刷新、趋势分析和告警联动的监控仪表盘。本章将介绍 Grafana 的安装、DolphinDB 数据源插件的部署,以及在 Grafana 中配置数据源和创建实时监控面板的基本流程。

接入流程概览

  • 安装 Grafana 与 DolphinDB 数据源插件:完成 Grafana 安装,并将 DolphinDB 数据源插件部署到 Grafana 插件目录。

    1. Grafana 官网下载并安装最新的开源版本(OSS)。安装完成后,确认 Grafana 服务可以正常启动。

    2. 下载并部署 DolphinDB 数据源插件。可在 GitHub Releases 页面获取最新插件压缩包,例如 dolphindb-datasource.v2.0.900.zip。下载后,将压缩包中的 dolphindb-datasource 文件夹解压到 Grafana 插件目录。

      常见插件目录如下:

      Windows:<grafana 安装目录>\data\plugins\

      Linux:/var/lib/grafana/plugins/

      注:如果不存在 plugins 目录,可手动创建。

    3. 修改 Grafana 配置文件,允许加载未签名的 dolphindb-datasource 插件。在 [plugins] 配置段中取消注释 allow_loading_unsigned_plugins,并将其设置为 dolphindb-datasource

      allow_loading_unsigned_plugins = dolphindb-datasource

      注:每次修改配置项后,需重启 Grafana。

    4. 重启 Grafana 服务,使插件和配置变更生效。Windows 环境可在“任务管理器 > 服务”中找到 Grafana 服务并重启;Linux 环境可根据部署方式使用相应的服务管理命令重启。

  • 在 Grafana 中配置 DolphinDB 数据源并创建监控面板:完成登录、数据源连接测试、Panel 查询配置和自动刷新设置。

    1. 打开 http://localhost:3000。初始登入名以及密码均为 admin。打开 http://localhost:3000/datasources,或点击左侧导航的 Configuration > Data sources 添加数据源,搜索并选择 “dolphindb”,配置数据源后点 Save & Test 以保存数据源。

    2. 打开或新建 Dashboard,编辑或新建 Panel,在 Panel 的 Data source 属性中选择上一步添加的数据源。输入 SQL 查询流表,例如 select ts, temp from clean_temp where deviceId='device_01'

    3. 设置刷新频率(如 5s),即可实现实时曲线展示。

    Grafana 支持基于查询结果配置告警规则,并可通过邮件、钉钉等渠道推送通知,便于及时发现设备异常并联动运维处理。DolphinDB 自带的 Web 管理界面也可用于基础表数据查看与简单监控;如果需要构建多维度、可交互、可自动刷新的复杂仪表盘,建议优先使用 Grafana。完整配置步骤请参考 DolphinDB Grafana DataSource Plugin,更多实践案例可参考 实现设备指标的采集监控和展示

7. 最佳实践

本案例以旋转机械振动监测为背景,演示从高频振动数据接入、实时指标计算、异常告警到 Grafana 可视化展示的完整工程链路。场景中模拟 16 台振动传感器,每台传感器以 1 ms 间隔采集加速度数据,并各写入 1000 条样例记录;系统需实时计算每秒振动有效值(RMS),判断 RMS 是否超过设定阈值,并将监控结果接入 Grafana 进行实时展示。

7.1 流表创建与数据接入

首先定义两张流表:vibration_raw用于接收原始数据,vibration_rms用于存储每秒有效值计算结果,vibration_alerts用于存储告警记录。

// 定义原始数据流表:时间戳、设备 ID、加速度值
share streamTable(
    1000000:0,
    ["timestamp", "deviceID", "accel"],
    [TIMESTAMP, SYMBOL, DOUBLE]
) as vibration_raw
// 定义 RMS 结果流表:窗口结束时间、设备 ID、RMS 值
share streamTable(
    1000000:0,
    ["timestamp", "deviceID", "rms"],
    [TIMESTAMP, SYMBOL, DOUBLE]
) as vibration_rms
// 定义告警流表:时间、设备 ID、触发值、规则 ID
share streamTable(
    1000000:0,
    ["timestamp", "deviceID", "value", "ruleID"],
    [TIMESTAMP, SYMBOL, DOUBLE, SYMBOL]
) as vibration_alerts

实际部署时,可通过 MQTT、Kafka 或 SDK 将传感器数据写入vibration_raw。本案例使用脚本模拟数据写入:

def writeSimulatedData(){
    deviceIDs = "deviceID_" + string(1..16)
    ts = now()
    for (id in deviceIDs) {
        data = table(take(ts + 0..999, 1000) as timestamp,
                take(id, 1000) as deviceID,  
                randNormal(0, 1, 1000) as signalnose) // 模拟振动加速度
        vibration_raw.append!(data)
    }
}
submitJob(
  "simulate_vibration", 
  "simulate vibration data", 
  writeSimulatedData
)

7.2 创建时序聚合引擎计算 RMS

RMS(有效值)的计算公式为:

3. 图 7-1 RMS 计算公式

在流式计算中,可通过时序聚合引擎每秒滑动一次窗口(窗口大小 1 秒)计算每个设备的 RMS。

// 创建时序聚合引擎:窗口 100 毫秒,滑动步长 10 毫秒,计算 accel 的 RMS
rms_engine = createTimeSeriesAggregator(
              name="rms_engine1", 
              windowSize=100, 
              step=10, 
              metrics=<[rms(accel)]>, 
              dummyTable=vibration_raw, 
              outputTable=vibration_rms, 
              timeColumn=`timestamp, 
              keyColumn=`deviceID,
              useSystemTime = true
              )
// 订阅原始流表,将数据输入引擎
subscribeTable(
              tableName="vibration_raw", 
              actionName="calc_rms", 
              offset=0, 
              handler=append!{rms_engine}, 
              msgAsTable=true
              );

设计说明

  • windowSize=100表示窗口长度为 100 毫秒;

  • step=10表示每 10 毫秒触发一次计算;

  • metrics中使用预定义好的函数rms(accel)计算 RMS 值;

  • useSystemTime = true表示使用系统时间作为时间基准,即毫秒,而非时间列 timestamp。

7.3 创建异常检测引擎实现超限告警

设振动烈度预警值为 1.0。需要检测 RMS 是否超过这个阈值,并将告警记录输出。

// 创建异常检测引擎,检测 rms 值
alert_engine = createAnomalyDetectionEngine(name="rms_alert1", 
                                             metrics=<[rms > 1.0]>, 
                                             dummyTable=vibration_rms, 
                                             outputTable=vibration_alerts, 
                                             keyColumn=`deviceID, 
                                             timeColumn=`timestamp
                                         )
// 订阅 RMS 结果流表,将数据输入异常检测引擎
subscribeTable(
    tableName="vibration_rms", 
    actionName="check_alert", 
    offset=0, 
    handler=append!{alert_engine}, 
    msgAsTable=true
    )

7.4 Grafana 实时监控展示

DolphinDB 提供针对 Grafana 的数据源插件,支持直接查询流表。关于 Grafana 插件的内容可参考 DolphinDB Grafana DataSource Plugin

完成数据源配置后,可在 Grafana 的 Query 面板中输入以下 SQL,实时查看设备 deviceID_1 的 RMS 指标变化:

select * from vibration_rms where deviceID='deviceID_1'

7.5 总结与扩展

本案例串联了振动数据接入、实时 RMS 计算、阈值告警和 Grafana 可视化展示等关键环节,代码可直接复制运行。实际应用中,可根据设备规模、采样频率、窗口参数和告警阈值进行调整。若需进一步了解随机振动信号分析、频域特征提取及更完整的工程实现,可参考 随机振动信号分析解决方案