尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Flink SQL DELETE 语句详解:行级删除原理、SupportsRowLevelDelete 接口与实战示例

Flink SQL DELETE 语句详解:行级删除原理、SupportsRowLevelDelete 接口与实战示例 大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/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才能消费行级删除产生的行数据未实现接口时抛异常若对未实现相关接口的表执行DELETEPlanner 会抛出异常此外截至目前 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 重写为SinkModifyOperationModifyType.DELETE把“删除哪些行”的问题转化为“查询出哪些行并交给 sink 消费”的问题。PlannerQueryOperation中显式抛出TableException(Delete statements are not SQL serializable.)说明该查询仅供内部重写使用不可序列化回 SQL。这一设计印证了接口 javadoc 中的优先级约定当表 sink 同时实现SupportsDeletePushDown与SupportsRowLevelDelete时只要applyDeleteFilters返回 truePlanner 总是优先使用删除下推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_ROWSsink 只收到匹配过滤条件、需要被删除的行。这些行统一携带RowKind.DELETE语义REMAINING_ROWSsink 收到的是删除后剩余的即不匹配过滤条件的行统一携带RowKind.INSERT语义。适合“整表重写”类存储例如以覆盖方式重写文件的连接器。以DELETE FROM t WHERE y 2;为例若返回DELETED_ROWSsink 会收到满足y 2的行若返回REMAINING_ROWSsink 会收到不满足y 2的行参见接口 javadoc 中的示例说明。序列化与反序列化RowLevelDeleteSpecPlanner 在将重写后的计划提交执行时需要把 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 进一步验证了行级删除、带分区列的删除、删除与插入混用的StatementSetstatementSet.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项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink Table SQL DELETE 语句完全指南行级删除、语法与连接器实现机制Flink Table SQL DELETE 语句完全指南行级删除、语法与连接器实现机制 DELETE 是 Flink Table API SQL 提供的大数据流处理批处理数据工程Flink 窗口去重Window DeduplicationSQL 详解语法、示例与实现原理Flink 窗口去重Window DeduplicationSQL 详解语法、示例与实现原理 窗口去重Window Deduplication是 Fl大数据流处理批处理数据工程STL到STEP转换引擎打破3D打印与精密制造间的格式壁垒STL到STEP转换引擎打破3D打印与精密制造间的格式壁垒 在数字化设计与制造领域工程师们长期面临着一个技术难题如何将3D打印中广泛使用的STL格式无缝转大数据流处理批处理数据工程上一篇星穹铁道智能工具技术赋能游戏效率提升的全自动化解决方案下一篇Tink C 集成 BoringCryptoBoringSSL FIPS 验证模块完全指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表