流计算sql使用问题

【TDengine 使用环境】
预生产环境

【TDengine 版本】

server_version()
3.4.2.8.community

【操作系统以及版本】

【部署方式】容器/非容器部署

docker容器

【集群节点数】

【集群副本数】

【描述业务影响】

我想实现一个源超级表ems下,一个字段totalMileage的计算(根据这个累计值去相减形成日/月年表)源子表最好和流形成的子表同名,如果不允许那就子表名加统一后缀区分,目前我的问题是流创建的不对。一些注意点我注意了(字段使用大小写的区分) 还有一些我不清楚 是需要先创建好插入的表bms_mileage_day吗?请哪位老师给我一个正确语法的sql!!

【问题复现路径/shan】做过哪些操作出现的问题

【遇到的问题:问题现象及影响】

【资源配置】

【报错完整截图】(不要大段的粘贴报错代码,论坛直接看报错代码不直观)

taos> CREATE STREAM st_mileage_day

AS
SELECT _wstart AS day_start,
_wend AS day_end,
FIRST(totalMileage) AS start_mileage,
LAST(totalMileage) AS end_mileage,
LAST(totalMileage) - FIRST(totalMileage) AS day_mileage,
COUNT(*) AS report_cnt,
COUNT(totalMileage) AS valid_cnt
FROM bms
PARTITION BY tbname
INTERVAL(1d)
INTO bms_mileage_day;

CREATE STREAM IF NOT EXISTS st_mileage_day INTERVAL(1d)
FROM
bms
PARTITION BY
tbname INTO bms_mileage_day AS
SELECT _wstart AS day_start,
_wend AS day_end,
FIRST(totalMileage) AS start_mileage,
LAST(totalMileage) AS end_mileage,
LAST(totalMileage) - FIRST(totalMileage) AS day_mileage
FROM %%tbname; 到底怎么写啊 哪些老师指定下

追加补充 我的源超级表bms 有tag----------- TAGS (bmsseq VARCHAR(50)) 哪位老师能帮忙指导下totalMileage这个字段的日月年计算

这个还是 3.3.6.x版本的建流语法,3.4.x 版本的建流语法已经不是这样的了。

是的 老师 能帮我写一下参考吗 语法有点看不懂 论坛也没有什么可参考的

一个字段totalMileage的计算: 但这个没有看明白,需要怎么计算?

源超级表ems 会有很多子表(需要计算的字段totalMileage) 我需要根据这个源去算出来这个累加的值的日月年增量 输出到目的表 比如 totalMileage 数据 从1 涨到 20 100 这样 我需要interver 1d 1n 1y这样去last-fist 算出增量

之前的低等级3.3.6.13版本搞出来了一个流 但是目标子表名会追加系统的编号导致我的目标子表名不可控 所以我升级了版本 但是现在流都创建不出来了 :sweat_smile:

好。是对每个子表都做同样的计算吧?输出的结果表就是 保存 窗口的开始时间、窗口的结束时间、totalMileage 在窗口的开始值、totalMileage在窗口的最后值,totalMileage的最后值-开始值,就这些吧?
另外,结果表有tag 字段吗?如果有的话,需要怎么设置tag的值呢?

源子表最好和流形成的子表同名,如果不允许那就流子表名加统一后缀区分 源bms表有tag 如果能带到流子表最好 没有不设置也行 我根据子表名去查数据即可

流计算的输出指标其实只需要2个核心的 1日/月/年(日月年会建3个流形成3个目标超表)时间 2日增量(totalMileage)

CREATE STREAM IF NOT EXISTS st_mileage_day
TRIGGER WINDOW_CLOSE
WATERMARK 3m
IGNORE EXPIRED 1
FILL_HISTORY 1
INTO bms_mileage_day
AS SELECT
_wstart AS day_start,
_wend AS day_end,
FIRST(totalMileage) AS start_mileage,
LAST(totalMileage) AS end_mileage,
CASE
WHEN LAST(totalMileage) IS NULL OR FIRST(totalMileage) IS NULL
THEN NULL
WHEN LAST(totalMileage) >= FIRST(totalMileage)
THEN LAST(totalMileage) - FIRST(totalMileage)
ELSE LAST(totalMileage)
END AS day_mileage,
COUNT(*) AS report_cnt,
COUNT(totalMileage) AS valid_cnt
FROM bms
PARTITION BY tbname, bmsseq
INTERVAL(1d); 这个是我之前3.3.6.13版本的流sql 这个就是在运算日表的

