管理 Flow
每一个 flow 是 GreptimeDB 中的一个持续聚合查询。
它根据传入的数据持续更新并聚合数据。
本文档描述了如何创建和删除一个 flow。
创建输入表
在创建 flow 之前,你需要先创建一张输入表来存储原始的输入数据,比如:
CREATE TABLE temp_sensor_data (
sensor_id INT,
loc STRING,
temperature DOUBLE,
ts TIMESTAMP TIME INDEX,
PRIMARY KEY(sensor_id, loc)
);
但是如果你不想存储输入数据,可以在创建输入表时设置表选项 WITH ('ttl' = 'instant') 如下:
CREATE TABLE temp_sensor_data (
sensor_id INT,
loc STRING,
temperature DOUBLE,
ts TIMESTAMP TIME INDEX,
PRIMARY KEY(sensor_id, loc)
) WITH ('ttl' = 'instant');
将 ttl 设置为 'instant' 会使得输入表成为一张临时的表,也就是说它会自动丢弃一切插入的数据,而表本身一直会是空的,插入数据只会被送到 flow 任务处用作计算用途。
创建 sink 表
在创建 flow 之前,你需要有一个 sink 表来存储 flow 生成的聚合数据。 虽然它与常规的时间序列表相同,但有一些重要的注意事项:
- 列的顺序和类型:确保 sink 表中列的顺序和类型与 flow 查询结果匹配。
- 时间索引:为 sink 表指定
TIME INDEX,通常使用时间窗口函数生成的时间列。 - 更新时间:Flow 引擎会自动将更新时间附加到每个计算结果行的末尾。此更新时间存储在
update_at列中。请确保在 sink 表的 schema 中包含此列。 - 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> ]
[ COMMENT '<string>' ]
AS
<SQL>;
当指定 OR REPLACE 时,如果已经存在同名的 flow,它将被更新为新 flow。请注意,这仅影响 flow 任务本身,source 表和 sink 表将不会被更改。当指定 IF NOT EXISTS 时,如果 flow 已经存在,它将不执行任何操作,而不是报告错误。还需要注意的是,OR REPLACE 不能与 IF NOT EXISTS 一起使用。
flow-name是目录级别的唯一标识符。sink-table-name是存储聚合数据的表名。 它可以是一个现有的表或一个新表。如果目标表不存在,flow将创建目标表。EXPIRE AFTER是一个可选的时间间隔,用于从 Flow 引擎中过期数据。 有关更多详细信息,请参考EXPIRE AFTER部分。COMMENT是 flow 的描述。SQL部分定义了用于持续聚合的查询。 它定义了为 flow 提供数据的源表。 每个 flow 可以有多个源表。 有关详细信息,请参考编写查询 部分。
一个创建 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 将每 10 秒计算一次 max(temperature) 并将结果存储在 my_sink_table 中。
所有在 1 小时内的数据都将用于 flow 计算。