管理 Flow
每一个 flow 是 GreptimeDB 中的一个持续聚合查询。
它根据传入的数据持续更新并聚合数据。
本文档描述了如何创建和删除一个 flow。
EVAL INTERVAL 用于调度 batching 计算,但不会决定执行模式;TQL workload 必须使用该子句。有关执行路由和 instant-TTL 限制,请参阅创建 flow。
创建 source 表
在创建 flow 之前,你需要先创建一张 source 表来存储原始数据,比如:
CREATE TABLE temp_sensor_data (
sensor_id INT,
loc STRING,
temperature DOUBLE,
ts TIMESTAMP TIME INDEX,
PRIMARY KEY(sensor_id, loc)
);
对于新的 workload,请避免在 Flow source 表上使用 WITH ('ttl' = 'instant')。这是旧的使用方式,不推荐用于新的聚合或 TQL workload。请为 source 数据设置合适的保留策略。
创建 sink 表
flow 把聚合结果写入 sink 表。如果 sink 表不存在,CREATE FLOW 会在能够从查询结果推断 schema 时自动创建它。
如果需要控制 schema 或布局,或者查询较复杂难以推断,请预先创建 sink 表。已有 sink 表会根据 flow 查询结果进行校验。
source 表和 sink 表不能是同一张表。
sink 表必须与 flow 查询结果兼容,即:
- 列的顺序和类型:对于预先创建的 SQL sink,列的顺序和类型应与查询输出匹配。
- 时间索引:为 sink 表指定
TIME INDEX,通常使用时间窗口函数生成的时间列。 - 更新时间:自动创建的 batching SQL sink 会添加
update_at列来记录更新时间。TQL sink 遵循查询输出,不会自动添加update_at。预先创建的 SQL sink 可以与查询输出列数一致,也可以在末尾额外包含一个用于更新时间的时间戳列。 - Tag:使用
PRIMARY KEY指定 Tag,与 time index 一起作为行数据的唯一标识,并优化查询性能。
例如:
/* 创建 sink 表 */
CREATE TABLE temp_alerts (
sensor_id INT,
loc STRING,
max_temp DOUBLE,
time_window TIMESTAMP TIME INDEX,
update_at TIMESTAMP,
PRIMARY KEY(sensor_id, loc)
);
CREATE FLOW temp_monitoring
SINK TO temp_alerts
AS
SELECT
sensor_id,
loc,
max(temperature) AS max_temp,
date_bin('10 seconds'::INTERVAL, ts) AS time_window
FROM temp_sensor_data
GROUP BY
sensor_id,
loc,
time_window
HAVING max_temp > 100;
sink 表包含列 sensor_id、loc、max_temp、time_window 和 update_at。
- 前四列分别对应 flow 的查询结果列
sensor_id、loc、max(temperature)和date_bin('10 seconds'::INTERVAL, ts)。 time_window列被指定为 sink 表的TIME INDEX。update_at列是 schema 中的最后一列,用于存储数据的更新时间。- 最后的
PRIMARY KEY指定sensor_id和loc作为 Tag 列。 这意味着 flow 将根据 Tagsensor_id和loc以及时间索引time_window插入或更新数据。
创建 flow
创建 flow 的语法是:
CREATE [ OR REPLACE ] FLOW [ IF NOT EXISTS ] <flow-name>
SINK TO <sink-table-name>
[ EXPIRE AFTER <expr> ]
[ EVAL INTERVAL <interval> ]
[ COMMENT '<string>' ]
[ WITH (<flow-option> = <value> [, ...]) ]
AS
<SQL>;
子句必须按上述顺序出现:EXPIRE AFTER 在 EVAL INTERVAL 之前。
EVAL INTERVAL 用于调度 batching 计算。TQL Flow 必须使用该子句。包含 Aggregate 或 Distinct 的 SQL 计划使用 batching,除非
instant-TTL source 选择旧的 streaming;普通投影和非聚合 join 也使用旧的 streaming。streaming 会忽略 EVAL INTERVAL。
对于 batching 时间窗口 SQL,计算可以是增量的,而不一定执行完整查询;没有可用时间窗口的 batching SQL 必须指定 EVAL INTERVAL,并执行完整查询;时间窗口聚合可以不指定该子句。
当指定 OR REPLACE 时,如果已经存在同名的 flow,它将被更新为新 flow。请注意,这仅影响 flow 任务本身,source 表和 sink 表将不会被更改。当指定 IF NOT EXISTS 时,如果 flow 已经存在,它将不执行任何操作,而不是报告错误。还需要注意的是,OR REPLACE 不能与 IF NOT EXISTS 一起使用。
flow-name是目录级别的唯一标识符。sink-table-name是存储聚合数据的表名。 它可以是一个现有的表或一个新表;有关创建和校验行为,请参阅创建 sink 表。EXPIRE AFTER是一个可选的时间间隔,用于使 Flow 引擎中的数据过期。有关详细信息,请参考EXPIRE AFTER部分。EVAL INTERVAL是一个可选的时间间隔,用于调度 batching 计算;streaming 会忽略它。COMMENT是 flow 的描述。WITH指定 flow 选项。 本文档介绍的用户 Flow 选项为defer_on_missing_source和实验性的experimental_enable_incremental_read。SQL部分定义了用于持续聚合的查询。 它定义了为 flow 提供数据的源表。 每个 flow 可以有多个源表。 有关详细信息,请参考编写 SQL 查询部分。
一个创建 flow 的简单示例:
CREATE FLOW IF NOT EXISTS my_flow
SINK TO my_sink_table
EXPIRE AFTER '1 hour'::INTERVAL
COMMENT 'My first flow in GreptimeDB'
AS
SELECT
max(temperature) as max_temp,
date_bin('10 seconds'::INTERVAL, ts) as time_window
FROM temp_sensor_data
GROUP BY time_window;
创建的 flow 会将 max(temperature) 按 10 秒时间窗口分组,并将结果存储在 my_sink_table 中。
最近 1 小时内的数据会用于 flow 计算。
EXPIRE AFTER
EXPIRE AFTER 子句指定数据将在 flow 引擎中过期的时间间隔。
对于包含可用时间窗口表达式的 Flow,source 表中早于指定间隔的数据会被排除在计算之外,sink 表中较早的行也不会被更新。这会限制时间窗口 Flow 的状态和重新计算范围,包括涉及 GROUP BY 的有状态查询。
不包含可用时间窗口表达式的 batching 计划必须指定 EVAL INTERVAL,并执行未过滤的快照,除非查询本身包含时间谓词;EXPIRE AFTER 不会额外添加时间过滤。它不会删除 source 表或 sink 表中的数据。若需删除表数据,请在创建表时通过 TTL 策略实现。
例如,如果 flow 引擎在 10:00:00 处理聚合,并且设置了 '1 hour'::INTERVAL,
当前时刻若输入数据的 Time Index 超过 1 小时(即早于 09:00:00),则会被判定为过期数据并被忽略。
仅时间戳为 09:00:00 及之后的数据会参与聚合计算,并更新到目标表。