ETL管理

更新时间:
复制 MD 格式

本文介绍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(保留原数据Null值,直接写入)

  • SKIP(跳过Null值,不写入)

NO_OP

udf.xxxx

STRING

配置UDF,需先上传UDF jar。参数格式如下:udf.<udfFunction> = <jarName>#<className>,其中udfFunction是使用udf的函数名,jarName是该udfjar包名,className是具体的类名。

stream.xxx

ANY

流引擎作业参数,例如:execution.checkpointing.interval

指定结果表

参数

是否必填

说明

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任务的状态,具体如下:

  • SUBMIT(任务已提交)

  • RUNNING(运行中)

  • DELETE(任务删除中)

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语句用于修改状态为RUNNINGETL任务。

语法

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(保留原数据Null值,直接写入)

  • SKIP(跳过Null值,不写入)

NO_OP

udf.xxxx

STRING

配置UDF,需先上传UDF jar。参数格式如下:udf.<udfFunction> = <jarName>#<className>,其中udfFunction是使用udf的函数名,jarName是该udfjar包名,className是具体的类名。

stream.xxx

ANY

流引擎作业参数,例如:execution.checkpointing.interval

指定结果表

参数

是否必填

说明

catalog_name

结果表的Catalog。

db_name

结果表所在数据库。

table_name

结果表的名称。

column_name

结果表的列名。

SQL查询语句(select_statement)

修改为新的SQL查询语句。

示例

假设宽表引擎中的源表sourcesink的结构如下:

-- 源表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;语句查看filter2ATTRIBUTES,确认是否修改成功。

DROP ETL

DROP ETL语句用于删除当前资源组下的ETL任务。

语法

drop_etl_statement ::= DROP ETL [ IF EXISTS ] etl_name

示例

删除当前资源组下名为filter2ETL任务。

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任务的状态,具体如下:

  • SUBMIT(任务已提交)

  • RUNNING(运行中)

  • DELETE(任务删除中)

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           |
+---------+---------------+---------+-----------------+-------------+-----------------+-----------+---------------+----------------+