本文介绍ETL管理的SQL语句,包括创建、查看、修改、删除和列出ETL任务的操作语法与示例。
引擎与版本
以下语句仅适用于流引擎。要求3.1.8及以上版本。
您可以通过控制台查看并升级小版本。
CREATE ETL
CREATE ETL语句用于在流引擎中创建ETL任务。
语法
create_etl_statement ::= CREATE ETL [IF NOT EXISTS] etl_name
[WITH etl_properties]
AS INSERT INTO [[catalog_name.]db_name.]table_name column_list
select_statement
etl_properties ::= '(' property_definition (',' property_definition)* ')'
property_definition ::= property_name '=' property_value
column_list ::= '(' column_name (',' column_name)* ')'使用说明
ETL名称(etl_name)
必填参数。ETL名称的设置需遵循以下规则:
可包含数字、大写英文字符、小写英文字符、半角句号(.)、中划线(-)和下划线(_)。
不能以半角句号(.)或中划线(-)开头。
长度为1~255字符。
ETL属性(etl_properties)
您可以通过WITH关键字添加以下ETL属性:
设置时,属性名前后需添加反引号(`),属性值前后需添加单引号(')。例如`parallelism` = '2'。
属性 | 数据类型 | 说明 | 默认值 |
parallelism | INTEGER | 任务并行度。 | 1 |
sink.ignore-update-before | BOOLEAN | Sink时是否忽略-U。 | false |
sink.ignore-delete | BOOLEAN | Sink时是否忽略-D。 | false |
sink.null-mode | STRING | Sink时是否写入Null值,取值如下:
| NO_OP |
udf.xxxx | STRING | 配置UDF,需先上传UDF jar。参数格式如下: | 无 |
stream.xxx | ANY | 流引擎作业参数,例如: | 无 |
指定结果表
参数 | 是否必填 | 说明 |
catalog_name | 否 | 结果表的Catalog。 |
db_name | 否 | 结果表所在数据库。 |
table_name | 是 | 结果表的名称。 |
column_name | 是 | 结果表的列名。 |
SQL查询语句(select_statement)
用于筛选数据的SQL语句,例如 SELECT p1, c1 FROM `lindorm_table`.`default`.`source` WHERE c1 > 10;。
示例
假设宽表引擎中的源表source和结果表sink的结构如下:
-- 源表source
CREATE TABLE source(p1 INT, c1 DOUBLE, PRIMARY KEY(p1));
-- 结果表1:sink
CREATE TABLE sink(p1 INT, c1 DOUBLE, PRIMARY KEY(p1));示例一:创建ETL filter1,将源表
source中符合条件的数据插入到结果表sink中。CREATE ETL IF NOT EXISTS filter1 AS INSERT INTO `lindorm_table`.`default`.`sink` (p1, c1) SELECT p1, c1 FROM `lindorm_table`.`default`.`source` WHERE c1 > 10;示例二:创建ETL filter2,将源表
source中符合条件的数据插入到结果表sink中,同时添加属性。CREATE ETL IF NOT EXISTS filter2 WITH ( `parallelism` = '2', `stream.execution.checkpointing.interval` = '30000' ) AS INSERT INTO `lindorm_table`.`default`.`sink` (p1, c1) SELECT p1, c1 FROM `lindorm_table`.`default`.`source` WHERE c1 > 10;
DESCRIBE ETL
DESCRIBE ETL语句用于打印指定ETL任务的详细信息。
语法
describe_etl_statement ::= { DESCRIBE | DESC } ETL etl_name使用说明
etl_name为指定ETL任务的名称。必须填写。
返回结果集说明
字段 | 说明 |
SOURCE_SCHEMA | 源表字段。 |
SINK_SCHEMA | 结果表字段。 |
CONTENT | ETL任务完整的数据插入逻辑,可通过ALTER ETL修改。 |
ATTRIBUTES | 创建ETL时设置的ETL属性(etl_properties)。如果创建时未设置,则显示为空。 |
STATUS | ETL任务的状态,具体如下:
|
RESOURCE_GROUP | 当前使用的资源组,可通过SET语句修改。 |
CREATE_USER | 创建该ETL任务的用户。 |
CREATE_TIME | ETL任务的创建时间。 |
示例
打印当前资源组下ETL任务filter1的详细信息。
DESC ETL filter1;返回结果:
+---------+---------------+------------------------------------------+----------------------------------------+-------------------------------------------------------------------------------------------------------------------------------------+------------+---------+------------------+-------------+------------------------------+
| ETL_ID | ETL_VERSION | SOURCE_SCHEMA | SINK_SCHEMA | CONTENT | ATTRIBUTES | STATUS | RESOURCE_GROUP | CREATE_USER | CREATE_TIME |
+---------+---------------+------------------------------------------+----------------------------------------+-------------------------------------------------------------------------------------------------------------------------------------+------------+---------+------------------+-------------+------------------------------+
| filter1 | 1747293558894 | `lindorm_table`.`default`.`source`:p1,c1 | `lindorm_table`.`default`.`sink`:p1,c1 | INSERT INTO `lindorm_table`.`default`.`sink` (`p1`, `c1`) SELECT `p1`, `c1` FROM `lindorm_table`.`default`.`source` WHERE `c1` > 10 | {} | RUNNING | lstream-e00s**** | r*** | Thu May 15 15:19:18 CST 2025 |
+---------+---------------+------------------------------------------+----------------------------------------+-------------------------------------------------------------------------------------------------------------------------------------+------------+---------+------------------+-------------+------------------------------+ALTER ETL
ALTER ETL语句用于修改状态为RUNNING的ETL任务。
语法
alter_etl_statement ::= ALTER ETL etl_name
[WITH etl_properties]
AS INSERT INTO [[catalog_name.]db_name.]table_name column_list
select_statement
etl_properties ::= '(' property_definition (',' property_definition)* ')'
property_definition ::= property_name '=' property_value
column_list ::= '(' column_name (',' column_name)* ')'使用说明
ETL名称(etl_name)
必填参数。指定需要修改的ETL任务。
ETL属性(etl_properties)
您可以通过WITH关键字添加以下ETL属性:
设置时,属性名前后需添加反引号(`),属性值前后需添加单引号(')。例如`parallelism` = '2'。
属性 | 数据类型 | 说明 | 默认值 |
parallelism | INTEGER | 任务并行度。 | 1 |
sink.ignore-update-before | BOOLEAN | Sink时是否忽略-U。 | false |
sink.ignore-delete | BOOLEAN | Sink时是否忽略-D。 | false |
sink.null-mode | STRING | Sink时是否写入Null值,取值如下:
| NO_OP |
udf.xxxx | STRING | 配置UDF,需先上传UDF jar。参数格式如下: | 无 |
stream.xxx | ANY | 流引擎作业参数,例如: | 无 |
指定结果表
参数 | 是否必填 | 说明 |
catalog_name | 否 | 结果表的Catalog。 |
db_name | 否 | 结果表所在数据库。 |
table_name | 是 | 结果表的名称。 |
column_name | 是 | 结果表的列名。 |
SQL查询语句(select_statement)
修改为新的SQL查询语句。
示例
假设宽表引擎中的源表source和sink的结构如下:
-- 源表source
CREATE TABLE source(p1 INT, c1 DOUBLE, PRIMARY KEY(p1));
-- 结果表1:sink
CREATE TABLE sink(p1 INT, c1 DOUBLE, PRIMARY KEY(p1));创建ETL filter2,同时添加属性。
CREATE ETL IF NOT EXISTS filter2
WITH (
`parallelism` = '2',
`stream.execution.checkpointing.interval` = '30000'
)
AS
INSERT INTO `lindorm_table`.`default`.`sink` (p1, c1)
SELECT p1, c1 FROM `lindorm_table`.`default`.`source` WHERE c1 > 10;修改ETL任务属性
将parallelism属性修改为4。
ALTER ETL filter2
WITH (`parallelism` = '4')
AS
INSERT INTO `lindorm_table`.`default`.`sink` (p1, c1)
SELECT p1, c1 FROM `lindorm_table`.`default`.`source`;结果验证
您可以通过DESC ETL filter2;语句查看filter2的ATTRIBUTES,确认是否修改成功。
DROP ETL
DROP ETL语句用于删除当前资源组下的ETL任务。
语法
drop_etl_statement ::= DROP ETL [ IF EXISTS ] etl_name示例
删除当前资源组下名为filter2的ETL任务。
DROP ETL IF EXISTS filter2;结果验证
您可以通过SHOW ETLS;语句验证是否删除成功。
SHOW ETLS
SHOW ETLS语句用于列出当前资源组下所有ETL任务或名称符合匹配规则的ETL任务。
语法
show_etls_statement ::= SHOW ETLS [ LIKE string_literal ]使用说明
匹配表达式(LIKE string_literal)
LIKE关键字后的查找表达式是一个字符串常量,系统将根据该字符串常量模糊匹配系统属性。该字符串常量仅支持以下通配符:
%:替代0个或多个字符。_:替代一个字符。
返回结果集说明
字段 | 说明 |
ETL_ID | ETL任务名称。 |
ETL_VERSION | ETL版本号,仅显示最新版本。 |
STATUS | ETL任务的状态,具体如下:
|
JOB_RESTART_NUM | ETL任务重启次数。 |
NUM_RECORDS | ETL任务处理的数据条数。 |
RECORDS_PER_SEC | ETL任务每秒处理的数据条数。 |
NUM_BYTES | ETL任务处理的总数据量,单位为字节(byte)。 |
BYTES_PER_SEC | ETL任务每秒处理的数据量,单位为字节(byte)。 |
EVENT_TIME_LAG | ETL任务处理数据的延迟时间,单位为毫秒(ms)。如果结果显示为负数,表示无数据流入。 |
示例
展示所有ETL任务
列出当前资源组下所有ETL任务。
SHOW ETLS;返回结果:
+---------+---------------+---------+-----------------+-------------+-----------------+-----------+---------------+----------------+
| ETL_ID | ETL_VERSION | STATUS | JOB_RESTART_NUM | NUM_RECORDS | RECORDS_PER_SEC | NUM_BYTES | BYTES_PER_SEC | EVENT_TIME_LAG |
+---------+---------------+---------+-----------------+-------------+-----------------+-----------+---------------+----------------+
| filter1 | 1753940184757 | RUNNING | 0 | 0.0 | 0.0 | 0.0 | 0.0 | -1.0 |
| filter2 | 1753945414775 | RUNNING | 0 | 0.0 | 0.0 | 0.0 | 0.0 | -1.0 |
+---------+---------------+---------+-----------------+-------------+-----------------+-----------+---------------+----------------+模糊匹配ETL任务名称
列出当前资源组下的所有名称以filter开头的ETL任务。
SHOW ETLS LIKE 'filter%';返回结果:
+---------+---------------+---------+-----------------+-------------+-----------------+-----------+---------------+----------------+
| ETL_ID | ETL_VERSION | STATUS | JOB_RESTART_NUM | NUM_RECORDS | RECORDS_PER_SEC | NUM_BYTES | BYTES_PER_SEC | EVENT_TIME_LAG |
+---------+---------------+---------+-----------------+-------------+-----------------+-----------+---------------+----------------+
| filter1 | 1753940184757 | RUNNING | 0 | 0.0 | 0.0 | 0.0 | 0.0 | -1.0 |
| filter2 | 1753945414775 | RUNNING | 0 | 0.0 | 0.0 | 0.0 | 0.0 | -1.0 |
+---------+---------------+---------+-----------------+-------------+-----------------+-----------+---------------+----------------+