TDengine 流计算详解与实战案例

1. 流计算概述

流计算是一种实时处理动态数据流的技术,能够对持续到达的数据进行即时分析、聚合和转换。与传统批处理不同,流计算强调低延迟事件驱动持续运行,广泛应用于物联网监控、实时告警、在线统计等场景。

TDengine 作为专为物联网时序数据优化的数据库,内置了流计算引擎,支持在数据写入时自动触发计算,无需部署额外的流处理框架(如 Flink、Spark Streaming)。用户通过 SQL 定义流任务,系统自动管理窗口、状态和结果输出。

2. TDengine 流计算核心特性

  • 多种窗口类型:计数窗口、事件窗口、滑动窗口、会话窗口、定时窗口
  • 灵活触发机制:基于时间或数据行数触发计算
  • 分区与分组:支持按超级表子表名(tbname)、标签、列值进行分组计算
  • 历史数据处理:可回溯处理存量历史数据(FILL_HISTORY_FIRST
  • 结果输出:可写入普通表或超级表,也可仅触发通知(WebSocket)
  • 乱序容忍:支持配置最大延迟时间(MAX_DELAY)处理迟到数据
  • 低资源消耗:计算在数据库内部完成,无需外部依赖

3. 窗口类型说明

窗口类型 触发条件 典型应用场景
计数窗口 每累积指定行数触发一次 每采集 N 个数据点计算一次平均值
事件窗口 条件持续满足一定时间后触发,条件结束则关闭窗口 温度超过阈值持续 10 分钟则告警
滑动窗口 按固定时间间隔滚动或滑动 每 5 分钟计算最近 5 分钟的平均值
会话窗口 不活动超时后关闭窗口 检测设备离线时长
定时触发 固定时间周期触发,不依赖数据 每小时统计一次总数据量

4. 数据模型与示例表结构

本次案例基于以下超级表和子表:

-- 超级表(已存在,包含所有子表的公共字段)
CREATE STABLE properties_hjjcz (
    _ts TIMESTAMP,
    a07 FLOAT,    -- 环境温度
    a05 FLOAT,
    a08 FLOAT,
    a01 FLOAT,    -- PM2.5 浓度
    a09 FLOAT,
    a04 FLOAT,
    a06 FLOAT,
    a03 FLOAT,
    createtime BIGINT,
    a02 FLOAT     -- 电流/功率
) TAGS (deviceid BINARY(32));

-- 创建子表(用于分区计算)
CREATE TABLE properties_hjjcz_hjjcz01 USING properties_hjjcz TAGS ('hjjcz01');
CREATE TABLE properties_hjjcz_hjjcz02 USING properties_hjjcz TAGS ('hjjcz02');
CREATE TABLE properties_hjjcz_hjjcz03 USING properties_hjjcz TAGS ('hjjcz03');
CREATE TABLE properties_hjjcz_hjjcz04 USING properties_hjjcz TAGS ('hjjcz04');

示例插入数据:

INSERT INTO properties_hjjcz_hjjcz01 (_ts, a07, a01, a02, deviceid)
VALUES ('2025-12-09 20:35:06.956', 0.15, 121.12, 180.96, 'hjjcz01');

5. 流计算实战案例

5.1 计数窗口 – 每写入 3 行数据计算 a01(PM2.5)平均值

业务需求:设备每上传 3 条数据,立即计算这 3 条数据的 a01 平均值,并保存到结果表。

CREATE STREAM count_avg_pm25
COUNT_WINDOW(3)
FROM properties_hjjcz_hjjcz01
INTO avg_pm25_count
AS
SELECT _twstart, AVG(a01) AS avg_pm25
FROM %%trows;
  • count_avg_pm25:船舰的流式的表,在taos中通过show streams 可以查看
  • _twstart:窗口起始时间戳
  • %%trows:表示触发时读取到的当前窗口内的数据行(相当于子查询)
  • avg_pm25_count:自动创建目标生成的表,不需要手动创建
    在这里插入图片描述
    在这里插入图片描述

5.2 事件窗口 – 温度超过 80 度持续 1 分钟,计算平均温度

业务需求:当 a07 > 80 的条件连续满足 10 分钟以上,则触发计算该时段内的平均温度。

CREATE STREAM high_temp_event
EVENT_WINDOW(START WITH a07 > 80 END WITH a07 <= 80)
TRUE_FOR(10m)
FROM properties_hjjcz_hjjcz01
INTO high_temp_avg
AS
SELECT _twstart, _twend, AVG(a07) AS avg_temp
FROM %%trows;
  • START WITH a07 > 80 END WITH a07 <= 80 需要满足这个条件才会触发,也就是才会生成一条记录
  • TRUE_FOR(10m):表示条件必须持续成立至少 10 分钟才输出结果
  • 窗口结束条件为 a07 <= 80,此时计算窗口内的平均值

在这里插入图片描述
在这里插入图片描述

5.3 滑动窗口 – 超级表每个子表每 5 分钟计算 pm10 平均值

业务需求:超级表 properties_hjjcz 下有多个子表(不同设备),要求每个设备单独每 5 分钟计算一次最近 5 分钟内的 a01 平均值,结果写入超级表 avg_pm10_stb 的不同子表中。

CREATE STREAM sliding_a01_single
INTERVAL(5m) SLIDING(5m)
FROM properties_hjjcz
PARTITION BY tbname
INTO avg_a01_stb
AS
SELECT _twstart, AVG(a01) AS avg_pm10
FROM %%tbname
WHERE _ts >= _twstart AND _ts < _twend;
  • PARTITION BY tbname:按子表名分组,每个子表独立计算
  • %%tbname:代表当前触发分组对应的子表名,动态替换
  • INTERVAL(5m):窗口大小(时间跨度)
  • SLIDING(5m)):窗口滑动步长,每 5 分钟 计算一次
  • PARTITION BY:分组计算(独立窗口)按子表名分组
    在这里插入图片描述

