【流式计算】聚合结果全部子表数据都都是同一张子表问题啥情况

【TDengine 使用环境】
测试环境

【TDengine 版本】

3.4.1.13

【描述业务影响】

  • 业务期望:想把各个子表的按照【小时频率】进行统计聚合后写入聚合表中

SELECT
_wstart AS ts,
tbname AS monitor_common_code,
CAST(AVG(current_value) AS DECIMAL(16,4)) AS avg_value,
max(current_value) AS max_value,
min(current_value) AS min_value,
CAST(SUM(current_value) AS DECIMAL(16,4)) AS sum_value,
cols(max(current_value), ts) AS max_time,
cols(min(current_value), ts) AS min_time

FROM ops_monitor.st_monitor_data_realtime
WHERE ts >= ‘2026-01-01’ AND ts < ‘2026-12-01’
PARTITION BY tbname
INTERVAL(1h) SLIDING(1h);

  • 流式计算创建命令:

CREATE STREAM opsmonitor.stream_monitor_stat_hour
INTERVAL(1h) SLIDING(1h)
FROM ops_monitor.st_monitor_data_realtime
PARTITION BY tbname, station_id, line_id, device_id, energy_type, energy_cal_config_flag, monitor_type, monitor_type_code
STREAM_OPTIONS(
WATERMARK(1m) | FILL_HISTORY_FIRST(‘2026-01-01 00:00:00’)
)
INTO ops_monitor.st_monitor_stat_hour
OUTPUT_SUBTABLE(concat('hour
’, %%1))
TAGS(
tbname BINARY(128), AS %%1,
station_id BIGINT AS %%2,
line_id BIGINT AS %%3,
device_id BIGINT AS %%4,
energy_type INT AS %%5,
energy_cal_config_flag INT AS %%6,
monitor_type INT AS %%7,
monitor_type_code BINARY(32) AS %%8
)
AS SELECT
_twstart AS ts,
CAST(AVG(current_value) AS DECIMAL(16,4)) AS avg_value,
max(current_value) AS max_value,
min(current_value) AS min_value,
CAST(SUM(current_value) AS DECIMAL(16,4)) AS sum_value,
cols(max(current_value), ts) AS max_time,
cols(min(current_value), ts) AS min_time
FROM ops_monitor.st_monitor_data_realtime
WHERE 1=1 AND energy_type = 0 AND ts >= _twstart AND ts < _twend;

  • 问题描述: 同一时间窗口下聚合后的各子表 avg/sum/max/min/maxTime/minTime 值相同???
    1. 是流式计算命令哪里有问题么?

核心问题:子查询 FROM 写错了,结果会是"全局聚合"而不是"每子表聚合"

流计算里,触发只决定何时算,分组只决定写到哪个输出子表,而计算的数据范围完全由 AS 子查询自己决定。你的子查询写的是:

FROM ops_monitor.st_monitor_data_realtime
WHERE 1=1 AND energy_type = 0 AND ts >= _twstart AND ts < _twend

这等于每个分组触发时,都对整张超级表(energy_type=0 且在窗口内)做一次全局聚合,然后把完全相同的一行结果写进每个分组的输出子表。结果就是:所有 hourXXX 子表里的数据一模一样,根本不是"各个子表分别统计"。

官方推荐模式是 FROM %%tbname(当前分组对应的子表)+ 窗口边界过滤:

FROM %%tbname WHERE _c0 >= _twstart AND _c0 <= _twend —— 来自 best-practices 文档的窗口触发示例

%%tbname 在 PARTITION BY 含 tbname 时可用,正好你的分组里有。

修正后的建流 SQL

只改一处 FROM,其余语法均合法:

CREATE STREAM opsmonitor.stream_monitor_stat_hour
INTERVAL(1h) SLIDING(1h)
FROM ops_monitor.st_monitor_data_realtime
PARTITION BY tbname, station_id, line_id, device_id, energy_type, energy_cal_config_flag, monitor_type, monitor_type_code
STREAM_OPTIONS(
  WATERMARK(1m) | FILL_HISTORY_FIRST('2026-01-01 00:00:00')
)
INTO ops_monitor.st_monitor_stat_hour
OUTPUT_SUBTABLE(concat('hour', %%1))
TAGS(
  tbname BINARY(128) AS %%1,
  station_id BIGINT AS %%2,
  line_id BIGINT AS %%3,
  device_id BIGINT AS %%4,
  energy_type INT AS %%5,
  energy_cal_config_flag INT AS %%6,
  monitor_type INT AS %%7,
  monitor_type_code BINARY(32) AS %%8
)
AS SELECT
  _twstart AS ts,
  CAST(AVG(current_value) AS DECIMAL(16,4)) AS avg_value,
  max(current_value) AS max_value,
  min(current_value) AS min_value,
  CAST(SUM(current_value) AS DECIMAL(16,4)) AS sum_value,
  cols(max(current_value), ts) AS max_time,
  cols(min(current_value), ts) AS min_time
FROM %%tbname
WHERE energy_type = 0 AND ts >= _twstart AND ts < _twend;

逐项核对(都有文档依据)

  • STREAM_OPTIONS 里的 |:合法。语法定义为 STREAM_OPTIONS(stream_option [|...]),官方示例 STREAM_OPTIONS(MAX_DELAY(1m) | FILL_HISTORY_FIRST) 即用 | 分隔。FILL_HISTORY_FIRST('2026-01-01 00:00:00') 的 start_time 参数也支持,语义是优先回算历史再转实时,与你的回填需求一致。

  • PARTITION BY tbname, station_id, ... 多列混用子表+标签:合法。“指定触发的分组列,支持多列,目前只支持按照子表和标签进行分组。”

  • OUTPUT_SUBTABLE(concat('hour', %%1)):合法。tbname_expr 是任意输出字符串表达式,可用分组列占位符 %%n(n 为分组列下标,%%1=tbname);CONCAT 是标准标量函数(VARCHAR/NCHAR,2~8 参数)。注意表名超过上限会被截断(表名上限 192 字节),tbname 太长时可能产生重名冲突。