ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Flink Table SQL DELETE 语句完全指南:行级删除、语法与连接器实现机制

Flink Table SQL DELETE 语句完全指南:行级删除、语法与连接器实现机制 大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载DELETE是 Flink Table API SQL 提供的行级删除语句用于按照WHERE条件批量删除目标表中的数据。本指南以 Flink 官方 SQL 文档docs/content.zh/docs/dev/table/sql/delete.md为骨架结合flink-table模块的连接器能力接口与集成测试源码讲解 DELETE 语句的语法、四种执行方式Java / Scala / Python / SQL CLI、批模式限制以及SupportsRowLevelDelete与SupportsDeletePushDown两条底层实现路径。读完本文你将掌握 DELETE 语句的正确用法、触发条件与常见异常并理解如何为自己的自定义连接器接入行级删除能力。DELETE 语句概述与适用前提DELETE语句用于根据条件删除目标表中的数据。与关系型数据库中的 DELETE 不同Flink 的 DELETE 是一个提交即运行的批作业语句通过TableEnvironment提交后立即触发一个 Flink 作业执行。使用 DELETE 语句必须满足两个前提仅支持批模式当前 Flink 的 DELETE 语句只能在批模式下执行流模式下不适用。目标表连接器必须实现SupportsRowLevelDelete接口该接口定义在 flink-table/flink-table-common/src/main/java/org/apache/flink/table/connector/sink/abilities/SupportsRowLevelDelete.java只有实现了该接口的动态表 SinkDynamicTableSink才具备行级删除能力。如果一个表没有实现SupportsRowLevelDelete接口却执行了 DELETEFlink 会直接抛出异常。目前 Flink 内置维护的连接器如 Filesystem、JDBC、Hive 等均尚未实现该接口因此 DELETE 语句当前主要面向自定义连接器或第三方连接器开放。如果确实需要删除数据官方文档给出的替代方案是使用INSERT OVERWRITE重写整个表等方式实现等价的全量覆盖效果。DELETE 语句语法DELETE 语句的完整语法如下DELETE FROM [catalog_name.][db_name.]table_name [ WHERE condition ]语法要点说明catalog_name、db_name均可省略省略时使用当前会话的默认 Catalog 与默认数据库。table_name为必填的目标表名。WHERE condition为可选的过滤条件指定条件时仅删除满足条件的行条件删除省略条件时删除表中的全部数据全表删除。关键字不区分大小写字段名如需与关键字冲突可加反引号如user。从语法上看DELETE 与标准 SQL 高度一致但其底层执行机制与普通关系型数据库有本质区别具体见下文底层实现机制章节。执行 DELETE 语句DELETE 语句可以通过TableEnvironment的executeSql()方法Python 中为execute_sql()执行也可以在 SQL CLI 中直接输入。executeSql()执行 DELETE 语句时会立即提交一个 Flink 作业并返回一个TableResult对象通过TableResult.getJobClient()可以获取JobClient来方便地操作如取消、查询状态已提交的作业。Java 示例EnvironmentSettings settings EnvironmentSettings.newInstance().inBatchMode().build(); TableEnvironment tEnv TableEnvironment.create(settings); // 注册一个 Orders 表 tEnv.executeSql(CREATE TABLE Orders (user STRING, product STRING, amount INT) WITH (...)); // 插入原始数据 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 // 根据 where 条件删除 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 // 全表删除 tEnv.executeSql(DELETE FROM Orders).await(); tEnv.executeSql(SELECT * FROM Orders).print(); // Empty setScala 示例val env StreamExecutionEnvironment.getExecutionEnvironment() val settings EnvironmentSettings.newInstance().inBatchMode().build() val tEnv StreamTableEnvironment.create(env, settings) // 注册一个 Orders 表 tEnv.executeSql(CREATE TABLE Orders (user STRING, product STRING, amount INT) WITH (...)) // 插入原始数据 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 // 根据 where 条件删除 tEnv.executeSql(DELETE FROM Orders WHERE user Lili).await() tEnv.executeSql(SELECT * FROM Orders).print() // 2 rows in set // 全表删除 tEnv.executeSql(DELETE FROM Orders).await() tEnv.executeSql(SELECT * FROM Orders).print() // Empty setPython 示例env_settings EnvironmentSettings.in_batch_mode() table_env TableEnvironment.create(env_settings) # 注册一个 Orders 表 table_env.executeSql(CREATE TABLE Orders (user STRING, product STRING, amount INT) WITH (...)) # 插入原始数据 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 # 根据 where 条件删除 table_env.executeSql(DELETE FROM Orders WHERE user Lili).wait() table_env.executeSql(SELECT * FROM Orders).print() # 2 rows in set # 全表删除 table_env.executeSql(DELETE FROM Orders).wait() table_env.executeSql(SELECT * FROM Orders).print() # Empty set注意 Python API 中阻塞等待作业完成的方法是.wait()对应 Java/Scala 的.await()。SQL CLI 示例在 SQL CLI 中先通过SET语句切换到批模式再依次执行建表、插入与删除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从 CLI 输出可以看到DELETE 与 INSERT 一样是一条提交 SQL 更新语句Submitting SQL update statement的作业最终通过Job ID在集群上执行。底层实现机制行级删除的两条路径文档只要求目标表实现SupportsRowLevelDelete接口但在实际源码中Flink 为 DELETE 提供了两条实现路径理解这两条路径对连接器开发者至关重要。路径一SupportsDeletePushDown过滤器下推在 SupportsDeletePushDown.java 中Flink 允许把WHERE子句分解出的过滤器合取范式直接下推给 Sink由 Sink 直接删除数据规划阶段调用applyDeleteFilters(ListResolvedExpression filters)Sink 返回是否接受全部过滤器若返回true执行阶段调用executeDeletion()真正执行删除并返回预估删除行数未知时返回Optional.empty()。例如语句DELETE FROM t WHERE (a 1 OR a 2) AND b IS NOT NULL;会被分解为两个过滤器[a 1 OR a 2]与[b IS NOT NULL]Sink 只有能同时接受这两个过滤器时才会返回true。重要优先级规则当 Sink 同时实现了SupportsDeletePushDown与SupportsRowLevelDelete时只要applyDeleteFilters()返回trueplanner 总是优先使用SupportsDeletePushDown见 SupportsRowLevelDelete.java 的类注释。路径二SupportsRowLevelDelete行级删除当过滤器无法下推例如包含子查询、applyDeleteFilters()返回false或 Sink 未实现下推接口时若 Sink 实现了SupportsRowLevelDeleteFlink 会把 DELETE 语句重写为查询产出要删除的行或删除后的剩余行交给 Sink 消费。该接口的核心方法为RowLevelDeleteInfo applyRowLevelDelete(Nullable RowLevelModificationScanContext context);其中RowLevelDeleteInfo指导 planner 如何重写 DELETE 语句包含两个默认实现requiredColumns()Sink 执行删除所需的列集合返回Optional.empty()时表示需要全部列getRowLevelDeleteMode()返回删除模式默认DELETED_ROWS。RowLevelDeleteMode枚举SupportsRowLevelDelete.java定义了两种模式模式语义Sink 收到的数据行类型RowKindDELETED_ROWSSink 只收到需要被删除的行匹配过滤条件例DELETE FROM t WHERE y 2时收到满足y 2的行RowKind#DELETEREMAINING_ROWSSink 只收到删除后剩余的行不匹配过滤条件例DELETE FROM t WHERE y 2时收到不满足y 2的行RowKind#INSERT另外applyRowLevelDelete()的参数RowLevelModificationScanContext由实现了SupportsRowLevelModificationScan的表 Source 生成并传递用于在编译期实现 Source 与 Sink 之间的协调该上下文接口定义于 RowLevelModificationScanContext.java本身为空标记接口连接器可自行扩展。若 Source 未实现对应接口则该参数为null。规划期到执行期的规格化RowLevelDeleteSpec行级删除能力在规划期会被编码为 Sink 能力规格RowLevelDeleteSpecflink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/abilities/sink/RowLevelDeleteSpec.java它将rowLevelDeleteMode删除模式与requiredPhysicalColumnIndices所需物理列索引序列化为 JSON在执行期apply()时将RowLevelModificationScanContext重新回传给实现了SupportsRowLevelDelete的 Sink若 Sink 未实现该接口则抛出TableException。正是通过这个规格对象行级删除的相关信息才能随作业图序列化、分发到各 TaskManager 并在执行端恢复。常见异常与边界场景源码测试实证DeleteTableITCase.java 以test-update-delete测试连接器工厂类见 TestUpdateDeleteTableFactory.java支持delete-mode、support-delete-push-down、mix-delete、required-columns-for-delete、only-accept-equal-predicate等测试配置覆盖了大量 DELETE 场景是理解行为边界的最佳参考过滤器可下推但 Sink 不接受时抛异常测试中当表仅支持下推且WHERE含非等值谓词如a 1时DELETE 抛出UnsupportedOperationException错误信息为Cant perform delete operation of the table ... because the corresponding dynamic table sink has not yet implemented SupportsRowLevelDelete 全类名——这与文档未实现接口则抛异常的描述一致。行级删除支持子查询条件DELETE FROM t WHERE a (select count(1) from t where c 1)可以通过行级删除正常执行说明过滤器无法下推子查询时重写机制仍然有效。支持指定必需列与部分主键删除通过required-columns-for-delete配置如a;c验证了requiredColumns()语义复合主键表PRIMARY KEY (a, c)下配置a;b也能正确执行删除。混合模式mix-delete true下子查询条件走行级删除无WHERE的全表删除则回退到下推删除。StatementSet 限制StatementSet中不允许同时包含 INSERT 与 DELETE 语句会抛出TableException: Unsupported SQL query! Only accept a single SQL statement of type DELETE.同理compilePlanSql()只接受 INSERT 语句对 DELETE 也会抛异常。Legacy Sink 限制若目标表使用的是遗留TableSink如测试中的connector COLLECTIONDELETE 会抛出TableException提示需实现DynamicTableSink。此外RowLevelDeleteTest.java 通过verifyExplainInsert(DELETE FROM ...)对两种删除模式分别校验了 DELETE 语句重写后的执行计划可作为连接器开发者调试重写逻辑的参考。为自定义连接器接入 DELETE 能力结合以上源码分析自定义连接器要支持 DELETE可遵循以下步骤优先实现SupportsDeletePushDown若 Sink 能直接按过滤器删除数据实现applyDeleteFilters()返回过滤器接受情况接受时在executeDeletion()中执行实际删除。实现SupportsRowLevelDelete处理无法下推的场景实现applyRowLevelDelete()通过RowLevelDeleteInfo声明所需列与删除模式DELETED_ROWS或REMAINING_ROWS并消费对应 RowKind 的行数据完成删除。可选实现SupportsRowLevelModificationScanSource 侧需要向 Sink 传递扫描上下文信息时通过RowLevelModificationScanContext与 Sink 协同。注意批模式限定DELETE 仅在批模式下生效连接器接入时需确认作业运行模式。总结Flink 的 DELETE 语句提供了标准 SQL 形态的行级删除能力其核心约束是仅批模式 Sink 必须实现SupportsRowLevelDelete或可下推的SupportsDeletePushDown。本文从官方文档出发结合 SupportsRowLevelDelete.java、SupportsDeletePushDown.java 与 DeleteTableITCase.java 等源码与测试完整梳理了语法、四语言执行示例、两条底层实现路径与常见异常边界。由于 Flink 内置连接器尚未实现行级删除当前 DELETE 的主要受众是具备删除能力的自定义连接器与第三方连接器开发者对于内置表可通过INSERT OVERWRITE等方式实现数据覆盖更新。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐ScyllaDB CQL DELETE 语句完全指南行/列删除、范围删除与条件删除ScyllaDB CQL DELETE 语句完全指南行/列删除、范围删除与条件删除 导读 本文以 ScyllaDB 官方 CQL 文档中的 DELETE 章节数据库分布式数据库后端大数据TDengine 数据删除DELETE完全指南语法、删除标记机制与 SECURE_DELETE 安全删除TDengine 数据删除DELETE完全指南语法、删除标记机制与 SECURE_DELETE 安全删除 DELETE 语句是 TDengine 中按时间数据库时序数据库物联网大数据实时分析云原生现代C删除函数掌握delete语法的终极指南现代C删除函数掌握delete语法的终极指南 现代C删除函数delete语法是C11引入的强大特性它允许开发者显式禁用类的特定成员函数文档教程上一篇如何高效使用Qwen2-1.5B-Instruct10个实用技巧提升AI对话质量下一篇React Aria性能优化组件渲染性能与内存管理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表