【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 已经同时做了两件事:
- 计算分组:每个
device_id 独立开窗口、独立计算聚合函数
- 路由输出:每个分区的结果写入对应的输出子表
当子查询内又写了一次 PARTITION BY device_id 时,流计算的执行引擎会混淆分区作用域。实际结果是:子查询的 PARTITION BY 被忽略或执行异常,导致窗口函数(LAST_ROW、FIRST、TWA)都在全量数据上计算,然后流级别再把这个相同的全量计算结果分别路由到各个子表。所以所有子表的数据都一样。
尝试以下语句:
删除子查询中的 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; -- 推荐用 < 而非 <=