ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Flink SQL ALTER 语句完全指南:ALTER TABLE / VIEW / DATABASE / FUNCTION / CATALOG 实战与原理

Flink SQL ALTER 语句完全指南:ALTER TABLE / VIEW / DATABASE / FUNCTION / CATALOG 实战与原理 Flink SQL ALTER 语句完全指南ALTER TABLE / VIEW / DATABASE / FUNCTION / CATALOG 实战与原理【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink导读在 Flink Table / SQL 生态中ALTER语句用于修改一个已经在 Catalog 中注册的表、视图、函数或数据库的定义也可以直接修改 Catalog 本身的属性。它是元数据治理、Schema 演化和流作业平滑迭代的核心工具你无需删除并重建对象就能动态地增删列、调整主键、更新 watermark 策略、管理分区、切换表选项甚至替换 UDF 实现。读完本文你将掌握 Flink SQL 全部五种 ALTER 语句ALTER TABLE、ALTER VIEW、ALTER DATABASE、ALTER FUNCTION、ALTER CATALOG的完整语法、Java/Scala/Python/SQL CLI 四种执行方式以及从解析器到 Catalog 底层实现的源码级原理。概览Flink SQL 支持的 ALTER 语句ALTER语句用于修改一个已经在 Catalog 中注册的表、视图或函数定义或 catalog 本身的定义。当前 Flink SQL 支持以下五种 ALTER 语句ALTER TABLEALTER VIEWALTER DATABASEALTER FUNCTIONALTER CATALOG执行 ALTER 语句TableEnvironment 编程接口在 Java 和 Scala 中使用TableEnvironment的executeSql()方法执行 ALTER 语句若执行成功方法返回OK否则抛出异常。Python 中对应的方法是execute_sql()。此外也可以在 SQL CLI 中直接交互式执行 ALTER 语句。以下是一个覆盖全部 ALTER 场景的完整示例Java / Scala / Python 三端 API 等价TableEnvironment tableEnv TableEnvironment.create(...); // 注册名为 “Orders” 的表 tableEnv.executeSql(CREATE TABLE Orders (user BIGINT, product STRING, amount INT) WITH (...)); // 字符串数组 [Orders] String[] tables tableEnv.listTables(); // or tableEnv.executeSql(SHOW TABLES).print(); // 新增列 order 并置于第一位 tableEnv.executeSql(ALTER TABLE Orders ADD order INT COMMENT order identifier FIRST); // 新增更多列, 以及主键和 watermark tableEnv.executeSql(ALTER TABLE Orders ADD (ts TIMESTAMP(3), category STRING AFTER product, PRIMARY KEY(order) NOT ENFORCED, WATERMARK FOR ts AS ts - INTERVAL 1 HOUR)); // 修改列类型, 注释及 watermark 策略 tableEnv.executeSql(ALTER TABLE Orders MODIFY (amount DOUBLE NOT NULL, category STRING COMMENT category identifier AFTER order, WATERMARK FOR ts AS ts)); // 删除 watermark tableEnv.executeSql(ALTER TABLE Orders DROP WATERMARK); // 删除列 tableEnv.executeSql(ALTER TABLE Orders DROP (amount, ts, category)); // 重命名列 tableEnv.executeSql(ALTER TABLE Orders RENAME order TO order_id); // Orders 的表名改为 NewOrders tableEnv.executeSql(ALTER TABLE Orders RENAME TO NewOrders); // 字符串数组[NewOrders] String[] tables tableEnv.listTables(); // or tableEnv.executeSql(SHOW TABLES).print(); // 注册名为 cat2 的 catalog tableEnv.executeSql(CREATE CATALOG cat2 WITH (typegeneric_in_memory)); // 增加属性 default-database tableEnv.executeSql(ALTER CATALOG cat2 SET (default-databasedb));Scala 版本语法一致val tableEnv TableEnvironment.create(...) // 注册名为 “Orders” 的表 tableEnv.executeSql(CREATE TABLE Orders (user BIGINT, product STRING, amount INT) WITH (...)) // 新增列 order 并置于第一位 tableEnv.executeSql(ALTER TABLE Orders ADD order INT COMMENT order identifier FIRST) // 新增更多列, 以及主键和 watermark tableEnv.executeSql(ALTER TABLE Orders ADD (ts TIMESTAMP(3), category STRING AFTER product, PRIMARY KEY(order) NOT ENFORCED, WATERMARK FOR ts AS ts - INTERVAL 1 HOUR)) // 修改列类型, 注释, 以及主键和 watermark tableEnv.executeSql(ALTER TABLE Orders MODIFY (amount DOUBLE NOT NULL, category STRING COMMENT category identifier AFTER order, WATERMARK FOR ts AS ts)) // 删除 watermark tableEnv.executeSql(ALTER TABLE Orders DROP WATERMARK) // 删除列 tableEnv.executeSql(ALTER TABLE Orders DROP (amount, ts, category)) // 重命名列 tableEnv.executeSql(ALTER TABLE Orders RENAME order TO order_id) // 字符串数组 [Orders] val tables tableEnv.listTables() // or tableEnv.executeSql(SHOW TABLES).print() // rename Orders to NewOrders tableEnv.executeSql(ALTER TABLE Orders RENAME TO NewOrders) // 字符串数组[NewOrders] val tables tableEnv.listTables() // or tableEnv.executeSql(SHOW TABLES).print() // 注册名为 cat2 的 catalog tableEnv.executeSql(CREATE CATALOG cat2 WITH (typegeneric_in_memory)) // 增加属性 default-database tableEnv.executeSql(ALTER CATALOG cat2 SET (default-databasedb))Python 版本语法一致table_env TableEnvironment.create(...) # 字符串数组 [Orders] tables table_env.list_tables() # or table_env.execute_sql(SHOW TABLES).print() # 新增列 order 并置于第一位 table_env.execute_sql(ALTER TABLE Orders ADD order INT COMMENT order identifier FIRST) # 新增更多列, 主键及 watermark table_env.execute_sql(ALTER TABLE Orders ADD (ts TIMESTAMP(3), category STRING AFTER product, PRIMARY KEY(order) NOT ENFORCED, WATERMARK FOR ts AS ts - INTERVAL 1 HOUR)) # 修改列类型, 列注释, 主键及 watermark table_env.execute_sql(ALTER TABLE Orders MODIFY (amount DOUBLE NOT NULL, category STRING COMMENT category identifier AFTER order, WATERMARK FOR ts AS ts)) # 删除 watermark table_env.execute_sql(ALTER TABLE Orders DROP WATERMARK) # 删除列 table_env.execute_sql(ALTER TABLE Orders DROP (amount, ts, category)) # 重命名列 table_env.execute_sql(ALTER TABLE Orders RENAME order TO order_id) # 把 Orders 的表名改为 NewOrders table_env.execute_sql(ALTER TABLE Orders RENAME TO NewOrders) # 字符串数组[NewOrders] tables table_env.list_tables() # or table_env.execute_sql(SHOW TABLES).print() # 注册名为 cat2 的 catalog table_env.execute_sql(CREATE CATALOG cat2 WITH (typegeneric_in_memory)) # 增加属性 default-database table_env.execute_sql(ALTER CATALOG cat2 SET (default-databasedb))SQL CLI 交互式执行在 SQL CLI 中可以逐条执行 ALTER 语句并用DESCRIBE、SHOW TABLES、DESC CATALOG EXTENDED实时验证变更结果。完整的交互过程如下Flink SQL CREATE TABLE Orders (user BIGINT, product STRING, amount INT) WITH (...); [INFO] Execute statement succeeded. Flink SQL ALTER TABLE Orders ADD order INT COMMENT order identifier FIRST; [INFO] Execute statement succeeded. Flink SQL DESCRIBE Orders; ----------------------------------------------------------------- | name | type | null | key | extras | watermark | comment | ----------------------------------------------------------------- | order | INT | TRUE | | | | order identifier | | user | BIGINT | TRUE | | | | | | product | STRING | TRUE | | | | | | amount | INT | TRUE | | | | | ----------------------------------------------------------------- 4 rows in set Flink SQL ALTER TABLE Orders ADD (ts TIMESTAMP(3), category STRING AFTER product, PRIMARY KEY(order) NOT ENFORCED, WATERMARK FOR ts AS ts - INTERVAL 1 HOUR); [INFO] Execute statement succeeded. Flink SQL DESCRIBE Orders; --------------------------------------------------------------------------------------------------------- | name | type | null | key | extras | watermark | comment | --------------------------------------------------------------------------------------------------------- | order | INT | FALSE | PRI(order) | | | order identifier | | user | BIGINT | TRUE | | | | | | product | STRING | TRUE | | | | | | category | STRING | TRUE | | | | | | amount | INT | TRUE | | | | | | ts | TIMESTAMP(3) *ROWTIME* | TRUE | | | ts - INTERVAL 1 HOUR | | --------------------------------------------------------------------------------------------------------- 6 rows in set Flink SQL ALTER TABLE Orders MODIFY (amount DOUBLE NOT NULL, category STRING COMMENT category identifier AFTER order, WATERMARK FOR ts AS ts); [INFO] Execute statement succeeded. Flink SQL DESCRIBE Orders; --------------------------------------------------------------------------------------------- | name | type | null | key | extras | watermark | comment | --------------------------------------------------------------------------------------------- | order | INT | FALSE | PRI(order) | | | order identifier | | category | STRING | TRUE | | | | category identifier | | user | BIGINT | TRUE | | | | | | product | STRING | TRUE | | | | | | amount | DOUBLE | FALSE | | | | | | ts | TIMESTAMP(3) *ROWTIME* | TRUE | | | ts | | --------------------------------------------------------------------------------------------- 6 rows in set Flink SQL ALTER TABLE Orders DROP WATERMARK; [INFO] Execute statement succeeded. Flink SQL DESCRIBE Orders; ----------------------------------------------------------------------------------- | name | type | null | key | extras | watermark | comment | ----------------------------------------------------------------------------------- | order | INT | FALSE | PRI(order) | | | order identifier | | category | STRING | TRUE | | | | category identifier | | user | BIGINT | TRUE | | | | | | product | STRING | TRUE | | | | | | amount | DOUBLE | FALSE | | | | | | ts | TIMESTAMP(3) | TRUE | | | | | ----------------------------------------------------------------------------------- 6 rows in set Flink SQL ALTER TABLE Orders DROP (amount, ts, category); [INFO] Execute statement succeeded. Flink SQL DESCRIBE Orders; ------------------------------------------------------------------------- | name | type | null | key | extras | watermark | comment | ------------------------------------------------------------------------- | order | INT | FALSE | PRI(order) | | | order identifier | | user | BIGINT | TRUE | | | | | | product | STRING | TRUE | | | | | ------------------------------------------------------------------------- 3 rows in set Flink SQL ALTER TABLE Orders RENAME order to order_id; [INFO] Execute statement succeeded. Flink SQL DESCRIBE Orders; ----------------------------------------------------------------------------- | name | type | null | key | extras | watermark | comment | ----------------------------------------------------------------------------- | order_id | INT | FALSE | PRI(order_id) | | | order identifier | | user | BIGINT | TRUE | | | | | | product | STRING | TRUE | | | | | ----------------------------------------------------------------------------- 3 rows in set Flink SQL SHOW TABLES; ------------ | table name | ------------ | Orders | ------------ 1 row in set Flink SQL ALTER TABLE Orders RENAME TO NewOrders; [INFO] Execute statement succeeded. Flink SQL SHOW TABLES; ------------ | table name | ------------ | NewOrders | ------------ 1 row in set Flink SQL CREATE CATALOG cat2 WITH (typegeneric_in_memory); [INFO] Execute statement succeeded. Flink SQL ALTER CATALOG cat2 SET (default-databasedb_new); [INFO] Execute statement succeeded. Flink SQL DESC CATALOG EXTENDED cat2; -------------------------------------------- | info name | info value | -------------------------------------------- | name | cat2 | | type | generic_in_memory | | comment | | | option:default-database | db_new | -------------------------------------------- 4 rows in set从上述过程可以直观看到ADD/MODIFY/DROP之后DESCRIBE的type、null、key、watermark、comment列会即时反映 Schema 变化*ROWTIME*标记的ts列在删除 watermark 后也会恢复为普通时间戳列。ALTER TABLE当前支持的 ALTER TABLE 语法如下ALTER TABLE [IF EXISTS] table_name { ADD { schema_component | (schema_component [, ...]) | [IF NOT EXISTS] partition_component [partition_component ...] | distribution } | MODIFY { schema_component | (schema_component [, ...]) | distribution } | DROP {column_name | (column_name, column_name, ....) | PRIMARY KEY | CONSTRAINT constraint_name | WATERMARK | [IF EXISTS] partition_component [, ...] | DISTRIBUTION } | RENAME old_column_name TO new_column_name | RENAME TO new_table_name | SET (key1val1, ...) | RESET (key1, ...) } schema_component: { column_component | constraint_component | watermark_component } column_component: column_name column_definition [FIRST | AFTER column_name] constraint_component: [CONSTRAINT constraint_name] PRIMARY KEY (column_name, ...) NOT ENFORCED watermark_component: WATERMARK FOR rowtime_column_name AS watermark_strategy_expression column_definition: { physical_column_definition | metadata_column_definition | computed_column_definition } [COMMENT column_comment] physical_column_definition: column_type metadata_column_definition: column_type METADATA [ FROM metadata_key ] [ VIRTUAL ] computed_column_definition: AS computed_column_expression partition_component: PARTITION (key1val1, key2val2, ...) [WITH (key1val1, key2val2, ...)] distribution: { DISTRIBUTION BY [ { HASH | RANGE } ] (bucket_column_name1, bucket_column_name2, ...) ] [INTO n BUCKETS] | DISTRIBUTION INTO n BUCKETS }IF EXISTS若表不存在则不进行任何操作不抛出异常。ADD使用ADD语句向已有表中增加 columns、constraints、watermark、partitions 和 distribution。向表新增列时可通过FIRST或AFTER col_name指定位置不指定位置时默认追加在最后。-- 新增一列 ALTER TABLE MyTable ADD category_id STRING COMMENT identifier of the category; -- 新增列主键和 watermark ALTER TABLE MyTable ADD ( log_ts STRING COMMENT log timestamp string FIRST, ts AS TO_TIMESTAMP(log_ts) AFTER log_ts, PRIMARY KEY (id) NOT ENFORCED, WATERMARK FOR ts AS ts - INTERVAL 3 SECOND ); -- 新增一个分区 ALTER TABLE MyTable ADD PARTITION (p11,p2a) with (k1v1); -- 新增两个分区 ALTER TABLE MyTable ADD PARTITION (p11,p2a) with (k1v1) PARTITION (p11,p2b) with (k2v2); -- add new distribution using a hash on uid into 4 buckets ALTER TABLE MyTable ADD DISTRIBUTION BY HASH(uid) INTO 4 BUCKETS; -- add new distribution on uid into 4 buckets CREATE TABLE MyTable ADD DISTRIBUTION BY (uid) INTO 4 BUCKETS; -- add new distribution on uid. CREATE TABLE MyTable ADD DISTRIBUTION BY (uid); -- add new distribution into 4 buckets CREATE TABLE MyTable ADD DISTRIBUTION INTO 4 BUCKETS;注意指定列为主键列时会隐式修改该列的 nullability 为 false。源码佐证ADD 语义在解析层面由SqlAlterTableAdd、SqlAlterTableAddConstraint、SqlAlterTableSchema等解析树节点表达见 flink-sql-parser 的 ddl 包最终被转换为 TableChange 变更对象。TableChange接口为每一种 ALTER TABLE 子操作定义了对应的静态工厂方法add(Column)/add(Column, ColumnPosition)对应 ADD 列add(UniqueConstraint)对应 ADD 主键add(WatermarkSpec)对应 ADD watermarkadd(TableDistribution)对应 ADD distribution。其中ColumnPosition为 null 表示追加在列尾为FIRST表示置于首位为AFTER表示跟在指定列之后。该接口被标记为PublicEvolving说明它同时是 Catalog 与 Planner 之间的公共变更契约。MODIFY使用MODIFY语句修改列的位置、类型、注释、nullability主键或 watermark。可使用FIRST或AFTER col_name将已有列移动至指定位置不指定时默认保持位置不变。-- modify a column type, comment and position ALTER TABLE MyTable MODIFY measurement double COMMENT unit is bytes per second AFTER id; -- modify definition of column log_ts and ts, primary key, watermark. They must exist in table schema ALTER TABLE MyTable MODIFY ( log_ts STRING COMMENT log timestamp string AFTER id, -- reorder columns ts AS TO_TIMESTAMP(log_ts) AFTER log_ts, PRIMARY KEY (id) NOT ENFORCED, WATERMARK FOR ts AS ts -- modify watermark strategy );注意指定列为主键列时会隐式修改该列的 nullability 为 false。源码佐证TableChange中对应modify(oldColumn, newColumn, columnPosition)及细粒度的modifyPhysicalColumnType改物理列类型构造时通过Preconditions.checkArgument(oldColumn.isPhysical())强制要求目标必须是物理列、modifyColumnName重命名列时按物理/元数据/计算列分别重建新列并保留原注释、modifyColumnComment、modifyColumnPosition、modify(UniqueConstraint)、modify(TableDistribution)、modify(WatermarkSpec)。因此 MODIFY 实际是旧列定义 新列定义 新位置三元组的一次性原子替换。DROP使用DROP语句删除列、主键、分区或 watermark。-- 删除一个列 ALTER TABLE MyTable DROP measurement; -- 删除多个列 ALTER TABLE MyTable DROP (col1, col2, col3); -- 删除主键 ALTER TABLE MyTable DROP PRIMARY KEY; -- 删除一个分区 ALTER TABLE MyTable DROP PARTITION (id 1); -- 删除两个分区 ALTER TABLE MyTable DROP PARTITION (id 1), PARTITION (id 2); -- 删除 watermark ALTER TABLE MyTable DROP WATERMARK; -- drop distribution ALTER TABLE MyTable DROP DISTRIBUTION;源码佐证解析层为每种 DROP 子句都定义了独立节点SqlAlterTableDropColumn、SqlAlterTableDropPrimaryKey、SqlAlterTableDropConstraint、SqlAlterTableDropWatermark、SqlAlterTableDropDistribution见 flink-sql-parser 的 ddl 包在 SqlNodeToOperationConversion.java 中分别转换为TableChange.dropColumn(...)、dropConstraint(...)、dropWatermark()、dropDistribution()。对 DROP 的约束性校验在测试中有明确覆盖例如删除被计算列引用的列、被用作分区键的列、被用作 distribution key 的列都会抛错见 SqlDdlToOperationConverterTest.java错误信息形如The column \e is referenced by computed column g, j.与The column a is used as the partition keys.。RENAME使用RENAME语句修改列名或表名。-- rename column ALTER TABLE MyTable RENAME request_body TO payload; -- rename table ALTER TABLE MyTable RENAME TO MyTable2;源码佐证列重命名与表重命名在解析层分属SqlAlterTableRenameColumn与SqlAlterTableRename两个节点其中列重命名仍转换为TableChange.ModifyColumnName而表重命名在SqlNodeToOperationConversion.convertAlterTable中单独生成AlterTableRenameOperation携带新旧表标识符与IF EXISTS标志不会经过TableChange变更管道。SET为指定的表设置一个或多个属性。若个别属性已经存在于表中则使用新值覆盖旧值。-- set rows-per-second ALTER TABLE DataGenSource SET (rows-per-second 10);RESET为指定的表重置一个或多个属性使其回退到默认值。-- reset rows-per-second to the default value ALTER TABLE DataGenSource RESET (rows-per-second);源码佐证与限制SET/RESET由SqlAlterTableOptions、SqlAlterTableReset节点解析并转换为TableChange.set(key, value)/TableChange.reset(key)最终封装为AlterTableChangeOperation。测试同时验证了重要的边界约束ALTER TABLE RESET不允许修改connector等关键选项ALTER TABLE RESET does not support changing connector且不允许空 keydoes not support empty key。这一点同样适用于 ALTER CATALOG 的 RESET不允许变更type不允许空 key意味着connector/type这类决定实现类型的属性是不可通过 ALTER 动态切换的。ALTER VIEWALTER VIEW [catalog_name.][db_name.]view_name RENAME TO new_view_name将给定视图重命名为同一 catalog 和 database 下的新名称。ALTER VIEW [catalog_name.][db_name.]view_name AS new_query_expression将定义给定视图的底层查询替换为新的查询表达式。使用提示重命名视图只影响元数据名称视图定义保持不变而AS new_query_expression形式用于视图定义演进例如底层表新增字段、查询逻辑调整替换后新查询会立即生效下游引用该视图的语句无需修改。ALTER DATABASEALTER DATABASE [catalog_name.]db_name SET (key1val1, key2val2, ...)在数据库中设置一个或多个属性。若个别属性已经在数据库中设定将会使用新值覆盖旧值。源码佐证该语句由SqlAlterDatabase解析后在 SqlNodeToOperationConversion.java 的convertAlterDatabase中构造AlterDatabaseOperation最终落到 Catalog.alterDatabase 接口void alterDatabase(String name, CatalogDatabase newDatabase, boolean ignoreIfNotExists)。它会把整个CatalogDatabase定义含 properties 集合替换为新定义并以ignoreIfNotExists决定数据库不存在时是静默忽略还是抛DatabaseNotExistException。ALTER FUNCTIONALTER [TEMPORARY|TEMPORARY SYSTEM] FUNCTION [IF EXISTS] [catalog_name.][db_name.]function_name AS identifier [LANGUAGE JAVA|SCALA|PYTHON]修改一个有 catalog 和数据库命名空间的 catalog function需要指定一个新的identifier并可指定 language tag。若函数不存在且未指定IF EXISTS操作会抛出异常。如果 language tag 是JAVA或SCALA则identifier是 UDF 实现类的全限定名。关于 JAVA/SCALA UDF 的实现参考 自定义函数。如果 language tag 是PYTHON则identifier是 UDF 对象的全限定名例如pyflink.table.tests.test_udf.add。关于 PYTHON UDF 的实现参考 Python UDFs。TEMPORARY修改一个有 catalog 和数据库命名空间的临时 catalog function并覆盖原有的 catalog function。TEMPORARY SYSTEM修改一个没有数据库命名空间的临时系统 catalog function并覆盖系统内置的函数。IF EXISTS若函数不存在则不进行任何操作。LANGUAGE JAVA|SCALA|PYTHONLanguage tag 用于指定 Flink runtime 如何执行这个函数。目前只支持 JAVA、SCALA 和 PYTHON且函数的默认语言为 JAVA。源码佐证convertAlterFunction在 SqlNodeToOperationConversion.java 中区分系统函数与非系统函数系统函数走AlterSystemFunctionOperation普通函数走AlterCatalogFunctionOperation并携带ifExists与isTemporary标志。注意原文档中的描述若函数不存在删除会抛出异常在语法注释中保留实际语义以IF EXISTS控制。ALTER CATALOGALTER CATALOG catalog_name SET (key1val1, ...) | RESET (key1, ...) | COMMENT commentSET为指定的 catalog 设置一个或多个属性。若个别属性已经存在则使用新值覆盖旧值。-- set default-database ALTER CATALOG cat2 SET (default-databasedb);RESET为指定的 catalog 重置一个或多个属性。-- reset default-database ALTER CATALOG cat2 RESET (default-database);COMMENT为指定的 catalog 设置注释。若注释已经存在则使用新值覆盖旧值。ALTER CATALOG cat2 COMMENT comment for catalog cat2;SQL 中连续两个单引号表示转义的单引号因此上述语句实际写入的注释为comment for catalog cat2。源码佐证测试用例验证了 ALTER CATALOG 的完整行为见 SqlDdlToOperationConverterTest.javaSET支持一次设置多组 key/value同名 key 后者覆盖前者k2 v2被k2 v2_new覆盖RESET不允许变更typeALTER CATALOG RESET does not support changing type且不允许空 keyCOMMENT会生成独立的 summary 字符串ALTER CATALOG cat2 COMMENT comment for catalog cat2。深层原理从 SQL 文本到 Catalog 变更的完整链路理解 ALTER 语句的底层实现有助于把握各语句的适用边界与限制。整条链路在仓库中清晰可见解析层flink-sql-parser基于 Calcite 的 SQL 解析器将ALTER TABLE文本解析为SqlAlterTable及其子类节点SqlAlterTableAdd、SqlAlterTableModify、SqlAlterTableSchema、SqlAlterTableDropColumn、SqlAlterTableRename、SqlAlterTableOptions、SqlAlterTableReset等见 flink-sql-parser ddl 包。SqlAlterTable抽象基类见 SqlAlterTable.java统一持有tableIdentifier、可选的partitionSpec和ifTableExists标志并提供fullTableName()、getPartitionKVs()等工具方法其中ifTableExists()的语义注释明确当指定IF EXISTS时返回 true用于表不存在时忽略错误。转换层flink-table-plannerSqlNodeToOperationConversion是 SQL AST 到 Operation 的中枢。在convertAlterTable中按子节点类型分派RENAME 走AlterTableRenameOperationSET/RESET 走AlterTableChangeOperation其余 schema 类变更统一解析为ListTableChange后封装为AlterTableChangeOperation。ALTER DATABASE 生成AlterDatabaseOperationALTER FUNCTION 生成AlterCatalogFunctionOperation/AlterSystemFunctionOperation。变更模型层flink-table-commonTableChange 是 ALTER TABLE 的统一变更语言。它把 ADD / MODIFY / DROP / SET / RESET 全部归一化为类型安全的变更对象每个静态工厂方法都在 Javadoc 中标注了其等价的 SQL 语句。例如TableChange.modifyColumnName(oldColumn, newName)等价于ALTER TABLE table_name RENAME old_column_name TO new_column_nameTableChange.set(key, value)等价于ALTER TABLE table_name SET key value。这种设计让 Catalog 实现只需接收目标新表 变更列表而无需理解 SQL 语法。Catalog 落地层flink-table-commonCatalog 接口定义了元数据修改的底层方法族alterDatabase(String name, CatalogDatabase newDatabase, boolean ignoreIfNotExists)、alterTable(ObjectPath tablePath, CatalogBaseTable newTable, boolean ignoreIfNotExists)、alterPartition(...)、alterFunction(...)、alterTableStatistics(...)、alterTableColumnStatistics(...)、alterPartitionStatistics(...)等。任何实现该接口的 Catalog如GenericInMemoryCatalog、Hive Catalog 等都通过这些方法把变更持久化到对应元数据存储。从工程视角看这一分层带来的收益是语法解析、语义校验、变更建模与存储实现完全解耦。新增一种 ALTER 子句只需在解析器与转换器各加一个节点Catalog 实现则通过统一的TableChange/Operation契约复用变更逻辑同时IF EXISTS、隐式 nullability 调整、RESET 保护 key 等一致性约束都在转换层集中校验对应测试见 SqlDdlToOperationConverterTest.java确保所有 Catalog 行为一致。总结与最佳实践Flink SQL 的 ALTER 语句族覆盖了表结构、分区、分布、视图定义、数据库属性、UDF 实现与 Catalog 属性五类对象的运行时修改是元数据治理与 Schema 演化的必备能力。实际使用中建议遵循以下几点尽量使用IF EXISTS实现幂等变更无论是ALTER TABLE IF EXISTS还是ALTER FUNCTION IF EXISTS都能让重复执行的运维脚本避免因对象缺失而中断。善用DESCRIBE/SHOW TABLES/DESC CATALOG EXTENDED校验变更SQL CLI 中这些命令与 ALTER 语句配合可以即时核对列类型、null 约束、主键、watermark 与注释是否符合预期。区分持久化与临时对象ALTER FUNCTION默认修改的是 catalog 中持久化的函数TEMPORARY/TEMPORARY SYSTEM仅作用于当前 Session适合快速试验或覆盖内置函数。警惕 RESET 的保护性约束connector、catalogtype这类决定实现类型的属性不可通过RESET重置遇到此类需求应重新 CREATE 对象而非试图 ALTER。理解变更的级联影响删除被计算列引用、被分区键或 distribution key 引用的列会被校验拒绝将列设为主键会隐式将其改为 NOT NULL这些约束都是 Planner 在转换阶段集中强制执行的。【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进