- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】flink
本篇文章系统讲解 Apache Flink 中 SQLDELETE语句的使用方法与底层实现。DELETE是 Flink Table 模块提供的行级删除能力,当前仅支持批(Batch)模式,要求目标表连接器实现SupportsRowLevelDelete接口。读完本文,你将掌握在 Java / Scala / Python 与 SQL CLI 中编写和运行DELETE语句的完整姿势,理解 Planner 如何将一条 DELETE 重写为删除/保留行集合的查询,并了解其与SupportsDeletePushDown下推的取舍关系。
概述与适用条件
DELETE语句用于根据可选过滤条件(filter)对目标表执行行级删除(row-level deletion),其整体能力由 SupportsRowLevelDelete.java 接口承载。
使用DELETE前必须清楚以下三个约束:
- 仅支持 Batch 模式:当前
DELETE语句只在批执行模式下可用,流模式下执行会报错; - 连接器必须实现
SupportsRowLevelDelete接口:该接口是 sink 能力的声明点,只有实现了该接口的DynamicTableSink才能消费行级删除产生的行数据; - 未实现接口时抛异常:若对未实现相关接口的表执行
DELETE,Planner 会抛出异常;此外,截至目前 Flink 官方维护的连接器中还没有一个内置支持DELETE(即官方连接器尚未内置实现该接口)。
注意:删除操作不可逆,执行前请务必确认过滤条件与目标表,避免误删数据。
运行一条 DELETE 语句
DELETE语句可以通过TableEnvironment的executeSql()方法执行。executeSql()会立即提交一个 Flink 作业,并返回与该作业关联的TableResult实例。Python 侧对应execute_sql()方法,SQL CLI 中则直接输入 SQL。
Java 示例
EnvironmentSettings settings = EnvironmentSettings.newInstance().inBatchMode().build(); TableEnvironment tEnv = TableEnvironment.create(settings); // register a table named "Orders" tEnv.executeSql("CREATE TABLE Orders (`user` STRING, product STRING, amount INT) WITH (...)"); // insert values tEnv.executeSql("INSERT INTO Orders VALUES ('Lili', 'Apple', 1), ('Jessica', 'Banana', 2), ('Mr.White', 'Chicken', 3)").await(); tEnv.executeSql("SELECT * FROM Orders").print(); // +--------------------------------+--------------------------------+-------------+ // | user | product | amount | // +--------------------------------+--------------------------------+-------------+ // | Lili | Apple | 1 | // | Jessica | Banana | 2 | // | Mr.White | Chicken | 3 | // +--------------------------------+--------------------------------+-------------+ // 3 rows in set // delete by filter tEnv.executeSql("DELETE FROM Orders WHERE `user` = 'Lili'").await(); tEnv.executeSql("SELECT * FROM Orders").print(); // +--------------------------------+--------------------------------+-------------+ // | user | product | amount | // +--------------------------------+--------------------------------+-------------+ // | Jessica | Banana | 2 | // | Mr.White | Chicken | 3 | // +--------------------------------+--------------------------------+-------------+ // 2 rows in set // delete entire table tEnv.executeSql("DELETE FROM Orders").await(); tEnv.executeSql("SELECT * FROM Orders").print(); // Empty set示例展示了两种典型用法:按过滤条件删除(DELETE FROM Orders WHERE \user` = 'Lili')与**删除全表数据**(不带WHERE的DELETE FROM Orders)。由于user是 SQL 保留字,示例中使用反引号 ```` 对其转义,这是 Flink SQL 中处理保留字的规范写法。
Scala 示例
val env = StreamExecutionEnvironment.getExecutionEnvironment() val settings = EnvironmentSettings.newInstance().inBatchMode().build() val tEnv = StreamTableEnvironment.create(env, settings) // register a table named "Orders" tEnv.executeSql("CREATE TABLE Orders (`user` STRING, product STRING, amount INT) WITH (...)"); // insert values tEnv.executeSql("INSERT INTO Orders VALUES ('Lili', 'Apple', 1), ('Jessica', 'Banana', 2), ('Mr.White', 'Chicken', 3)").await(); tEnv.executeSql("SELECT * FROM Orders").print(); // delete by filter tEnv.executeSql("DELETE FROM Orders WHERE `user` = 'Lili'").await(); tEnv.executeSql("SELECT * FROM Orders").print(); // delete entire table tEnv.executeSql("DELETE FROM Orders").await(); tEnv.executeSql("SELECT * FROM Orders").print(); // Empty setScala 场景下,需要基于StreamExecutionEnvironment构建StreamTableEnvironment,并同样通过EnvironmentSettings.newInstance().inBatchMode().build()显式切换到批模式。
Python 示例
env_settings = EnvironmentSettings.in_batch_mode() table_env = TableEnvironment.create(env_settings) # register a table named "Orders" table_env.executeSql("CREATE TABLE Orders (`user` STRING, product STRING, amount INT) WITH (...)"); # insert values table_env.executeSql("INSERT INTO Orders VALUES ('Lili', 'Apple', 1), ('Jessica', 'Banana', 2), ('Mr.White', 'Chicken', 3)").wait(); table_env.executeSql("SELECT * FROM Orders").print(); # 3 rows in set # delete by filter table_env.executeSql("DELETE FROM Orders WHERE `user` = 'Lili'").wait(); table_env.executeSql("SELECT * FROM Orders").print(); # 2 rows in set # delete entire table table_env.executeSql("DELETE FROM Orders").wait(); table_env.executeSql("SELECT * FROM Orders").print(); # Empty setPython API 中对应方法名为execute_sql(),返回结果的同步等待使用wait()而不是await()。
SQL CLI 示例
Flink SQL> SET 'execution.runtime-mode' = 'batch'; [INFO] Session property has been set. Flink SQL> CREATE TABLE Orders (`user` STRING, product STRING, amount INT) with (...); [INFO] Execute statement succeeded. Flink SQL> INSERT INTO Orders VALUES ('Lili', 'Apple', 1), ('Jessica', 'Banana', 1), ('Mr.White', 'Chicken', 3); [INFO] Submitting SQL update statement to the cluster... [INFO] SQL update statement has been successfully submitted to the cluster: Job ID: bd2c46a7b2769d5c559abd73ecde82e9 Flink SQL> SELECT * FROM Orders; user product amount Lili Apple 1 Jessica Banana 2 Mr.White Chicken 3 Flink SQL> DELETE FROM Orders WHERE `user` = 'Lili'; user product amount Jessica Banana 2 Mr.White Chicken 3在 SQL CLI 中,需要先通过SET 'execution.runtime-mode' = 'batch'将运行模式切换为批模式,再执行 CREATE / INSERT / DELETE 语句。注意DELETE是一个 DML(数据操纵语言)语句,执行时会向集群提交一个 Flink 作业,返回结果中的Job ID可用于在 Web UI 或日志中追踪作业状态。
DELETE ROWS 语法
DELETE FROM [catalog_name.][db_name.]table_name [ WHERE condition ]语法说明:
- 目标表标识:
table_name前可带可选的catalog_name与db_name两级命名空间前缀,用于定位不同 catalog / database 下的表,缺省时使用当前会话的默认 catalog 与默认 database; - 过滤条件:
WHERE condition为可选。省略WHERE时表示删除表中全部数据;带WHERE时仅删除满足条件的行; - 条件形式:
condition可以是任意合法表达式,支持等值/比较/逻辑组合,甚至可以包含子查询(见下文源码测试佐证),例如WHERE a >= (SELECT count(1) FROM t WHERE c > 1)。
底层实现:Planner 如何处理一条 DELETE
要理解 DELETE 的语义,需要回到 Planner 的语句转换链路。在 SqlNodeToOperationConversion.java 中,convertDelete(SqlDelete sqlDelete)负责把 Calcite 的SqlDelete语法树转换为 Table 层的 Operation:
- 标记修改类型:通过
RowLevelModificationContextUtils.setModificationType(...)将本次操作标记为DELETE,该上下文会传递给实现了SupportsRowLevelModificationScan的 source,使扫描阶段感知到“这是一次删除操作”; - 解析目标表:从
LogicalTableModify中取出表的限定名,并通过CatalogManager解析出ContextResolvedTable; - 优先尝试删除下推:调用
DeletePushDownUtils.getDynamicTableSink(...)获取表对应的DynamicTableSink。若 sink 实现了SupportsDeletePushDown且其applyDeleteFilters(filters)返回 true,则直接构造DeleteFromFilterOperation,由连接器在自身层面完成过滤删除,无需扫描全表; - 回退到行级删除:当下推不可用时,将 DELETE 重写为
SinkModifyOperation(ModifyType.DELETE),把“删除哪些行”的问题转化为“查询出哪些行并交给 sink 消费”的问题。PlannerQueryOperation中显式抛出TableException("Delete statements are not SQL serializable."),说明该查询仅供内部重写使用,不可序列化回 SQL。
这一设计印证了接口 javadoc 中的优先级约定:当表 sink 同时实现SupportsDeletePushDown与SupportsRowLevelDelete时,只要applyDeleteFilters返回 true,Planner 总是优先使用删除下推(SupportsRowLevelDelete.java)。
深度解析 SupportsRowLevelDelete 接口
作为行级删除的扩展点,SupportsRowLevelDelete 是一个@PublicEvolving接口,包含如下核心成员:
applyRowLevelDelete(context)
RowLevelDeleteInfo applyRowLevelDelete(@Nullable RowLevelModificationScanContext context);Planner 在重写 DELETE 语句前调用该方法,向 sink 询问“你期望以何种方式消费删除操作”。参数context由实现了SupportsRowLevelModificationScan的 table source 传入(若 source 未实现该接口则为null),用于在扫描阶段与删除阶段之间传递信息。
RowLevelDeleteInfo
该内部接口用来指导 Planner 如何重写 DELETE 语句,包含两个可覆写方法:
requiredColumns():返回 sink 执行行级删除所需的列。若返回Optional.empty(),表示需要全部列;否则 sink 消费到的行将按返回的列顺序排列。这在“删除只需主键/分区键”的场景下可显著减少跨网络传输的数据量;getRowLevelDeleteMode():返回删除模式,决定 Planner 将 DELETE 重写为“被删除行集合”还是“删除后的剩余行集合”,默认值为DELETED_ROWS。
RowLevelDeleteMode 枚举
enum RowLevelDeleteMode { DELETED_ROWS, REMAINING_ROWS }- DELETED_ROWS:sink 只收到匹配过滤条件、需要被删除的行。这些行统一携带
RowKind.DELETE语义; - REMAINING_ROWS:sink 收到的是删除后剩余的(即不匹配过滤条件的)行,统一携带
RowKind.INSERT语义。适合“整表重写”类存储(例如以覆盖方式重写文件的连接器)。
以DELETE FROM t WHERE y = 2;为例:若返回DELETED_ROWS,sink 会收到满足y = 2的行;若返回REMAINING_ROWS,sink 会收到不满足y = 2的行(参见接口 javadoc 中的示例说明)。
序列化与反序列化:RowLevelDeleteSpec
Planner 在将重写后的计划提交执行时,需要把 sink 能力序列化进 JSON 执行计划。这由 RowLevelDeleteSpec.java 完成:它以@JsonTypeName("RowLevelDelete")标识自身,序列化rowLevelDeleteMode与requiredPhysicalColumnIndices(所需物理列索引数组),并在apply(DynamicTableSink)时校验 sink 是否实现了SupportsRowLevelDelete——若未实现,则抛出TableException,这与文档中“未实现接口则报错”的描述一一对应。
测试佐证:Planner 层的行级删除行为
仓库中的测试用例可以帮你直观确认 DELETE 的行为边界:
- RowLevelDeleteTest.java 以参数化方式覆盖
DELETED_ROWS与REMAINING_ROWS两种模式,测试了:无条件删除(DELETE FROM t)、带过滤条件删除(DELETE FROM t where a = 1 and b = '123')、带子查询的删除(DELETE FROM t where b = '123' and a = (select count(*) from t))、指定自定义必需列('required-columns-for-delete' = 'b;c')以及元数据列删除等场景; - 测试中使用的
test-update-delete连接器(TestUpdateDeleteTableFactory.java)暴露了三个测试参数:required-columns-for-delete(必需列)、delete-mode(删除模式)与support-delete-push-down(是否支持删除下推),可用于本地复现验证; - 运行时集成测试 DeleteTableITCase.java 进一步验证了行级删除、带分区列的删除、删除与插入混用的
StatementSet(statementSet.addInsertSql("DELETE FROM t"))等端到端行为,说明 DELETE 可以与其他 DML 语句一起在StatementSet中批量提交。
常见问题与注意事项
- 流模式能否使用 DELETE?不能。
DELETE目前只支持批模式,流作业中执行会失败; - 官方连接器为什么用不了 DELETE?当前 Flink 官方维护的连接器尚未实现
SupportsRowLevelDelete,因此对官方连接器建的表执行 DELETE 会抛异常。如需使用,需要自行实现该接口,或等待官方/第三方连接器支持; - DELETE 与 SupportsDeletePushDown 的区别?
SupportsDeletePushDown由连接器直接在过滤层面完成删除(更高效),SupportsRowLevelDelete则由 Planner 重写查询、将需要删除(或剩余)的行交给 sink。两者同时存在时,Planner 优先尝试下推; executeSql()的返回:DELETE 通过executeSql()执行会立即提交 Flink 作业并返回TableResult,可通过await()(Java/Scala)或wait()(Python)同步等待作业完成,再执行后续查询验证结果。
- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】flink
相关推荐
Flink Table SQL DELETE 语句完全指南:行级删除、语法与连接器实现机制
Flink Table SQL DELETE 语句完全指南:行级删除、语法与连接器实现机制 DELETE 是 Flink Table API & SQL 提供的
大数据流处理批处理数据工程Flink 窗口去重(Window Deduplication)SQL 详解:语法、示例与实现原理
Flink 窗口去重(Window Deduplication)SQL 详解:语法、示例与实现原理 窗口去重(Window Deduplication)是 Fl
大数据流处理批处理数据工程STL到STEP转换引擎:打破3D打印与精密制造间的格式壁垒
STL到STEP转换引擎:打破3D打印与精密制造间的格式壁垒 在数字化设计与制造领域,工程师们长期面临着一个技术难题:如何将3D打印中广泛使用的STL格式无缝转
大数据流处理批处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考