CREATE STREAM IF NOT EXISTS st_mileage_day interval(1d) sliding(1d)
FROM bms PARTITION BY tbname
INTO bms_mileage_day
OUTPUT_SUBTABLE(CONCAT(‘st_mileage_day_’, tbname))
AS
SELECT _twstart AS day_start,
_twend AS day_end,
FIRST(totalMileage) AS start_mileage,
LAST(totalMileage) AS end_mileage,
LAST(totalMileage) - FIRST(totalMileage) AS day_mileage
from %%tbname;

你试试这个吧。这个是3.4版本的建流语法。

输出的流结果子表的名称是: OUTPUT_SUBTABLE(CONCAT(‘st_mileage_day_’, tbname)) 来控制的。这里是在源子表基础上加了一个 前缀:st_mileage_day_ , 你也可以安装这种模式改成后缀。

CREATE STREAM IF NOT EXISTS st_mileage_day interval(1d) sliding(1d)
FROM bms PARTITION BY tbname
INTO bms_mileage_day
OUTPUT_SUBTABLE(CONCAT(“st_mileage_day_”, tbname))
AS
SELECT _twstart AS day_start,
_twend AS day_end,
FIRST(totalMileage) AS start_mileage,
LAST(totalMileage) AS end_mileage,
LAST(totalMileage) - FIRST(totalMileage) AS day_mileage
from %%tbname; 当前已经按照这个创建成功流了 但是没有生成数据出来

我不太理解您这个sliding(1d) 如果我是月和年 咋办? 这个机制您能给我讲一下吗 能否支持成之前流的那种实时计算比如 WATERMARK 3m
IGNORE EXPIRED 1
FILL_HISTORY 1 并且支持历史数据等选项啊

我追加了几条新数据 也没有输出 当前流只产生了目的超表 没有子表记录 是这个窗口有关系吗?一天结束才触发计算啊? 能改成实时触发计算吗

对啊,你这个是1天的滑动窗口。你可以写入超过1天的数据啊,比如:
2026-09-01 08:00:00.000 …
2026-09-01 09:00:00.000 …
2026-09-02 08:00:00.000 …
2026-09-02 09:00:00.000 …
2026-09-03 08:00:00.000 …
写入这几条记录,会输出2条结果记录。

WATERMARK 3m
IGNORE EXPIRED 1
FILL_HISTORY 1
这些选项都有啊,

[STREAM_OPTIONS(stream_option [|…])]

stream_option: {WATERMARK(duration_time) | EXPIRED_TIME(exp_time) | IGNORE_DISORDER | DELETE_RECALC | DELETE_OUTPUT_TABLE | FILL_HISTORY[(start_time)] | FILL_HISTORY_FIRST[(start_time)] | CALC_NOTIFY_ONLY | LOW_LATENCY_CALC | PRE_FILTER(expr) | FORCE_OUTPUT | MAX_DELAY(delay_time) | EVENT_TYPE(event_types) | IGNORE_NODATA_TRIGGER | IDLE_TIMEOUT(duration_time) | FLUSH_ON_OUTER_CLOSE}

你要仔细看一下TDengine的官网文档,里面有详细的介绍说明:

CREATE STREAM IF NOT EXISTS st_mileage_day

INTERVAL(1d) sliding(1d)

FROM bms

PARTITION BY tbname

STREAM_OPTIONS (

WATERMARK(1m) | FILL_HISTORY_FIRST(“2026-01-01 00:00:00”)

)

INTO bms_mileage_day

OUTPUT_SUBTABLE(CONCAT(“st_mileage_day_”, tbname))

AS

SELECT _twstart AS day_start,

_twend AS day_end,

FIRST(`totalMileage`) AS start_mileage,

LAST(`totalMileage`) AS end_mileage,

LAST(`totalMileage`) - FIRST(`totalMileage`) AS day_mileage

FROM %%tbname

WHERE ts >= _twstart AND ts < _twend; 当前我已经能够形成流计算数据了 但是他不像是实时来一条触发算一下 而是在大窗口下一天才触发一次吗? 我这个天的维度每天会有不停的更新数据 非要到第二天的数据来了 才会触发形成今天的结果吗

之前的版本是用的AT_ONCE这些来控制 当前3.4X使用什么来控制啊 我看INTERVAL(1d) sliding(1d)必须连用是这样吗 那还支持来一条触发算一条吗? 老师清楚吗