流计算结果不符合预期

【TDengine 使用环境】
测试

【TDengine 版本】

3.4.1.0

【操作系统以及版本】

centos

【部署方式】容器

创建了一个流计算。语句如下 CREATE STREAM IF NOT EXISTS nbq_pr

    **INTERVAL**(1d) SLIDING(1d)

    **FROM** product_nbq

    **PARTITION** **BY** device_id

    **INTO** stream_nbq_pr

    OUTPUT_SUBTABLE(CONCAT('stream_nbq_pr\_', **CAST**(device_id **AS** **VARCHAR**)))

    **AS**

    **SELECT**

        \_twend **AS** ts,

        LAST_ROW(ljfdl) - **FIRST**(ljfdl) **AS** ac_energy,

        TWA(zlgl) \* 24 **AS** dc_energy

    **FROM** product_nbq

    **WHERE** \_c0 >= \_twstart **AND** \_c0 <= \_twend

    **PARTITION** **BY** device_id

结果发现各个子表的结果数据都是一样的。

问题分析

你在流计算的 两个位置 都写了 PARTITION BY device_id,这是导致所有子表结果相同的原因。

你的写法:

FROM product_nbq PARTITION BY device_id     ← ① 流级别分区(控制输出子表)
INTO stream_nbq_pr OUTPUT_SUBTABLE(...)
AS
SELECT ... FROM product_nbq
WHERE _c0 >= _twstart AND _c0 <= _twend
PARTITION BY device_id                       ← ② 子查询内分区(控制计算分组)

根因

在 TDengine 流计算中,PARTITION BY 只能出现在流级别FROM trigger_table PARTITION BY ... INTO ...),不应当出现在 AS 后的子查询中。

流级别的 PARTITION BY 已经同时做了两件事:

  1. 计算分组:每个 device_id 独立开窗口、独立计算聚合函数
  2. 路由输出:每个分区的结果写入对应的输出子表

当子查询内又写了一次 PARTITION BY device_id 时,流计算的执行引擎会混淆分区作用域。实际结果是:子查询的 PARTITION BY 被忽略或执行异常,导致窗口函数(LAST_ROWFIRSTTWA)都在全量数据上计算,然后流级别再把这个相同的全量计算结果分别路由到各个子表。所以所有子表的数据都一样。

尝试以下语句:

删除子查询中的 PARTITION BY device_id,同时建议补充 STREAM_OPTIONS(FORCE_OUTPUT)

CREATE STREAM IF NOT EXISTS nbq_pr
    INTERVAL(1d) SLIDING(1d)
    FROM product_nbq
    PARTITION BY device_id
    STREAM_OPTIONS(FORCE_OUTPUT)                        -- 分区流推荐加此项
    INTO stream_nbq_pr
    OUTPUT_SUBTABLE(CONCAT('stream_nbq_pr\_', CAST(device_id AS VARCHAR)))
AS
SELECT
    _twend AS ts,
    LAST_ROW(ljfdl) - FIRST(ljfdl) AS ac_energy,
    TWA(zlgl) * 24 AS dc_energy
FROM product_nbq
WHERE _c0 >= _twstart AND _c0 < _twend;                 -- 推荐用 < 而非 <=