欢迎光临
我们一直在努力

Apache IoTDB 连续查询(CQ)全解析:从语法到实战,手把手教你玩转实时数据计算

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

Apache IoTDB 连续查询(CQ)全解析:从语法到实战,手把手教你玩转实时数据计算

本文围绕 IoTDB 的连续查询(CQ)展开全面解析,先通俗阐释其定义,即对实时数据周期性自动执行查询并存储结果的 “自动计算器”,可实现滑动窗口流式计算。接着详细拆解语法,包括各参数含义、格式及注意事项,如 WHERE 子句不可加时间过滤等。随后通过 5 个实战案例,从不同配置场景演示 CQ 创建、执行逻辑与结果,覆盖周期配置、窗口设置、空值填充等情况。还介绍了 CQ 的管理操作(查询、删除、修改)、适用场景(数据降采样、预计算昂贵查询、替代子查询)及性能优化配置参数,最终帮助读者掌握 CQ 用法,提升实时数据计算效率。

在这里插入图片描述

在时序数据库的日常使用中,我们经常会遇到这样的需求:对实时涌入的传感器数据,每隔一段时间就自动计算一次统计指标,比如每小时的平均温度、每分钟的电压最大值,还要把结果存到新的序列里方便后续分析。要是每次都手动写查询语句执行,不仅麻烦还容易出错。这时候,IoTDB的连续查询(Continuous Query,简称CQ)就派上大用场了。今天这篇文章,就从连续查询的基本概念入手,一步步带你吃透它的语法规则、实战用法、管理技巧和适用场景,最后再补充配置参数的细节,让你看完就能上手用。

一、什么是连续查询?用大白话讲清楚

先抛开复杂的定义,咱们用通俗的话解释下连续查询。它就像一个“自动工作的计算器”——你提前设定好规则,比如“每隔20秒,查一下过去40秒内所有温度传感器的最大值,然后把结果存到带‘_max’后缀的序列里”,之后它就会按照这个规则周期性地自动执行,不用你再手动操作。

具体来说,连续查询是对实时数据周期性自动执行的查询,执行结果会自动写入你指定的时间序列。它最核心的作用是实现滑动窗口流式计算,比如计算某个设备每10秒的平均湿度、某条生产线每5分钟的设备运行次数等。而且,你还能通过自定义RESAMPLE子句来创建不同的滑动窗口,即便数据偶尔乱序到达,也能有一定的容忍度,不用怕计算结果出错。

举个简单的例子:假设你有一批温度传感器,每秒都会上传一次数据,你需要看每小时的平均温度。如果不用连续查询,你可能得每小时手动写一次SELECT AVG(temperature) FROM … WHERE time > …,还要手动把结果存起来。但用了连续查询,只需要创建一次任务,它就会每小时自动完成查询、计算、存储这一整套流程,省心又高效。

二、连续查询语法详解:每个参数都讲透

要用好连续查询,首先得把它的语法规则弄明白。IoTDB的连续查询语法有固定结构,咱们先看完整的语法模板,再逐个拆解里面的关键部分,最后强调几个容易踩坑的注意事项。 在这里插入图片描述 在这里插入图片描述

1. 完整语法模板

先看一下创建连续查询的完整SQL结构,后面会逐个部分详细解释:

CREATE (CONTINUOUS QUERY | CQ) <cq_id>
[RESAMPLE
[EVERY <every_interval>]
[BOUNDARY <execution_boundary_time>]
[RANGE <start_time_offset>[, end_time_offset]]
]
[TIMEOUT POLICY BLOCKED|DISCARD]
BEGIN
SELECT CLAUSE — 要计算的指标,比如MAX(temperature)、AVG(voltage)
INTO CLAUSE — 结果要写入的目标序列,比如root.ln.wf02.wt02.temp_max
FROM CLAUSE — 数据来源的序列,支持通配符,比如root.ln.*.*
[WHERE CLAUSE] — 注意:这里不能加时间过滤条件!
[GROUP BY(<group_by_interval>[, <sliding_step>]) [, level = <level>]] — 按时间分组,比如每10秒一组
[HAVING CLAUSE] — 对分组结果过滤,比如只保留MAX(temperature) > 100的数据
[FILL {PREVIOUS | LINEAR | constant}] — 填充空值,比如用100.0填充没有数据的时间段
[LIMIT rowLimit OFFSET rowOffset] — 限制返回结果行数,比如只取前100行
[ALIGN BY DEVICE] — 按设备对齐结果,方便按设备查看数据
END

2. 关键参数逐个拆解

上面的语法里,很多参数是可选的,但理解每个参数的作用是用好连续查询的关键。下面用表格的形式,把每个参数的含义、用法、注意事项都列清楚,方便你对照查看:

