一、先说说我们遇到的问题
在物联网项目里,最常见的动作就是往数据库里不停地写入传感器数据。温度、湿度、电表读数、机器转速……每秒钟都有成千上万条记录进来。这些数据很有价值,但也很让人头疼。
头疼在哪呢?比如老板想看一下过去一个月里某台设备的平均温度曲线。按最朴素的办法,去原始数据表里找到这个设备的全部记录,然后一条条算平均值。听起来简单,可数据量一大就不行了。假设设备每秒上报一条,一天就是八万六千多条,一个月就有两百多万条记录。如果线上有几百台设备,原始表里的数据就是几亿行。查询一次要扫这么多行,再一个个做聚合,数据库累得够呛,用户等得也不耐烦。往往一个页面加载要好几秒,甚至几十秒。
有人说,加索引、加缓存呗。但对于时序数据库来说,数据是持续不断增长的,查询范围又往往很大,光靠索引很难从根本上解决聚合慢的问题。更关键的是,我们常见的那些图表,比如“过去一个月的平均温度”,其实是一个很固定的统计需求。既然每次都要算同样的东西,为什么不能提前算好存起来呢?这就是连续聚合能解决的事。
二、什么是降采样,为啥能帮忙
2.1 降采样的核心思路
降采样,说白了就是把高频数据变成低频数据。比如原本一秒一条原始记录,我把它变成一分钟一条,每条数据代表这一分钟的平均值、最大值、最小值或者总和。这样一来,数据量就缩小了六十倍。查询一分钟粒度的趋势时,只需要读那一分钟一条的汇总表,速度自然快了很多。
不过,如果我们每次都手动去算,那和查询时实时算没有区别。真正聪明的方法,是让数据在写入的时候就顺便聚合好。这个“顺便聚合好”的过程,就是连续聚合。
2.2 连续聚合和普通查询的区别
普通查询是“等用户来问才干活”。查询来了,数据库才去扫原始数据,算完就完。连续聚合则是“数据一进来就干活”。数据库按照约定好的时间窗口,把新数据聚合成一条结果,存到另一张预计算结果表里。用户后续的查询直接走预计算结果,完全不碰原始数据。
打个比方,原始数据就像仓库里堆满的散货。连续聚合就像是一个流水线,每来一批货就自动打包成标准箱放好。用户要货时,直接去取标准箱就行,不用再翻散货堆。这样仓库门口就不会排长队了。
三、TDengine里怎么用连续聚合
TDengine是专为时序数据设计的数据库,自带流式计算能力。用它做连续聚合,只需要写一条 CREATE STREAM 语句,告诉它数据从哪来、怎么聚合、结果放哪。接下来它自己会维护。
3.1 建库建表并写入示例数据
假设我们做了一个智能电表管理平台,电表每秒上报一次功率。我们想每分钟算一次平均功率和峰值功率。先建库建表,再插入几条数据感受一下。
下面使用的技术栈是 TDengine SQL,你可以通过 taos 命令行工具或者 REST 接口来执行这些语句。
-- 创建一个数据库,保留100天的数据
CREATE DATABASE iot_data KEEP 100;
-- 切换到该数据库
USE iot_data;
-- 创建超级表。超级表是一类设备的模板,下面可以挂多个子表
CREATE STABLE meters (
ts TIMESTAMP, -- 记录时间
power FLOAT, -- 当前功率,单位kW
device_id NCHAR(20) -- 设备编号
) TAGS (
location NCHAR(20) -- 设备所在位置
);
-- 创建两个子表,同时插入数据
-- USING meters TAGS(...) 表示基于超级表创建子表,并指定标签值
INSERT INTO meter_dev001 USING meters TAGS('北京机房') VALUES
('2025-01-06 10:00:00.000', 32.5, 'DEV-001'),
('2025-01-06 10:00:01.000', 35.2, 'DEV-001'),
('2025-01-06 10:00:02.000', 30.1, 'DEV-001');
INSERT INTO meter_dev002 USING meters TAGS('上海机房') VALUES
('2025-01-06 10:00:00.000', 18.7, 'DEV-002'),
('2025-01-06 10:00:01.000', 20.3, 'DEV-002');
说明一下,TDengine里的超级表本身不直接存数据,实际存储都在子表里。上面的 INSERT 语句在插入的同时自动创建了子表,并打上了标签。真实项目里,也可以提前手工创建子表,再把数据写进去。
3.2 用时间窗口做降采样查询
如果现在就想手工算一下每分钟的平均功率,可以用 INTERVAL 子句。它会把时间戳划分成固定长度的时间窗口,在每个窗口内做聚合。
-- 按1分钟窗口计算每个设备的平均功率和最大功率
SELECT
_wstart, -- 窗口开始时间
device_id,
AVG(power) AS avg_power, -- 平均功率
MAX(power) AS max_power -- 峰值功率
FROM iot_data.meters
WHERE ts >= '2025-01-06 10:00:00.000'
AND ts < '2025-01-06 10:01:00.000'
INTERVAL(1m);
手动执行这种查询,在数据量小的时候响应很快。但当数据积攒到几千万条时,每次都要全表扫描计算,压力会越来越大。所以我们需要连续聚合,让数据库自己把结果维护好。
3.3 创建连续聚合,自动维护预计算结果
TDengine里创建连续聚合任务,核心就是 CREATE STREAM 语句。下面这个例子,把原始表 meters 的数据按1分钟窗口聚合,结果写入 meter_1m 表。
-- 创建连续聚合任务,命名为 stream_meter_1m
-- 结果输出到 iot_data.meter_1m,如果表不存在,TDengine会自动创建
CREATE STREAM stream_meter_1m
INTO iot_data.meter_1m
(
ts,
device_id,
avg_power,
max_power
) AS
SELECT
_wstart AS ts, -- 用窗口开始时间作为结果表的时间戳
device_id,
AVG(power) AS avg_power, -- 计算平均功率
MAX(power) AS max_power -- 计算峰值功率
FROM iot_data.meters
WHERE power IS NOT NULL -- 过滤掉空值
INTERVAL(1m); -- 窗口大小是1分钟
创建完任务之后,TDengine会持续监听原始表的新数据。每当有新的电表数据写入,流计算就会把对应的功率值累加到相应的1分钟窗口里。窗口结束时,把聚合结果写进 meter_1m 表。我们不需要写任何定时任务,也不用关心调度细节。
3.4 查询预计算结果表
几分钟后,我们来看一下预计算结果:
-- 直接查询连续聚合生成的结果表,关注 DEV-001 设备的情况
SELECT
ts,
device_id,
avg_power,
max_power
FROM iot_data.meter_1m
WHERE device_id = 'DEV-001'
AND ts >= '2025-01-06 10:00:00.000'
AND ts <= '2025-01-06 10:05:00.000'
ORDER BY ts;
因为结果表里每分钟只有一条记录,数据量比原始表小了几个数量级,所以这个查询速度非常快。前端做趋势图时,直接依赖这个结果表就好。
3.5 滚动窗口与滑动窗口的补充
上面用到的 INTERVAL(1m) 是滚动窗口,意思是从整分钟开始,1分钟一个窗口,窗口之间不重叠。比如10:00:00到10:00:59是一个窗口,10:01:00到10:01:59是下一个窗口。
有些场景需要滑动窗口,比如每秒都要显示“最近一分钟的平均值”,此时窗口之间会重叠很多。TDengine支持用 SLIDING 子句实现滑动窗口。下面是一个滑动窗口的示例:
-- 每30秒滑动一次,窗口长度1分钟
-- 这表示每30秒计算一次最近1分钟的平均功率
SELECT
_wstart,
_wend,
device_id,
AVG(power) AS avg_power
FROM iot_data.meters
WHERE power IS NOT NULL
INTERVAL(1m) SLIDING(30s);
注意,在连续聚合中使用滑动窗口时,结果表里可能会产生大量重叠的记录。实际项目中,滚动窗口更常用,因为数据量更可控。大家可以根据业务需求选择。
3.6 多种时间粒度的组合
有些场景我们既要看分钟级趋势,也要看小时级趋势。最简单的方式是创建多个连续聚合任务,分别输出到不同的结果表。比如再创建一个小时级的聚合:
-- 创建小时级连续聚合任务
CREATE STREAM stream_meter_1h
INTO iot_data.meter_1h
(
ts,
device_id,
avg_power,
max_power
) AS
SELECT
_wstart AS ts,
device_id,
AVG(power) AS avg_power,
MAX(power) AS max_power
FROM iot_data.meters
WHERE power IS NOT NULL
INTERVAL(1h);
这样分钟表和小时表并存。前端缩放的时候,按照不同的时间范围去查不同的表。比如用户只看最近一小时,查分钟表就可以了;用户看一个月,直接查小时表甚至天表。
四、实际性能对比
为了直观地说明效果,我们做个假设。假设原始表里有三千万条记录。你要查某台设备最近一个月每分钟的平均功率。如果直接去原始表实时聚合,数据库要扫描几百万行,可能得花好几秒钟。而查预先生成的分钟聚合表,这个表里只有大约四万条记录(一个月30天乘24小时乘60分钟等于43200条),扫描起来非常轻快,响应时间通常在几十毫秒内。
当然,真实性能取决于服务器配置、数据分布、并发查询量等因素。但从原理上讲,查询结果表需要的计算量和IO次数都远低于原始表。所以这个方案能大大缩短查询响应时间,属于“空间换时间”的经典套路。
五、应用场景
5.1 物联网设备监控
机房温湿度、智能电表、工厂设备状态,这类数据上报频率高,而且需要长时间看趋势。用连续聚合生成分钟、小时、天级别的汇总数据,最适合不过。运维人员看大屏时,图表能非常快地刷出来。
5.2 金融交易行情
行情数据每秒都在变,但图表展示用1分钟K线基本够了。提前把每分钟的开高低收算好,查询K线图就不用再扫所有tick数据。这里需要注意,TDengine的流计算可以很好地处理这类高频写入和低频聚合。
5.3 运维监控指标
服务器CPU、内存、网络流量,通常几十秒采集一次。运维平台查询历史趋势时,如果每次都实时聚合,数据量一大就慢。用连续聚合后,监控页面变得非常顺畅,告警故事也更容易讲清楚。
5.4 其他高基数时序场景
只要有“高频写入,低频分析”的需求,都可以考虑。比如车联网轨迹数据、环境监测数据、燃气表数据、水表数据等。甚至电商网站的点击流日志,只要带时间戳,都能用这个思路做预聚合。
六、技术优缺点
6.1 好处
- 查询响应快。预计算结果表数据量小,扫描快,用户体验好。
- 降低数据库计算压力。聚合是持续增量完成的,不会在查询高峰期集中占用CPU。
- 代码简单。一条SQL就能定义聚合规则,不用自己开发定时任务和调度逻辑。
- 实时性好。数据流入后很快就能在聚合表中查到,通常延迟在秒级以内。
- 可扩展性好。TDengine的流计算可以水平扩展,配合分布式架构能支撑更大的数据量。
6.2 需要留意的代价
- 额外的存储空间。虽然结果表比原始表小很多,但总归要多占一点磁盘。
- 原始精度丢失。聚合后的数据无法还原成原始值,所以如果需要看明细数据,原始表还是要保留。
- 维护成本。任务多了以后,需要管理流任务的状态、失败重试等。TDengine提供了一些状态查询机制,但作为使用者,我们仍然需要了解底层原理。
七、注意事项
- 窗口设计要合理。窗口太小,结果表还是很大,查询提升有限;窗口太大,趋势不够平滑,细节丢失。通常根据业务需求来定,比如分钟级、小时级、天级。
- 时间基准。TDengine的窗口以服务端时区对齐,默认从
1970-01-01 00:00:00开始。如果你的业务需要按自然日或整点,要记得检查时区设置。 - 数据补历史。如果原始表里已经有了历史数据,再创建连续聚合任务时,流计算通常只处理新写入的数据,不会自动回填历史。如果需要历史聚合,可以手动写一条带
INTERVAL的查询,把结果先写入目标表,再启用流任务。 - 注意标签和过滤条件。超级表上有很多子表,如果连续聚合任务不加过滤,会对所有设备生效。如果你只关心某几个设备,最好在SQL里把
WHERE条件写清楚。 - 定期清理。原始数据占用的空间最大,如果业务并不需要保存很长的明细,可以设置数据过期策略,让TDengine自动清理。同时保留聚合结果表,这样既能节省成本,又能保证查询效率。
- 结果表命名。建议在结果表名里带上时间粒度,比如
meter_1m、meter_1h、meter_1d。这样团队里其他人一看就明白这张表是什么粒度,避免拿错表。
八、总结
降采样和连续聚合是处理时序数据查询慢的一对好搭档。TDengine把这两件事结合得非常自然,写一条 CREATE STREAM 语句,数据库就自动帮你维护好了预计算结果。用户查询时直接面向小表,响应时间大幅缩短。当然,它并不是银弹,设计时要考虑窗口大小、历史数据回填、存储开销等细节。但只要规划好,这个方案能让你在数据量翻倍的情况下,依然保持丝滑的查询体验。
实际落地时,可以从一个最简单的分钟级窗口开始,熟悉流任务的状态和结果表结构,再逐步增加小时级、天级聚合。数据预计算这件事,做起来并不复杂,但带来的收益非常可观。希望这篇文章能帮你迈出第一步。
评论
围绕“连续聚合落地:使用TDengine降采样功能减少查询响应时间的方案”参与讨论