TDengine 流式计算快速上手:用 SQL 实现分钟级聚合、降采样与实时告警
发布时间:2026/9/12 2:56:06 锦皓数字建站

TDengine 流式计算快速上手用 SQL 实现分钟级聚合、降采样与实时告警【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengineTDengine面向工业物联网 IIoT 场景的高性能时序数据库内置了流式计算引擎你无需再单独部署 Kafka、Flink 等流处理系统只需一条 SQL 即可定义数据写入后自动触发的实时计算逻辑结果可写入目标表也可通过 WebSocket 通知外部应用。本文基于快速入门章节的test库与taosBenchmark -y生成的meters电表数据带你完整走一遍确认数据 → 创建流 → 查看结果 → 写入新数据观察更新 → 清理示例的全流程并在此基础上结合仓库源码与配套文档深入讲解触发模式、分组、输出表、控制选项等核心机制使你既能立刻跑通第一个流任务也能理解其底层工作原理为后续设计生产级流计算打下基础。为什么需要内置流式计算在时序数据处理中以下需求非常常见分层存储与智能降采样工业设备每秒可能产生上万条原始数据全量保存会导致存储成本飙升、查询效率下降、历史趋势分析响应缓慢。流式计算可以在数据写入时就把分钟级、小时级指标算好让后续查询只扫描聚合结果。预计算加速实时决策当用户查询全量数据集时系统可能需要扫描数百亿条记录几乎无法实时返回结果导致看板与报表滞后。提前生成通用统计量可以大幅降低宽范围查询的延迟。异常检测与低延迟告警监控与告警要求基于预定义规则以极低延迟获取特定数据传统批处理往往存在分钟级延迟。传统时序方案通常要部署 Kafka、Flink 等流处理系统但这类系统的复杂度带来了高昂的开发与运维成本。TDengine 的流式计算引擎位于仓库 source/libs/new-stream 目录提供了实时处理写入数据流的能力用 SQL 定义实时变换后数据一旦写入流的源表就会按定义自动处理并根据触发模式把结果推送到目标表——这是对复杂流处理系统的轻量化替代即使在高吞吐数据写入下也能提供毫秒级结果延迟。从v3.3.7.0起TDengine 提供了全新的流式计算能力完整能力概述见 流式计算总览。触发与计算分离的设计思想与传统的流处理相比TDengine 的流式计算采用触发trigger与计算compute解耦的策略仍然基于连续无界数据流但带来了三方面扩展处理目标扩展传统流处理中事件触发源与计算目标通常是同一个数据集。TDengine 允许触发源事件驱动源与计算源分离——触发表与计算源表可以不同甚至可以没有触发表被处理的数据集在列和时间范围上都可以灵活变化。触发机制扩展除标准的数据写入触发外还支持窗口触发基于事件时间、定时触发基于系统时间等多种触发模式在事件触发前还可以对触发数据进行预过滤只有满足条件的数据才参与触发判定。计算范围扩展计算可以作用于触发表也可以作用于其他库表计算类型不受限制任何查询语句都支持计算结果既可以作为通知发送也可以写入输出表或两者同时进行。核心能力一览在动手实践前先对这套能力有一个整体认知触发模式支持周期触发PERIOD、滑动触发SLIDING、时间窗口INTERVAL、会话窗口SESSION、状态窗口STATE_WINDOW、事件窗口EVENT_WINDOW、计数窗口COUNT_WINDOW以及嵌套窗口等触发可以按组划分触发数据可以预过滤。详见 流语法。计算与结果输出计算可以是任意查询。结果可以写入输出表INTO、作为通知发送NOTIFY或两者兼有。详见 流语法。控制选项通过STREAM_OPTIONS配置历史数据回放FILL_HISTORY/FILL_HISTORY_FIRST、乱序水位线WATERMARK、最大延迟触发MAX_DELAY、低延迟计算LOW_LATENCY_CALC等在结果新鲜度与资源负载之间取得平衡。详见 流语法。运维与限制流任务运行在 snode 上涉及高可用、权限、手动重算、乱序/更新/删除等非典型写入。详见 操作与限制。部署与设计部署方式、配置参数、创建流之前的设计要点与典型示例。详见 部署与设计。下面从一个INTERVAL时间窗口流开始动手。前置条件开始之前请确认以下三点TDengine 服务已运行且能用taos客户端连接。集群中已部署 snode流任务运行在 snode 上。已按快速入门章节执行过taosBenchmark -y生成了test库与meters超级表。如果还没执行请先在终端运行taosBenchmark -y。在客户端中查看 snodeSHOW SNODES;如果没有 snode先查看 dnode 列表再在其中某个 dnode 上创建 snode把示例中的1替换为实际的 dnode IDSHOW DNODES; CREATE SNODE ON DNODE 1;关于 snode 部署的更完整指导可参见 操作与限制部署 snode。从源码看 snode 的角色流式处理的架构是计算与存储分离除数据读取外流处理的所有功能都只在 snode 上执行。每个 dnode 最多承载一个 snode每个 snode 有多个执行线程集群中至少需要一个 snode。为保证流处理高可用建议在多台物理节点上部署多个 snode流任务在多个 snode 间负载均衡每对 snode 互为副本并保存流状态与进度信息——如果集群只部署单个 snode则无法保证高可用。出于资源隔离考虑强烈建议把 snode 部署在专用 dnode 上详见 操作与限制。准备示例数据进入客户端后切换到test库USE test;确认meters超级表已经有数据SELECT tbname, ts, current, voltage FROM meters WHERE voltage 250 and tbname d1 ORDER BY ts DESC LIMIT 5;结果与下面类似子表名和数值可能因taosBenchmark版本或随机数据而略有差异tbname | ts | current | voltage | d1 | 2017-07-14 10:40:09.998 | 11.7984 | 253 | d1 | 2017-07-14 10:40:09.998 | 11.7984 | 253 | d1 | 2017-07-14 10:40:09.998 | 11.7984 | 253 | d1 | 2017-07-14 10:40:09.998 | 11.7984 | 253 | d1 | 2017-07-14 10:40:09.998 | 11.7984 | 253 | Query OK, 5 row(s) in settest.meters中的时间戳落在2017-07-14 10:40:00.000到2017-07-14 10:40:09.999之间。使用下面的 1 分钟窗口时历史数据落在同一个窗口内后续的写入示例会插入下一分钟的数据从而让你观察到新窗口的产生。创建流每分钟计算每块电表的平均电流为了让示例可以反复运行先清理可能存在的同名流与输出表DROP STREAM IF EXISTS avg_current_stream; DROP STABLE IF EXISTS avg_current_stb;执行下面的 SQL 创建流每 1 分钟一个窗口按子表计算平均电流结果写入输出超级表avg_current_stb。输出子表名与源表对应例如d0→avg_d0并保留源表的groupId标签。CREATE STREAM avg_current_stream INTERVAL(1m) SLIDING(1m) FROM meters PARTITION BY tbname, groupId STREAM_OPTIONS(FILL_HISTORY_FIRST | MAX_DELAY(3s)) INTO avg_current_stb OUTPUT_SUBTABLE(CONCAT(avg_, tbname)) TAGS ( groupId INT AS groupId ) AS SELECT _twstart AS ts, _twend AS window_end, AVG(current) AS avg_current FROM %%trows;各子句的含义如下子句含义INTERVAL(1m) SLIDING(1m)以 1 分钟为窗口、1 分钟为滑动步长触发计算。窗口从 Unix 时间 01970-01-01 00:00:00 UTC起对齐划分可通过INTERVAL(interval_val[, interval_offset])中的偏移量改变划分起点FROM meters PARTITION BY tbname, groupId按子表触发把groupId也放入分区列表同时会将groupId作为输出标签STREAM_OPTIONS(FILL_HISTORY_FIRST \| MAX_DELAY(3s))优先计算已写入的历史数据打开的窗口在打开约 3 秒后也会触发一次方便快速演示INTO avg_current_stb把结果写入输出超级表。有触发分组时输出表为超级表无分组时输出表为普通表OUTPUT_SUBTABLE(CONCAT(avg_, tbname))根据源表名构造输出子表名例如d0→avg_d0TAGS (groupId INT AS groupId)把源分组标签复制到输出超级表便于按标签分组查询%%trows当前触发窗口内的行集合只能作为查询表名出现在FROM子句中关于创建流的完整语法与参数见 流语法CREATE STREAM。关于几个关键机制的原理解读分组PARTITION BY是最小执行单元在流式计算中组是最小的执行单元。逻辑上可以把每个组看作一个独立的流任务拥有各自的输出表和事件通知。如果不指定分组定时触发场景下也可不指定触发表整个流只有一个组对应一张输出表和一条通知。由于每个组独立运行它们的计算进度、输出频率等行为可以互不相同。这也是一设备一表设计哲学在流计算中的体现——一次建流即可为所有子表分别计算结果。输出表的数量等于触发表的分组数量若未指定分组则只创建一张普通输出表。%%trows的使用限制%%trows引用触发表中每个组满足当前触发条件的数据集仅适用于WINDOW_CLOSE触发且只能用于FROM子句使用%%trows的查询不支持对%%trows加WHERE条件过滤或做 JOIN。若要基于窗口时间范围取全量数据而非触发时刻的快照可改用FROM %%tbname WHERE _c0 _twstart AND _c0 _twend的形式——前者只使用触发时捕获的窗口数据后者在计算时重新读取窗口时间范围内、对应组表的数据两者并不等价对比示例见 部署与设计滑动触发。FILL_HISTORY_FIRST与FILL_HISTORY的区别FILL_HISTORY[(start_time)]从指定时间默认最早记录开始计算历史数据FILL_HISTORY_FIRST[(start_time)]则优先计算历史数据历史处理完成前不开始实时计算适合必须严格按时间顺序处理历史的场景例如COUNT_WINDOW必须优先处理历史否则窗口可能无法对齐。两者互斥且都不支持PERIOD定时触发模式。MAX_DELAY(delay_time)的行为定义窗口尚未关闭时强制触发的最大等待时间处理时间。从窗口打开起若窗口仍处于打开状态则每隔delay_time产生一次触发。delay_time支持秒s、分钟m、小时h、天d最小值 3 秒精度容差约 1 秒。注意WATERMARK先于窗口判定求值可能导致设置了MAX_DELAY却不触发的情况因为窗口从未真正打开。非窗口触发会自动忽略该选项。查看流创建完成后列出当前库中的流SHOW STREAMS;结果类似如下stream_name | status | message | db_name | avg_current_stream | Idle | Current deploy times: 0 | test | Query OK, 1 row(s) in set如需更详细的信息查询系统表SELECT * FROM information_schema.ins_streams WHERE stream_name avg_current_stream;流的调度是异步的创建后请等待几秒再查询输出表。由于设置了FILL_HISTORY_FIRST约 1 万个子表的历史窗口需要回放首次计算可能会稍慢。information_schema.ins_streams中还提供了实时处理延迟realtime_lag_ms、输入/输出速率input_rows_per_sec_1m/output_rows_per_sec_1m、结果形成延迟runner_result_latency_avg_1m_ms、历史处理进度history_progress_pct等运行时指标可用于监控与排查详见 可观测性与故障排查。从源码看任务模型一个运行中的流会以多个任务task的形式执行任务类型包括 Reader读取触发数据、Trigger判定触发条件、Runner执行计算与写出结果等。仓库 source/libs/new-stream/src 中的streamRunner.c计算执行、streamTriggerTask.c触发任务、streamTriggerMerger.c触发合并、streamReader.c/streamReaderExt.c数据读取、streamCheckpoint.c检查点等文件即对应这一流水线的各个环节。流级视图定位问题后可用SELECT * FROM information_schema.ins_stream_tasks下钻到任务级视图通过status、last_update、node_id和任务级指标定位受影响的节点或任务。查看计算结果查询输出超级表avg_current_stbSELECT tbname, groupId, ts, window_end, avg_current FROM avg_current_stb ORDER BY ts, tbname LIMIT 5;每一行是某块电表在 1 分钟窗口内的平均电流。tbname是输出子表名例如avg_d0。历史数据落在10:40:00–10:41:00窗口内如下所示tbname | groupId | ts | window_end | avg_current | avg_d0 | 1 | 2017-07-14 10:40:00.000 | 2017-07-14 10:41:00.000 | 10.208475 | avg_d1 | 7 | 2017-07-14 10:40:00.000 | 2017-07-14 10:41:00.000 | 10.208475 | avg_d2 | 2 | 2017-07-14 10:40:00.000 | 2017-07-14 10:41:00.000 | 10.208475 | avg_d3 | 4 | 2017-07-14 10:40:00.000 | 2017-07-14 10:41:00.000 | 10.208475 | avg_d4 | 3 | 2017-07-14 10:40:00.000 | 2017-07-14 10:41:00.000 | 10.208475 |具体的groupId值取决于taosBenchmark为表分配的标签可能与上表略有出入。关于输出表结构的要点输出表中的第一列即输出表的主键列必须是合法的TIMESTAMP主键值执行时若出现NULL主键对应的计算结果会被丢弃。窗口聚合的起始时间_twstart、结束时间_twend等触发上下文通过占位符注入查询窗口触发类型可用的占位符包括_twstart当前窗口开始时间戳、_twend当前打开窗口的结束时间戳仅WINDOW_CLOSE触发使用、_twduration窗口时长、_twrownum窗口行数另有适用于所有触发的_tgrpid组 ID、_tlocaltime当前触发系统时间、%%n分组列引用、%%tbname触发表引用等。完整占位符表见 流语法占位符。写入新数据并观察结果更新窗口默认在关闭时触发。如果只在10:41:00–10:42:00窗口内写数据该窗口尚未关闭本例设置了MAX_DELAY(3s)因此窗口打开约 3 秒后仍会产生一次结果。向该窗口内为d0插入一行数据INSERT INTO d0 VALUES (2017-07-14 10:41:30, 12.4, 221, 147);也可以在下一个窗口开始时写入一行用事件时间关闭上一个窗口此时即使没有MAX_DELAY也会触发INSERT INTO d0 VALUES (2017-07-14 10:42:00, 12.5, 220, 147);等待几秒后查询d0对应输出子表avg_d0的最新流计算结果SELECT tbname, groupId, ts, window_end, avg_current FROM avg_current_stb WHERE tbname avg_d0 ORDER BY ts DESC LIMIT 3;结果类似如下10:41:00窗口已出现在输出表中tbname | groupId | ts | window_end | avg_current | avg_d0 | 1 | 2017-07-14 10:41:00.000 | 2017-07-14 10:42:00.000 | 12.400000 | avg_d0 | 1 | 2017-07-14 10:40:00.000 | 2017-07-14 10:41:00.000 | 10.208475 |流会持续运行只要meters中有新数据写入满足触发条件的窗口就会继续计算并写入输出表。事件类型EVENT_TYPE与更多触发时机默认只监听WINDOW_CLOSE窗口关闭事件但通过STREAM_OPTIONS(EVENT_TYPE(...))可以选择WINDOW_OPEN窗口打开、WINDOW_CLOSE、IDLE组空闲、RESUME组恢复等事件类型多个类型用|分隔。其中IDLE组超过IDLE_TIMEOUT时长未收到新数据时触发一次需要同时配置IDLE_TIMEOUTRESUME空闲组再次收到数据时立即触发同样需要IDLE_TIMEOUTIDLE_TIMEOUT的取值范围为[1s, 10d]空闲检测基于处理时间并使用单调时钟避免受系统时钟调整影响。这为数据停止到达这类无数据事件也提供了触发与通知能力。清理示例如果不再需要本章创建的流与输出表执行DROP STREAM IF EXISTS avg_current_stream; DROP STABLE IF EXISTS avg_current_stb;注意DROP STREAM只删除流处理任务流已经写出的数据不会被删除。常见调整在快速入门阶段请留意以下常见调整只处理新写入数据从STREAM_OPTIONS中移除FILL_HISTORY_FIRST即可。希望在窗口关闭前尽快拿到结果保留或调整MAX_DELAY也可以在下一窗口开始处写数据用事件时间关闭当前窗口。存在乱序写入、更新或删除结合WATERMARK、重算机制与最佳实践来设计流。WATERMARK(duration_time)定义乱序数据容忍时长当前水位 已处理的最新事件时间 −WATERMARK时长只有事件时间早于当前水位的数据和窗口才参与触发求值超过水位容忍范围的乱序/更新/删除通过重算重新触发受影响数据范围的计算结果不删除而是再次写入保证正确性。注意WATERMARK不适用于PERIOD定时触发。乱序数据、数据更新、数据删除在各触发类型下的影响与处理策略对照表以及EXPIRED_TIME过期数据的语义详见 操作与限制非典型数据写入。把结果发送给外部应用使用NOTIFY创建通知流。通知通过 WebSocket 协议发送到外部应用NOTIFY(url [, ...]) ON (event_types)可指定事件类型与目标地址例如NOTIFY (ws://localhost:8080/notify, wss://192.168.1.1:8080/notify?keyfoo) ON (WINDOW_OPEN | WINDOW_CLOSE)通知消息体为 JSON 格式包含messageId、timestamp、streams含streamName与事件数组事件对象携带tableName、eventType、eventTime、triggerId、triggerType、groupId以及各窗口类型专属字段如时间窗口的windowStart/windowEnd和计算结果result。若只发通知不计算结果、或计算结果仅通知不落表可使用STREAM_OPTIONS(CALC_NOTIFY_ONLY)。完整的通知消息字段定义见 流语法通知机制。创建流之前的设计检查清单来自 部署与设计 的实践建议值得在正式建流前逐条过一遍按业务特征选触发类型每写入一条/若干条数据就处理→计数窗口满足窗口条件再处理→窗口触发按事件时间周期计算→滑动触发按处理时间周期计算→周期触发。按时效要求选窗口选项决定窗口打开、关闭或两者都触发是否用MAX_DELAY在窗口关闭前及时计算。选好触发表周期触发可省略触发表但若需要分组输出或在定时间隔内使用%%trows则必须指定触发表。保证源表与触发表的时间一致性两者不同时须确保触发表事件触发时源表已有有效数据。按业务决定分组方式全局聚合可不分组按标签聚合用标签分组需要单表结果用子表分组。分组过多会导致结果表过多可通过OUTPUT_SUBTABLE(tbname_expr)让多组合并到同一子表并配合输出表复合主键设计。评估乱序影响乱序严重时可考虑IGNORE_DISORDER对来自久远过去、已无时效价值的乱序记录可用EXPIRED_TIME标记过期。验证重算的有效性重算要求计算语句与源表独立于处理时间——同一触发即使执行多次也应产生有效结果。确认历史数据处理方式优先历史用FILL_HISTORY_FIRST否则用FILL_HISTORY。确认实时性要求要求极高时效性时使用LOW_LATENCY_CALC会消耗更多计算资源默认计算可能被延迟并批量执行以提升资源效率。部署与配置参考部署把 snode 部署在专用 dnode 上该 dnode 不应再承载 vnode、mnode、qnode在定义流之前预先创建多个 snode 以获得更好的负载均衡流计算负载较重时通过新增 snode 横向扩展。配置流处理相关配置参数完整说明见 taosd 组件文档包括配置项作用numOfMnodeStreamMgmtThreadsmnode 上的流管理线程数numOfStreamMgmtThreadsvnode/snode 上的流管理线程数numOfVnodeStreamReaderThreadsvnode 上的流读取线程数numOfStreamTriggerThreads流触发线程数numOfStreamRunnerThreads流执行线程数streamBufferSize流处理可用最大缓冲区单位 MB仅用于缓存%%trows的结果streamNotifyMessageSize事件通知消息大小控制streamNotifyFrameSize事件通知消息发送时的底层帧大小控制线程数与缓冲区大小应根据部署方式与负载调整线程越多 CPU 消耗越高负载重或并发流多时应增大缓冲区。更进一步的阅读本章只覆盖了一个最小可运行的流示例。要继续掌握 TDengine 流式计算的全部能力请继续阅读流式计算总览流处理概述、场景、触发/计算分离与能力扩展流语法CREATE STREAM、触发方式周期/滑动/时间窗口/会话/状态/事件/计数/嵌套窗口、结果输出、控制选项与通知语法操作与限制snode、权限、重算、乱序写入、配置参数与规则限制部署与设计部署、配置、建流前设计与典型示例可观测性与故障排查通过系统视图监控流延迟、吞吐、历史处理与重算进度数据查询窗口查询、聚合与结果解读仓库中与流处理相关的核心实现位于 source/libs/new-stream/src流管理streamMgmt.c、触发任务streamTriggerTask.c、计算执行streamRunner.c、窗口链streamWindowChain.c、检查点streamCheckpoint.c等配套测试用例集中在 test/cases/18-StreamProcessing可以结合源码与测试进一步验证本文涉及的语法与行为。【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
锦
锦皓数字建站
深耕本土企业品牌数字化升级,专注原创端正雅致商务官网,从视觉设计到稳定运维全程保驾护航。