5.4 定时触发 – 每 1 小时统计子表数据行数,写入毫秒时间戳

业务需求:每小时统计一次 properties_hjjcz_hjjcz01 表中的总数据量,结果保存到 hourly_stats 表,时间戳精确到毫秒。

CREATE STREAM hourly_row_count
PERIOD(1h)
FROM properties_hjjcz_hjjcz01
INTO hourly_stats
AS
SELECT CAST(_tlocaltime / 1000000 AS TIMESTAMP) AS ts, COUNT(*) AS row_count
FROM %%trows;

在这里插入图片描述

  • _tlocaltime:流计算触发时的系统时间(微秒),除以 1000000 转换为毫秒后转为 TIMESTAMP 类型

5.5 滑动窗口(历史优先 + 最大延迟)– 从最早数据开始,每 5 分钟计算 a02 平均值

业务需求:对于超级表每个子表,从历史最早数据开始回溯,每 5 分钟滚动窗口计算 a02 平均值。若窗口因为数据延迟未及时关闭,最多等待 1 分钟。

CREATE STREAM history_first_avg_a02
INTERVAL(5m) SLIDING(5m)
FROM properties_hjjcz
PARTITION BY tbname
STREAM_OPTIONS(MAX_DELAY(1m) | FILL_HISTORY_FIRST)
INTO avg_a02_stb
AS
SELECT _twstart, AVG(a02) AS avg_a02
FROM %%tbname
WHERE _ts >= _twstart AND _ts <= _twend;
  • FILL_HISTORY_FIRST:优先处理历史数据,再处理实时数据
  • MAX_DELAY(1m):允许迟到数据最多延迟 1 分钟,超过此时间的迟到数据被忽略
    在这里插入图片描述
    目前测试服务器3.3.8版本有些问题,只补全今天的数据!

5.6 计数窗口 + 通知 – 每 10 行且 a01>0 时计算平均值并发送 WebSocket

业务需求:每当子表中累积 10 条 a01 > 0 的数据,计算这 10 条数据的 a01 平均值,不保存结果,只通过 WebSocket 发送通知。

CREATE STREAM notify_avg_a01
COUNT_WINDOW(10, 1, a01)
FROM properties_hjjcz_hjjcz01
STREAM_OPTIONS(CALC_NOTIFY_ONLY | PRE_FILTER(a01 > 0))
NOTIFY('ws://localhost:8080/notify') ON (WINDOW_CLOSE)
AS
SELECT 
    NOW AS ts,
    LAST_ROW(a01) AS val
FROM 
    %%trows;
  • COUNT_WINDOW(10, 1, a01):基于 a01 列计数,每增加 10 行触发,步长 1(即每来一条新数据就检查是否满 10 条)
  • PRE_FILTER(a01 > 0):只计入 a01 > 0 的行
  • CALC_NOTIFY_ONLY:不保存计算结果到任何表
  • NOTIFY(...) ON (WINDOW_CLOSE):窗口关闭时向指定 WebSocket 地址发送 JSON 格式的通知

6. 流计算管理命令

  • 查看所有流任务SHOW STREAMS;
  • 删除流任务DROP STREAM stream_name;
  • 暂停/恢复流任务ALTER STREAM stream_name PAUSE; / ALTER STREAM stream_name RESUME;

7. 注意事项与最佳实践

  1. 时间戳列名:流计算中默认使用 _ts 作为时间戳列名,若实际列名为 ts 或其他名称,请在 SQL 中显式使用 _c0 或修改表结构。
  2. 超级表 vs 子表FROM 后跟超级表时需配合 PARTITION BY;跟子表时不需要。
  3. 窗口闭合:滑动窗口只有窗口关闭后才会输出结果,因此实时性要求高的场景应设置较小的 SLIDING
  4. 性能影响:过多的流任务会占用系统资源,建议合理控制并发窗口数。
  5. 通知失败处理:通过 NOTIFY_OPTIONS(ON_FAILURE_PAUSE) 可在通知失败时暂停流任务,防止数据丢失。

8. 总结

TDengine 流计算通过简洁的 SQL 语法实现了强大的实时数据处理能力。结合计数、事件、滑动等多种窗口模型,可以覆盖物联网领域 80% 以上的实时统计和监控需求。实际应用中,应根据数据特征和业务延时要求选择合适的窗口类型,并善用 STREAM_OPTIONS 调优乱序处理和历史数据回填。

我是一名工业互联网开发工程师,服务于纵横工业互联网团队,这是河南863旗下专注工业数字化的团队,也是深耕工业数字化转型领域的专业技术与解决方案服务商。我们聚焦工业企业智能化升级核心需求,打造出全栈式、可落地的工业互联网产品与服务体系,构建了自主可控的五大核心产品体系,包括面向产业集聚区/工业园区的产业集聚区工业互联网管理平台,以及面向工业企业全生产流程的物联网平台、能耗能碳管理平台、设备管理系统、MES生产制造执行系统。

Logo

AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。