首页/新闻资讯/正文详情

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

发布时间:2026/9/24 15:36:05 来源:云帆数科 栏目:资讯中心
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),仅供参考

相关推荐

F´ 框架中的 Svc::FatalHandler 组件:FATAL 事件处理与平台差异化实现解析
F´ 框架中的 Svc::FatalHandler 组件:FATAL 事件处理与平台差异化实现解析

嵌入式系统编程 【免费下载链接】fprime F - A flight software and embedded systems framework 项目地址: https://gitcode.com/gh_mirrors/fp/fprime 点击查看 免费下载 导读 Svc::FatalHandler 是 F(F Prime)飞行软件与嵌入式系统框架中… · 2026/9/24 15:35:52

OpenChamber 1.2.9 会话自动清理:可配置保留策略、全端同步与长期会话优化
OpenChamber 1.2.9 会话自动清理:可配置保留策略、全端同步与长期会话优化

AI Agent人工智能代码智能体交互助手 【免费下载链接】openchamber Agentic Development Environment based on OpenCode AI agent 项目地址: https://gitcode.com/gh_mirrors/op/openchamber 点击查看 免费下载 导读 OpenChamber 1.2.9 的核心主题是 会话生命周期… · 2026/9/24 15:35:46

Redwood Cells 完全指南:用声明式约定接管 GraphQL 数据获取与生命周期
Redwood Cells 完全指南:用声明式约定接管 GraphQL 数据获取与生命周期

后端前端Web框架开发工具 【免费下载链接】redwood RedwoodGraphQL 项目地址: https://gitcode.com/gh_mirrors/re/redwood 点击查看 免费下载 导读 Cells 是 Redwood 最具标志性的数据获取抽象:它用一套命名导出约定(QUERY、Loading、Empt… · 2026/9/24 15:35:46

Apache Thrift Rust crate 发布指南:从 crates.io 账户配置到 `cargo publish` 全流程
Apache Thrift Rust crate 发布指南:从 crates.io 账户配置到 `cargo publish` 全流程

Apache Thrift Rust crate 发布指南:从 crates.io 账户配置到 cargo publish 全流程 【免费下载链接】thrift Apache Thrift 项目地址: https://gitcode.com/gh_mirrors/thrift2/thrift Apache Thrift 的 Rust 运行时库(thrift crate)… · 2026/9/24 16:04:17

WeMod 专业版免费解锁:Wand-Enhancer 本地补丁一次性讲清
WeMod 专业版免费解锁:Wand-Enhancer 本地补丁一次性讲清

WeMod 专业版免费解锁:Wand-Enhancer 本地补丁一次性讲清 【免费下载链接】Wand-Enhancer Advanced UX and interoperability extension for Wand (WeMod) app 项目地址: https://gitcode.com/GitHub_Trending/we/Wand-Enhancer Wand-Enhancer 是一款开源的本… · 2026/9/24 16:04:17

wp-calypso 图片预加载组件 ImagePreloader:实现占位符过渡与加载状态管理
wp-calypso 图片预加载组件 ImagePreloader:实现占位符过渡与加载状态管理

前端CMS 【免费下载链接】wp-calypso The JavaScript and API powered WordPress.com 项目地址: https://gitcode.com/gh_mirrors/wp/wp-calypso 点击查看 免费下载 导读 ImagePreloader 是 wp-calypso(WordPress.com 的前端单页应用)内置的… · 2026/9/24 16:04:11

基于 task-provider-sample 深入解析 VS Code Task Provider API:从 Rakefile 自动检测到自定义构建任务
基于 task-provider-sample 深入解析 VS Code Task Provider API:从 Rakefile 自动检测到自定义构建任务

示例工程 【免费下载链接】vscode-extension-samples Sample code illustrating the VS Code extension API. 项目地址: https://gitcode.com/gh_mirrors/vs/vscode-extension-samples 点击查看 免费下载 导读 本篇文章以 vscode-extension-samples 仓库中的 task… · 2026/9/24 16:04:11

sinon.match.func 匹配器指南:在 Sinon 中断言函数类型参数
sinon.match.func 匹配器指南:在 Sinon 中断言函数类型参数

sinon.match.func 匹配器指南:在 Sinon 中断言函数类型参数 【免费下载链接】sinon Test spies, stubs and mocks for JavaScript. 项目地址: https://gitcode.com/gh_mirrors/si/sinon sinon.match.func 是 Sinon 匹配器(Matcher)体系… · 2026/9/24 16:04:04

一条15秒 AI 短剧,为什么能把创作者熬到凌晨三点?
一条15秒 AI 短剧,为什么能把创作者熬到凌晨三点?

观众只看 15 秒,创作者得“跑”一整晚。凌晨一点,手机屏幕上还在自动播放。上一集里,男主刚刚失忆;下一集里,女主已经带着证据冲进了会议室。你本来只想看一眼,结果抬头,充电线都快和枕头打结了… · 2026/9/24 16:03:58

基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程
基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程

简介:这是一套面向计算机、人工智能、自动化等专业学生与教师的毕业设计级项目资源,围绕YOLOv8实现渔船作业监控系统,可用于毕设、课程设计、大作业或项目立项演示。压缩包共97个文件,约24.21MB,以70个Python源码文件为… · 2026/9/24 0:00:13

1D-CNN时间序列建模实战:从Conv1d原理到工业落地
1D-CNN时间序列建模实战:从Conv1d原理到工业落地

简介:面向时间序列数据建模的一维卷积神经网络完整实现,适合深度学习入门者及需要快速验证时序模型的研究者,能够从音频、文本、传感器或股价等序列中挖掘局部特征与时间依赖。压缩包体积很小,只有3KB,内含3个Python脚… · 2026/9/24 0:00:26

柔软的L:汉语语流中被忽视的舌肌张力控制
柔软的L:汉语语流中被忽视的舌肌张力控制

1. 这个“L”不是字母表里的L,而是舌尖上的L最近在几个方言群和语音教学社群里,反复看到有人发一句:“也说字母L:柔软的长舌”。初看以为是英语发音课笔记,点开才发现全是方言爱好者、播音系学生、语言康复师甚至戏曲演… · 2026/9/24 0:00:44

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码