流计算sql疑问

【TDengine 使用环境】
测试

【TDengine 版本】

3.4.1.0

【操作系统以及版本】

centos7.9

【部署方式】容器

CREATE STREAM IF NOT EXISTS nbq_hourly
INTERVAL(1h) SLIDING(1h)
FROM product_nbq
PARTITION BY tbname, station_id, station_area_id, device_model_id, product_id, device_id
STREAM_OPTIONS (
    WATERMARK(1m) | FILL_HISTORY_FIRST('2026-01-01 00:00:00')
)
INTO stream_nbq_hourly
OUTPUT_SUBTABLE(CONCAT('stream_nbq_hourly_', CAST(device_id AS VARCHAR)))
TAGS (
    station_id BIGINT AS %%2,
    station_area_id BIGINT AS %%3,
    device_model_id BIGINT AS %%4,
    product_id BIGINT AS %%5,
    device_id BIGINT AS %%6
)
AS
SELECT
    _twend AS ts,
    last_row(ljfdl) AS ljfdl,
    last_row(yggl) AS yggl
FROM %%tbname
WHERE _c0 >= _twstart AND _c0 <= _twend

请问这个流计算的触发逻辑是什么样的,假设现在是下午15:30分,结果表里会出现 16:00 的数据吗。按第一理解应该是时间窗口结束触发一次计算,但是发现在15-16点的过程中,计算是一直产生的。具体是什么逻辑呀

假设现在是下午 15:30,结果表里会出现 16:00 的数据吗?
答案是:会。

而且不仅会出现 16:00 的,在 15:00 到 16:00 之间,结果表里已经产生了很多条 ts = 16:00 的数据了(只不过最后一条是最终结果)。

核心原因:为什么“窗口结束触发一次”的理解是错的?

你的第一理解(窗口结束触发一次)是批处理(Batch) 的思维。在流计算中,特别是设置了 SLIDING(1h) 且没有设置 DELAY 的情况下,触发逻辑是 “微批 + 水位线(Watermark)推进”

具体逻辑如下:

1. 触发机制是“数据驱动”,而非“时间驱动”

  • 窗口 [15:00, 16:00) 定义好了。
  • 只要有一条数据的时间戳(_c0 )落入这个窗口 ,系统就会立即计算这个窗口的当前聚合结果(Partial Result) ,并写入结果表(UPSERT 更新)。
  • 所以,在 15:00 到 16:00 之间,只要有设备上报数据,这个窗口就会持续不断地被触发 ,结果表中的 16:00 这条记录会被反复更新(覆盖)

2. WATERMARK(1m) 的作用是“截止”,不是“开始”

  • WATERMARK(1m) 意味着:系统允许数据最多迟到 1 分钟
  • 当系统收到一条时间戳为 15:59:59 的数据时,水位线会推进到 15:59:59 - 1m = 15:58:59
  • 因为 15:58:59 < 16:00:00 ,窗口还未关闭,继续更新。
  • 直到系统收到一条时间戳为 16:00:01 的数据 ,水位线推进到 16:00:01 - 1m = 15:59:01 。此时,水位线依然小于 16:00,窗口仍未关闭。
  • 直到系统收到时间戳为 16:01:00 的数据 ,水位线推进到 16:00:00 ,此时系统才判定:[15:00, 16:00) 窗口不再会有新数据了。此时,系统会输出该窗口的“最终结果(Final Result)”,并关闭该窗口的更新。

感谢解答~!

这个问题发生在我15:00-16:00期间创建这个流计算的时候,但是16:00 之后一直没有出现 17:00 的数据,这个是正常的吗?还是说要等到17:00后才会出现。那流计算不就是只计算一次吗

  1. 16:00之后没有出现17:00的数据,是正常的吗? —— 是绝对正常的。
  2. 流计算不就是只计算一次吗? —— 不是的 ,但你的观察恰恰暴露了**“历史数据回填(FILL_HISTORY_FIRST)”** 和**“实时数据”** 在触发逻辑上的巨大差异。

流计算不是“只计算一次”,而是“每个窗口只输出一次最终结果(在实时模式下)”。

但因为你用了 FILL_HISTORY_FIRST ,你经历了两个截然不同的阶段:

1. 历史回填阶段(15:00-16:00 你观察到的)

  • 触发逻辑批量扫描(Batch Scan)
  • 系统为了追历史,会强行把所有时间窗口都算一遍,不管有没有数据。这导致你看到 16:00 疯狂更新,且窗口刚结束就立马出现。
  • 这个阶段的“多次更新” ,是因为历史数据量巨大,系统在逐批提交(Checkpoint)造成的假象。

2. 实时阶段(16:00 之后 到 17:00)

  • 触发逻辑数据驱动(Data Driven)
  • 系统现在“佛系”了,只有收到新数据,才触发计算
  • 因为 17:00 的窗口是 [16:00:00, 17:00:00) ,在 16:00-17:00 之间,如果没有任何设备上报时间戳大于 16:00:00 的数据(或者数据还没到),系统不会 凭空创建 17:00 的记录。
  • 只有在收到第一条 16:01 的数据时17:00 这条记录才会被首次 Insert 进结果表。

我的现象是时间到了16:01, 且收到了 16:01 的数据,但是17:00的记录没有insert进结果表。结果表最新数据还是16:00的,即使收到了16:40分的数据,最新数据仍然是 16:00 的。

针对上面一些提问,来回答一下这个流相关的所有活动以及创建使用过程
这个流主要完成以下这个任务: 取每个小时窗口内的最后一条记录写入到对应的子表中

  1. 这是一个触发为时间窗口类型, 建立的窗口为[00:00 01:00)[01:00 02:00)[02:00 03:00)…
  2. 设置的数据过期参数为WATERMARK为1分钟,IGNORE_DISORDER这参数没有指定,说明不忽略乱序数据,并且指定FILL_HISTORY_FIRST这个从某个固定的时间第一次优先开始计算。
  3. 查询的是按子表的相关tag进行分组,去获取该分组的最后一条记录,然后写入到对应的结果子表中去。

执行过程:

  1. 当第一次创建启动的时候,由于指定了FILL_HISTORY_FIRST这个参数,他会根据这个参数的时间开始执行历史数据的计算,也就是创建了从这一天到当前的每个小时的最后的那条数据的记录。
  2. 历史记录计算完毕之后,就进入实时计算环节。在实际的应用中,可能产生的数据不是完全有序的,像上面的情况,就是产生了数据的乱序。那这个时候他是怎么算的呢?
    比如:当前数据库的窗口是[16:00 17:00),也就是一直在存储和计算16点到17点的数据,结果表中也只有16:00之前的数据结果。如果这个时候,突然到来一条17:10的数据,显然这条数据不在当前的窗口之中,说明该窗口出现关闭事件。该流计算会在结果中写入一条17:00的最后一条的记录。并且该窗口处于关闭状态,开启了另外一个[17:00 1800)的窗口。然后,后面又会来过16点之后的数据,但是该数据不会落在当前的窗口中,前面IGNORE_DISORDER这参数没有指定,说明不忽略乱序数据。到这个时候,就会产生一种情况和现场,他来一条16点多的数据,它就会进行计算和更新一次17点这个时间戳的结果集。
  3. 另外一点就是时间窗口跟当前计算机时间是没有关系的,主要是数据带上来的时间戳主键有关系。