参数名称作用说明格式要求默认值常见示例
<cq_id> 给连续查询起一个全局唯一的名字,后续管理(删除、查询)都靠它 自定义字符串,不能和已有的CQ重名 无,必须指定 cq_temp_max(温度最大值连续查询)、cq_voltage_avg(电压平均值连续查询)
EVERY <every_interval> 指定连续查询周期性执行的间隔,比如每隔20秒执行一次 支持的时间单位:ns(纳秒)、us(微秒)、ms(毫秒)、s(秒)、m(分)、h(时)、d(天)、w(周);值必须大于配置文件里的continuous_query_min_every_interval 等于GROUP BY子句中的group_by_interval EVERY 20s(每20秒执行一次)、EVERY 5m(每5分钟执行一次)
BOUNDARY <execution_boundary_time> 你期望的连续查询第一个周期任务的执行时间 时间格式(如2024-05-20T10:00:00+08:00),可以早于、等于或晚于当前时间 0(即当前时间开始计算第一个周期) BOUNDARY '2024-05-20T10:00:00+08:00'(希望从10点开始第一个周期)
RANGE <start_time_offset>[, end_time_offset] 定义每次查询的时间窗口范围:- start_time_offset:窗口开始时间,即now() – start_time_offset(当前时间减去这个偏移量)- end_time_offset:窗口结束时间,即now() – end_time_offset 两个参数都支持和EVERY一样的时间单位;且start_time_offset > end_time_offset – start_time_offset默认等于EVERY的值- end_time_offset默认等于0 RANGE 40s(窗口是过去40秒到当前时间)、RANGE 40s, 20s(窗口是过去40秒到过去20秒)
TIMEOUT POLICY 处理“前一个查询还没执行完,下一个周期的执行时间就到了”的情况,有两种策略可选 二选一:BLOCKED 或 DISCARD BLOCKED – BLOCKED:等前一个查询执行完,再执行下一个,所有窗口都会执行,但可能会延迟- DISCARD:直接跳过下一个周期的查询,不会延迟,但可能有部分窗口没执行
GROUP BY(<group_by_interval>[, <sliding_step>]) 按时间分组计算,比如每10秒一组统计最大值 – group_by_interval:分组间隔,必须大于0,且小于等于start_time_offset- sliding_step:可选,分组滑动步长,默认等于group_by_interval 无,若不指定GROUP BY,则必须显式指定EVERY GROUP BY(10s)(每10秒分一组)、GROUP BY(5m, 2m)(每5分钟一组,每2分钟滑动一次)
FILL 填充没有数据的时间区间,避免结果里出现大量NULL值 三种方式:- PREVIOUS:用前一个有值的结果填充- LINEAR:线性插值填充- constant:用指定的常量填充(如100.0) 不填充,默认显示NULL FILL(PREVIOUS)、FILL(100.0)
ALIGN BY DEVICE 按设备对齐查询结果,让同一设备的不同指标在同一行显示,方便按设备查看数据 无参数,直接写ALIGN BY DEVICE 不对齐,默认按时间和指标混合显示 在查询多个设备的多个指标时常用,比如SELECT temperature, voltage FROM … ALIGN BY DEVICE

3. 必看的语法注意事项

这几个坑很多新手都会踩,一定要记牢:

  • WHERE子句不能加时间过滤:因为连续查询会自动为每次执行的查询指定时间范围(就是你通过RESAMPLE设置的窗口),如果你手动加了WHERE time > '2024-05-20T10:00:00'这类时间条件,IoTDB会直接抛出异常。
  • GROUP BY TIME子句不能加时间范围参数:普通查询里GROUP BY TIME可以写GROUP BY([2024-05-20T10:00:00, 2024-05-20T11:00:00), 10s),但连续查询里不行!IoTDB会自动填充这个时间范围,你要是手动加了,同样会报错。
  • 必须有“周期”或“分组”:如果连续查询既没有GROUP BY TIME子句,也没指定EVERY子句,IoTDB会报错。简单说,就是得让连续查询知道“多久执行一次”或者“按多久分组计算”,不能两者都没有。
  • 时间参数的大小关系要注意:every_interval、start_time_offset、group_by_interval这三个值都必须大于0;而且group_by_interval要小于等于start_time_offset,不然分组间隔比窗口还大,计算就没意义了。
  • 三、连续查询实战:5个经典案例,代码+结果全解析

    光看语法不够直观,咱们结合实际数据来做实战演示。先说明下实验用的基础数据:有4个温度传感器,分别是root.ln.wf02.wt02.temperature、root.ln.wf02.wt01.temperature、root.ln.wf01.wt02.temperature、root.ln.wf01.wt01.temperature,数据从2021-05-11T22:18:14开始,每隔几秒上传一次,具体数据如下表:

    时间root.ln.wf02.wt02.temperatureroot.ln.wf02.wt01.temperatureroot.ln.wf01.wt02.temperatureroot.ln.wf01.wt01.temperature
    2021-05-11T22:18:14.598+08:00 121.0 72.0 183.0 115.0
    2021-05-11T22:18:19.941+08:00 0.0 68.0 68.0 103.0
    2021-05-11T22:18:24.949+08:00 122.0 45.0 11.0 14.0
    2021-05-11T22:18:29.967+08:00 47.0 14.0 59.0 181.0
    2021-05-11T22:18:34.979+08:00 182.0 113.0 29.0 180.0
    2021-05-11T22:18:39.990+08:00 42.0 11.0 52.0 19.0
    2021-05-11T22:18:44.995+08:00 78.0 38.0 123.0 52.0
    2021-05-11T22:18:49.999+08:00 137.0 172.0 135.0 193.0
    2021-05-11T22:18:55.003+08:00 16.0 124.0 183.0 18.0

    下面5个案例,覆盖了连续查询的常见用法,每个案例都包含“创建语句”“执行逻辑解释”“实际执行结果”三部分,帮你彻底理解。

    案例1:只配置执行周期(EVERY),计算每10秒的温度最大值

    需求

    每20秒执行一次查询,计算过去20秒内4个传感器每10秒的温度最大值,结果存到带“temperature_max”后缀的序列里。

    创建语句

    — 创建名为cq1的连续查询,每20秒执行一次
    CREATE CONTINUOUS QUERY cq1
    RESAMPLE EVERY 20s — 执行周期:20秒
    BEGIN
    — 计算temperature的最大值
    SELECT max_value(temperature)
    — 结果写入4个目标序列,路径前缀和原序列一致,后缀加temperature_max
    INTO root.ln.wf02.wt02(temperature_max),
    root.ln.wf02.wt01(temperature_max),
    root.ln.wf01.wt02(temperature_max),
    root.ln.wf01.wt01(temperature_max)
    — 数据来源:所有以root.ln开头的传感器(这里匹配4个温度传感器)
    FROM root.ln.*.*
    — 按10秒分组计算,每个分组的最大值就是该时间段的结果
    GROUP BY(10s)
    END

    执行逻辑解释
    • 执行周期:每20秒一次,所以第一次执行时间假设是22:18:40,第二次是22:19:00,以此类推。
    • 时间窗口:因为没指定RANGE,默认RANGE等于EVERY(20秒),所以每次查询的窗口是“当前时间-20秒”到“当前时间”。比如22:18:40执行时,窗口是[22:18:20, 22:18:40);22:19:00执行时,窗口是[22:18:40, 22:19:00)。
    • 分组计算:按10秒分组,所以每个窗口会分成2个组(20秒/10秒=2组),每个组计算一次最大值。
    执行结果
  • 22:18:40执行时,窗口是[22:18:20, 22:18:40),分成[22:18:20, 22:18:30)和[22:18:30, 22:18:40)两组:

    时间root.ln.wf02.wt02.temperature_maxroot.ln.wf02.wt01.temperature_maxroot.ln.wf01.wt02.temperature_maxroot.ln.wf01.wt01.temperature_max
    2021-05-11T22:18:20.000+08:00 122.0(22:18:24的122.0是该组最大值) 45.0(22:18:24的45.0) 59.0(22:18:29的59.0) 181.0(22:18:29的181.0)
    2021-05-11T22:18:30.000+08:00 182.0(22:18:34的182.0) 113.0(22:18:34的113.0) 52.0(22:18:39的52.0) 180.0(22:18:34的180.0)
  • 22:19:00执行时,窗口是[22:18:40, 22:19:00),分成[22:18:40, 22:18:50)和[22:18:50, 22:19:00)两组:

    时间root.ln.wf02.wt02.temperature_maxroot.ln.wf02.wt01.temperature_maxroot.ln.wf01.wt02.temperature_maxroot.ln.wf01.wt01.temperature_max
    2021-05-11T22:18:40.000+08:00 137.0(22:18:49的137.0) 172.0(22:18:49的172.0) 135.0(22:18:49的135.0) 193.0(22:18:49的193.0)
    2021-05-11T22:18:50.000+08:00 16.0(22:18:55的16.0) 124.0(22:18:55的124.0) 183.0(22:18:55的183.0) 18.0(22:18:55的18.0)
  • 最终查询结果(执行两次后的所有结果):

  • — 查询所有temperature_max序列的数据
    SELECT temperature_max from root.ln.*.*;

    返回结果:

    时间root.ln.wf02.wt02.temperature_maxroot.ln.wf02.wt01.temperature_maxroot.ln.wf01.wt02.temperature_maxroot.ln.wf01.wt01.temperature_max
    2021-05-11T22:18:20.000+08:00 122.0 45.0 59.0 181.0
    2021-05-11T22:18:30.000+08:00 182.0 113.0 52.0 180.0
    2021-05-11T22:18:40.000+08:00 137.0 172.0 135.0 193.0
    2021-05-11T22:18:50.000+08:00 16.0 124.0 183.0 18.0

    案例2:只配置时间窗口(RANGE),每10秒计算过去40秒的最大值

    需求

    不指定执行周期(默认等于GROUP BY的间隔),每10秒执行一次,计算过去40秒内4个传感器每10秒的温度最大值,结果存到带“temperature_max”后缀的序列里。

    创建语句

    — 创建名为cq2的连续查询,配置40秒的时间窗口
    CREATE CONTINUOUS QUERY cq2
    RESAMPLE RANGE 40s — 时间窗口:过去40秒到当前时间
    BEGIN
    SELECT max_value(temperature)
    INTO root.ln.wf02.wt02(temperature_max),
    root.ln.wf02.wt01(temperature_max),
    root.ln.wf01.wt02(temperature_max),
    root.ln.wf01.wt01(temperature_max)
    FROM root.ln.*.*
    GROUP BY(10s) — 分组间隔10秒,执行周期默认等于10秒
    END

    执行逻辑解释
    • 执行周期:没指定EVERY,默认等于GROUP BY的10秒,所以第一次执行22:18:40,第二次22:18:50,第三次22:19:00。
    • 时间窗口:RANGE=40秒,所以每次查询的窗口是“当前时间-40秒”到“当前时间”。比如22:18:40执行时,窗口是[22:18:00, 22:18:40);22:18:50执行时,窗口是[22:18:10, 22:18:50)。
    • 空值处理:IoTDB会自动跳过全是NULL的行。比如22:18:40执行时,[22:18:00, 22:18:10)这个时间段没有数据,结果是NULL,就不会写入目标序列。
    执行结果
  • 22:18:40执行,窗口[22:18:00, 22:18:40),分成4个10秒组,其中[22:18:00, 22:18:10)是NULL,不写入:

    时间root.ln.wf02.wt02.temperature_maxroot.ln.wf02.wt01.temperature_maxroot.ln.wf01.wt02.temperature_maxroot.ln.wf01.wt01.temperature_max
    2021-05-11T22:18:10.000+08:00 121.0(22:18:14的121.0) 72.0(22:18:14的72.0) 183.0(22:18:14的183.0) 115.0(22:18:14的115.0)
    2021-05-11T22:18:20.000+08:00 122.0 45.0 59.0 181.0
    2021-05-11T22:18:30.000+08:00 182.0 113.0 52.0 180.0
  • 22:18:50执行,窗口[22:18:10, 22:18:50),分成4个组:

    时间root.ln.wf02.wt02.temperature_maxroot.ln.wf02.wt01.temperature_maxroot.ln.wf01.wt02.temperature_maxroot.ln.wf01.wt01.temperature_max
    2021-05-11T22:18:10.000+08:00 121.0 72.0 183.0 115.0
    2021-05-11T22:18:20.000+08:00 122.0 45.0 59.0 181.0
    2021-05-11T22:18:30.000+08:00 182.0 113.0 52.0 180.0
    2021-05-11T22:18:40.000+08:00 137.0 172.0 135.0 193.0
  • 最终查询结果(执行三次后的所有非NULL数据):

  • SELECT temperature_max from root.ln.*.*;

    返回结果:

    时间root.ln.wf02.wt02.temperature_maxroot.ln.wf02.wt01.temperature_maxroot.ln.wf01.wt02.temperature_maxroot.ln.wf01.wt01.temperature_max
    2021-05-11T22:18:10.000+08:00 121.0 72.0 183.0 115.0
    2021-05-11T22:18:20.000+08:00 122.0 45.0 59.0 181.0
    2021-05-11T22:18:30.000+08:00 182.0 113.0 52.0 180.0
    2021-05-11T22:18:40.000+08:00 137.0 172.0 135.0 193.0
    2021-05-11T22:18:50.000+08:00 16.0 124.0 183.0 18.0

    这里要注意:因为窗口有重叠(每次窗口40秒,执行周期10秒),所以某些时间段(比如[22:18:10, 22:18:20))会被多次计算,但最终只会保留一份结果,不会重复写入。

    案例3:同时配置执行周期和时间窗口,并用常量填充空值

    需求

    每20秒执行一次,计算过去40秒内4个传感器每10秒的温度最大值,没有数据的时间段用100.0填充,结果存到带“temperature_max”后缀的序列里。

    创建语句

    — 创建名为cq3的连续查询,同时指定EVERY和RANGE,加FILL填充空值
    CREATE CONTINUOUS QUERY cq3
    RESAMPLE EVERY 20s RANGE 40s — 每20秒执行一次,窗口是过去40秒
    BEGIN
    SELECT max_value(temperature)
    INTO root.ln.wf02.wt02(temperature_max),
    root.ln.wf02.wt01(temperature_max),
    root.ln.wf01.wt02(temperature_max),
    root.ln.wf01.wt01(temperature_max)
    FROM root.ln.*.*
    GROUP BY(10s)
    FILL(100.0) — 空值用100.0填充
    END

    执行逻辑解释
    • 执行周期20秒,时间窗口40秒,所以22:18:40执行时窗口[22:18:00, 22:18:40),22:19:00执行时窗口[22:18:20, 22:19:00)。
    • 和案例2的区别是加了FILL(100.0),所以没有数据的时间段不会显示NULL,而是用100.0填充,并且会写入目标序列(不会跳过)。
    执行结果
  • 22:18:40执行,窗口[22:18:00, 22:18:40),[22:18:00, 22:18:10)用100.0填充:

    时间root.ln.wf02.wt02.temperature_maxroot.ln.wf02.wt01.temperature_maxroot.ln.wf01.wt02.temperature_maxroot.ln.wf01.wt01.temperature_max
    2021-05-11T22:18:00.000+08:00 100.0 100.0 100.0 100.0
    2021-05-11T22:18:10.000+08:00 121.0 72.0 183.0 115.0
    2021-05-11T22:18:20.000+08:00 122.0 45.0 59.0 181.0
    2021-05-11T22:18:30.000+08:00 182.0 113.0 52.0 180.0
  • 22:19:00执行,窗口[22:18:20, 22:19:00),无空值:

    时间root.ln.wf02.wt02.temperature_maxroot.ln.wf02.wt01.temperature_maxroot.ln.wf01.wt02.temperature_maxroot.ln.wf01.wt01.temperature_max
    2021-05-11T22:18:20.000+08:00 122.0 45.0 59.0 181.0
    2021-05-11T22:18:30.000+08:00 182.0 113.0 52.0 180.0
    2021-05-11T22:18:40.000+08:00 137.0 172.0 135.0 193.0
    2021-05-11T22:18:50.000+08:00 16.0 124.0 183.0 18.0
  • 最终查询结果:

  • SELECT temperature_max from root.ln.*.*;

    返回结果包含了填充的100.0:

    时间root.ln.wf02.wt02.temperature_maxroot.ln.wf02.wt01.temperature_maxroot.ln.wf01.wt02.temperature_maxroot.ln.wf01.wt01.temperature_max
    2021-05-11T22:18:00.000+08:00 100.0 100.0 100.0 100.0
    2021-05-11T22:18:10.000+08:00 121.0 72.0 183.0 115.0
    2021-05-11T22:18:20.000+08:00 122.0 45.0 59.0 181.0
    2021-05-11T22:18:30.000+08:00 182.0 113.0 52.0 180.0
    2021-05-11T22:18:40.000+08:00 137.0 172.0 135.0 193.0
    2021-05-11T22:18:50.000+08:00 16.0 124.0 183.0 18.0

    案例4:配置时间窗口的开始和结束偏移(RANGE两个参数)

    需求

    每20秒执行一次,计算“过去40秒到过去20秒”这个时间段内,4个传感器每10秒的温度最大值,空值用100.0填充,结果存到带“temperature_max”后缀的序列里。

    创建语句

    — 创建名为cq4的连续查询,RANGE指定两个参数:start_time_offset=40s,end_time_offset=20s
    CREATE CONTINUOUS QUERY cq4
    RESAMPLE EVERY 20s RANGE 40s, 20s — 窗口:now()-40s 到 now()-20s
    BEGIN
    SELECT max_value(temperature)
    INTO root.ln.wf02.wt02(temperature_max),
    root.ln.wf02.wt01(temperature_max),
    root.ln.wf01.wt02(temperature_max),
    root.ln.wf01.wt01(temperature_max)
    FROM root.ln.*.*
    GROUP BY(10s)
    FILL(100.0)
    END

    执行逻辑解释

    这个案例的关键是RANGE的两个参数,它定义的窗口不是“到当前时间”,而是“到过去20秒”,相当于计算“稍早一点的历史数据”,避免处理刚产生的、可能还没稳定的数据。

    • 22:18:40执行时,窗口是now()-40s=22:18:00到now()-20s=22:18:20,即[22:18:00, 22:18:20)。
    • 22:19:00执行时,窗口是22:18:20到22:18:40,即[22:18:20, 22:18:40)。
    • 每个时间段只会被计算一次,不会像案例2、3那样重复计算,因为窗口之间没有重叠。
    执行结果
  • 22:18:40执行,窗口[22:18:00, 22:18:20),分成2个组:

    时间root.ln.wf02.wt02.temperature_maxroot.ln.wf02.wt01.temperature_maxroot.ln.wf01.wt02.temperature_maxroot.ln.wf01.wt01.temperature_max
    2021-05-11T22:18:00.000+08:00 100.0 100.0 100.0 100.0
    2021-05-11T22:18:10.000+08:00 121.0 72.0 183.0 115.0
  • 22:19:00执行,窗口[22:18:20, 22:18:40),分成2个组:

    时间root.ln.wf02.wt02.temperature_maxroot.ln.wf02.wt01.temperature_maxroot.ln.wf01.wt02.temperature_maxroot.ln.wf01.wt01.temperature_max
    2021-05-11T22:18:20.000+08:00 122.0 45.0 59.0 181.0
    2021-05-11T22:18:30.000+08:00 182.0 113.0 52.0 180.0
  • 最终查询结果:

  • SELECT temperature_max from root.ln.*.*;

    返回结果:

    时间root.ln.wf02.wt02.temperature_maxroot.ln.wf02.wt01.temperature_maxroot.ln.wf01.wt02.temperature_maxroot.ln.wf01.wt01.temperature_max
    2021-05-11T22:18:00.000+08:00 100.0 100.0 100.0 100.0
    2021-05-11T22:18:10.000+08:00 121.0 72.0 183.0 115.0
    2021-05-11T22:18:20.000+08:00 122.0 45.0 59.0 181.0
    2021-05-11T22:18:30.000+08:00 182.0 113.0 52.0 180.0

    案例5:无GROUP BY子句,实现数据预处理(温度值+1)

    需求

    每20秒执行一次,将所有以root.ln开头的温度传感器数据加1(预处理),然后存到root.precalculated_sg这个新的数据库下,保持设备路径不变。

    创建语句

    — 创建名为cq5的连续查询,无GROUP BY,显式指定EVERY
    CREATE CONTINUOUS QUERY cq5
    RESAMPLE EVERY 20s — 必须显式指定EVERY,因为没有GROUP BY
    BEGIN
    — 对temperature值做预处理:加1
    SELECT temperature + 1
    — 结果写入root.precalculated_sg,设备路径和原序列一致(用::(temperature)匹配原序列的指标名)
    INTO root.precalculated_sg.::(temperature)
    — 数据来源:所有以root.ln开头的传感器
    FROM root.ln.*.*
    — 按设备对齐结果,方便查看每个设备的预处理后数据
    align by device
    END

    执行逻辑解释

    这个案例没有GROUP BY子句,属于“行级预处理”,不是统计计算。核心是将原始数据做简单转换后存储,适合提前处理需要频繁使用的中间数据。

    • 执行周期20秒,窗口是过去20秒到当前时间,每次执行会把窗口内的所有原始数据加1后存到新序列。
    • INTO root.precalculated_sg.::(temperature)中的::是通配符,表示“继承原序列的设备路径”。比如原序列root.ln.wf02.wt02.temperature,目标序列就是root.precalculated_sg.wf02.wt02.temperature。
    • ALIGN BY DEVICE让结果按设备分组显示,每个设备的所有数据在同一部分,更易读。
    执行结果
  • 22:18:40执行,窗口[22:18:20, 22:18:40),处理该时间段内的4个传感器数据(共4个设备,每个设备2-3条数据):
  • — 查询预处理后的结果,按设备对齐
    SELECT temperature from root.precalculated_sg.*.* align by device;

    返回部分结果(只展示前几个设备的几条数据):

    时间Devicetemperature
    2021-05-11T22:18:24.949+08:00 root.precalculated_sg.wf02.wt02 123.0(原122.0+1)
    2021-05-11T22:18:29.967+08:00 root.precalculated_sg.wf02.wt02 48.0(原47.0+1)
    2021-05-11T22:18:34.979+08:00 root.precalculated_sg.wf02.wt02 183.0(原182.0+1)
    2021-05-11T22:18:24.949+08:00 root.precalculated_sg.wf02.wt01 46.0(原45.0+1)
    2021-05-11T22:18:29.967+08:00 root.precalculated_sg.wf02.wt01 15.0(原14.0+1)

    (完整结果包含该窗口内所有4个设备的8条数据,每条数据都加1)

    2. 22:19:00执行,窗口[22:18:40, 22:19:00),处理该时间段内的12条数据(4个设备,每个设备3条数据),结果同样会自动写入`root.precalculated_sg`下的对应序列。

    ## 四、连续查询的管理:查询、删除、修改
    创建完连续查询后,还需要知道怎么管理它——比如查看当前有哪些CQ在运行、删除不需要的CQ、修改已有的CQ。下面逐个讲解对应的操作。

    ### 1. 查询系统中已有的连续查询
    如果你想知道当前IoTDB集群里有哪些连续查询,每个CQ的状态和具体语句,可以用`SHOW CONTINUOUS QUERIES`(或简写`SHOW CQS`)命令。

    #### 语法
    ```sql
    SHOW (CONTINUOUS QUERIES | CQS);

    示例

    假设我们已经创建了前面的cq1到cq5,执行查询命令:

    SHOW CONTINUOUS QUERIES;

    返回结果

    结果会按cq_id排序,包含3列:cq_id(CQ名称)、query(CQ的创建语句)、state(CQ的状态,active表示正在运行)。

    cq_idquerystate
    cq1 CREATE CONTINUOUS QUERY cq1 RESAMPLE EVERY 20s BEGIN SELECT max_value(temperature) INTO root.ln.wf02.wt02(temperature_max), root.ln.wf02.wt01(temperature_max), root.ln.wf01.wt02(temperature_max), root.ln.wf01.wt01(temperature_max) FROM root.ln.. GROUP BY(10s) END active
    cq2 CREATE CONTINUOUS QUERY cq2 RESAMPLE RANGE 40s BEGIN SELECT max_value(temperature) INTO …(省略部分内容)… GROUP BY(10s) END active
    …(cq3、cq4、cq5的创建语句)… active

    2. 删除已有的连续查询

    当某个连续查询不再需要时(比如业务需求变更,不需要再计算每10秒的温度最大值),可以用DROP CONTINUOUS QUERY(或简写DROP CQ)命令删除它。

    语法

    DROP (CONTINUOUS QUERY | CQ) <cq_id>;

    注意事项
    • 删除CQ后,它不会再自动执行,但之前已经写入目标序列的结果不会被删除(如果需要删除结果,需要手动删除目标序列的数据)。
    • DROP CQ命令执行成功后,不会返回任何结果集;如果指定的cq_id不存在,IoTDB会抛出异常,提示“CQ not found”。
    示例

    删除前面创建的cq1:

    DROP CONTINUOUS QUERY cq1;

    执行后,再用SHOW CQS查询,就看不到cq1了。

    3. 修改已有的连续查询

    这里要注意一个关键点:IoTDB目前不支持直接修改已有的连续查询。比如你想把cq2的执行周期从10秒改成15秒,不能用ALTER之类的命令直接改,只能用“先删除、再重新创建”的方式。

    操作步骤
  • 先删除原来的CQ:用DROP CQ <旧cq_id>命令删除不需要的CQ。
  • 再创建新的CQ:用CREATE CONTINUOUS QUERY <新cq_id>命令创建修改后的CQ(如果不需要改名称,可以用原来的cq_id)。
  • 示例

    将cq2的执行周期从默认的10秒改成15秒:

    — 步骤1:删除原来的cq2
    DROP CQ cq2;

    — 步骤2:创建新的cq2,执行周期改为15秒
    CREATE CONTINUOUS QUERY cq2
    RESAMPLE EVERY 15s RANGE 40s — 执行周期改成15秒
    BEGIN
    SELECT max_value(temperature)
    INTO root.ln.wf02.wt02(temperature_max),
    root.ln.wf02.wt01(temperature_max),
    root.ln.wf01.wt02(temperature_max),
    root.ln.wf01.wt01(temperature_max)
    FROM root.ln.*.*
    GROUP BY(10s)
    END

    这样就实现了“修改”cq2的效果。

    五、连续查询的适用场景:什么时候该用CQ?

    很多人学完语法和案例后,还是不知道什么时候该用连续查询。其实CQ的核心价值是“自动化周期性计算”和“预存储结果”,下面三个场景是它最常用的地方,帮你判断是否需要用CQ。

    场景1:数据降采样+不同保留策略

    在工业场景中,传感器通常会高频采集数据(比如每秒采集1000个点),但这些原始数据不需要长期保存——比如实时监控只需要看最近1天的高频数据,而月度报表只需要看每小时的低频数据。如果把所有高频数据都长期保存,会占用大量磁盘空间。

    这时候可以用连续查询做降采样:

  • 用CQ定期将高频原始数据(如每秒1000点)降采样成低频数据(如每小时1个点,取该小时的平均值或最大值)。
  • 将降采样后的低频数据存到另一个数据库(比如root.low_freq),并给这个数据库设置较长的TTL(Time To Live,数据存活时间),比如1个月;而原始高频数据所在的数据库(比如root.high_freq)设置较短的TTL,比如1天。
  • 这样既能满足实时监控的高频需求,又能满足长期分析的低频需求,还能节省磁盘空间,避免存储冗余数据。

    场景2:预计算代价昂贵的查询

    有些查询计算起来很耗时,比如“计算过去1年每个设备每天的平均运行时长”,如果每次需要这个数据都重新计算,会消耗大量CPU和IO资源,查询响应时间也会很长。

    这时候可以用连续查询预计算:

  • 创建一个CQ,周期性地计算这个“每日平均运行时长”(比如每天凌晨2点,计算前一天的平均值)。
  • 将计算结果存到目标序列(比如root.precomputed.device1.daily_avg_runtime)。
  • 后续需要这个数据时,直接查询目标序列,而不是重新计算,查询时间会从几秒甚至几分钟缩短到毫秒级。
  • 这种方式特别适合可视化工具(比如Grafana)渲染时序图——如果每次渲染都要计算大量数据,图表加载会很慢,用预计算的结果就能秒级加载。

    场景3:替代不支持的子查询

    目前IoTDB还不支持子查询,比如你想执行“计算过去1个月内,每天的温度最大值的平均值”,直接写SELECT AVG(day_max) FROM (SELECT MAX(temperature) AS day_max FROM … GROUP BY(1d))是会报错的。

    这时候可以用连续查询模拟子查询,分两步实现:

  • 第一步:创建CQ计算每日最大值(子查询部分)。— 创建CQ,每天计算一次当天的温度最大值,存到day_max序列
    CREATE CONTINUOUS QUERY cq_day_max
    RESAMPLE EVERY 1d — 每天执行一次
    BEGIN
    SELECT max_value(temperature) AS day_max
    INTO root.sensor.temp_day_max(day_max)
    FROM root.sensor.temperature
    GROUP BY(1d) — 按天分组,计算每天的最大值
    END
  • 第二步:查询CQ的结果,计算平均值(外层查询部分)。— 直接查询预计算好的day_max序列,计算过去1个月的平均值
    SELECT AVG(day_max)
    FROM root.sensor.temp_day_max
    WHERE time > now() 30d; — 时间范围:过去30天
  • 通过这种方式,就能实现原本需要子查询才能完成的功能。

    六、连续查询的配置参数:优化性能

    IoTDB有两个和连续查询相关的配置参数,在iotdb-system.properties文件中设置,主要用来控制CQ的执行线程和最小执行周期,合理配置能优化CQ的性能。

    1. continuous_query_submit_thread_count

    作用

    这个参数控制用于周期性提交连续查询执行任务的线程数量。简单说,就是有多少个线程负责把CQ的执行任务分配给worker线程去执行。

    类型

    int32(整数)

    默认值

    2

    调整建议
    • 如果你的IoTDB集群里有很多连续查询(比如上百个),默认的2个线程可能不够用,会导致CQ任务提交延迟,这时候可以适当增大这个值,比如改成4或8。
    • 如果CQ数量很少(比如只有几个),保持默认值2就够了,没必要增大,避免线程过多导致资源浪费。

    2. continuous_query_min_every_interval_in_ms

    作用

    这个参数指定系统允许的连续查询最小周期性时间间隔,单位是毫秒(ms)。也就是说,你创建CQ时指定的EVERY参数值,不能小于这个配置值,否则会创建失败。

    类型

    duration(时间长度)

    默认值

    1000(即1秒)

    调整建议
    • 默认情况下,CQ的最小执行周期是1秒,如果你需要更频繁的执行(比如每隔500毫秒执行一次),可以把这个值改小,比如改成500。
    • 但要注意:执行周期越小,CQ的执行频率越高,对系统资源(CPU、IO)的消耗也越大。如果不是必须,不建议把这个值设得太小,避免影响IoTDB的其他功能。
    配置示例

    如果你想把最小执行周期改成500毫秒,在iotdb-system.properties文件中添加或修改如下配置:

    continuous_query_min_every_interval_in_ms=500

    修改后需要重启IoTDB才能生效。

    七、总结:学好CQ,让实时数据计算更高效

    到这里,IoTDB连续查询的所有核心内容就讲完了。我们再来回顾一下:

  • 是什么:连续查询是周期性自动执行的查询,能实现滑动窗口流式计算,自动存储结果。
  • 怎么用:掌握语法规则(尤其是RESAMPLE、GROUP BY、FILL参数),结合实战案例理解不同场景的配置方式。
  • 怎么管:用SHOW查CQ、用DROP删CQ、用“删了再建”的方式改CQ。
  • 什么时候用:数据降采样、预计算昂贵查询、替代子查询这三个场景最常用。
  • 怎么优化:通过配置参数调整线程数和最小执行周期,适配你的业务需求。
  • 希望这篇文章能帮你彻底搞懂连续查询,在实际项目中用它来简化实时数据计算的流程,提高工作效率。如果在使用过程中遇到问题,可以再回头看看对应的案例和注意事项,也欢迎在评论区交流你的使用经验!

    🌐 附:IoTDB的各大版本

    📄 Apache IoTDB 是一款工业物联网时序数据库管理系统,采用端边云协同的轻量化架构,支持一体化的物联网时序数据收集、存储、管理与分析 ,具有多协议兼容、超高压缩比、高通量读写、工业级稳定、极简运维等特点。

    版本IoTDB 二进制包IoTDB 源代码发布说明
    2.0.5 – All-in-one- AINode- SHA512- ASC – 源代码- SHA512- ASC release notes
    1.3.5 – All-in-one- AINode- SHA512- ASC – 源代码- SHA512- ASC release notes
    0.13.4 – All-in-one- Grafana 连接器- Grafana 插件- SHA512- ASC – 源代码- SHA512- ASC release notes

    ✨ 去获取:https://archive.apache.org/dist/iotdb/

    联系博主

        xcLeigh 博主全栈领域优质创作者,博客专家,目前,活跃在CSDN、微信公众号、小红书、知乎、掘金、快手、思否、微博、51CTO、B站、腾讯云开发者社区、阿里云开发者社区等平台,全网拥有几十万的粉丝,全网统一IP为 xcLeigh。希望通过我的分享,让大家能在喜悦的情况下收获到有用的知识。主要分享编程、开发工具、算法、技术学习心得等内容。很多读者评价他的文章简洁易懂,尤其对于一些复杂的技术话题,他能通过通俗的语言来解释,帮助初学者更好地理解。博客通常也会涉及一些实践经验,项目分享以及解决实际开发中遇到的问题。如果你是开发领域的初学者,或者在学习一些新的编程语言或框架,关注他的文章对你有很大帮助。

        亲爱的朋友,无论前路如何漫长与崎岖,都请怀揣梦想的火种,因为在生活的广袤星空中,总有一颗属于你的璀璨星辰在熠熠生辉,静候你抵达。

         愿你在这纷繁世间,能时常收获微小而确定的幸福,如春日微风轻拂面庞,所有的疲惫与烦恼都能被温柔以待,内心永远充盈着安宁与慰藉。

        至此,文章已至尾声,而您的故事仍在续写,不知您对文中所叙有何独特见解?期待您在心中与我对话,开启思想的新交流。


         💞 关注博主 🌀 带你实现畅游前后端!

         🏰 大屏可视化 🌀 带你体验酷炫大屏!

         💯 神秘个人简介 🌀 带你体验不一样得介绍!

         🥇 从零到一学习Python 🌀 带你玩转Python技术流!

         🏆 前沿应用深度测评 🌀 前沿AI产品热门应用在线等你来发掘!

         💦 注:本文撰写于CSDN平台,作者:xcLeigh(所有权归作者所有) ,https://xcleigh.blog.csdn.net/,如果相关下载没有跳转,请查看这个地址,相关链接没有跳转,皆是抄袭本文,转载请备注本文原地址。


    在这里插入图片描述

         📣 亲,码字不易,动动小手,欢迎 点赞 ➕ 收藏,如 🈶 问题请留言(或者关注下方公众号,看见后第一时间回复,还有海量编程资料等你来领!),博主看见后一定及时给您答复 💌💌💌

    赞(0)
    未经允许不得转载:171主机测评 » Apache IoTDB 连续查询(CQ)全解析:从语法到实战,手把手教你玩转实时数据计算
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址