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

Flink SQL 集合操作(Set Operations)完整指南:UNION / INTERSECT / EXCEPT / IN / EXISTS 的语法、语义与流式状态治理

发布时间:2026/9/24 14:53:49 来源:云帆数科 栏目:资讯中心
Flink SQL 集合操作(Set Operations)完整指南:UNION / INTERSECT / EXCEPT / IN / EXISTS 的语法、语义与流式状态治理
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载Flink SQL 的集合操作Set Operations用于将多个查询结果按行集合的方式进行合并、取交、取差或存在性判断是编写多路数据对比、去重合并、白名单过滤等查询时最常用的 SQL 能力之一。本文以 Flink 官方文档 set-ops.md 为骨架完整讲解UNION、INTERSECT、EXCEPT、IN、EXISTS在 Batch 与 Streaming 两种模式下的语义、SQL 写法与输出示例并结合当前仓库源码深入剖析优化器对IN/EXISTS/INTERSECT的重写机制以及流式查询下状态无限增长问题与table.exec.state.ttl配置的实战应对方案。集合操作总览Batch 与 Streaming 同时支持集合操作是 SQL 标准中的经典语法Flink SQL 对其支持标注为{{ label Batch }} {{ label Streaming }}即批处理和流处理两种执行模式下均可使用。这意味着无论作业以有界数据集Batch还是无界数据流Streaming方式运行都可以直接使用下列操作符操作符语义去重行为UNION两个表行的并集只保留不重复的行UNION ALL两个表行的并集不去重保留全部行INTERSECT两个表行的交集只保留不重复的行INTERSECT ALL两个表行的交集不去重保留匹配次数EXCEPT出现在左表但不在右表的行只保留不重复的行EXCEPT ALL出现在左表但不在右表的行不去重保留剩余次数IN表达式是否存在于子查询结果中—EXISTS子查询是否至少返回一行—核心规律不带ALL的变体对结果做去重distinct带ALL的变体保留重复行。这一规律贯穿全部三个二元集合操作符。UNION取两表并集UNION和UNION ALL返回在任意一张表中出现的行。区别在于UNION只保留不重复的行UNION ALL不去除结果中的重复行。以下示例创建两个视图t1与t2分别包含重复值然后对比两种写法的输出Flink SQL create view t1(s) as values (c), (a), (b), (b), (c); Flink SQL create view t2(s) as values (d), (e), (a), (b), (b); Flink SQL (SELECT s FROM t1) UNION (SELECT s FROM t2); --- | s| --- | c| | a| | b| | d| | e| --- Flink SQL (SELECT s FROM t1) UNION ALL (SELECT s FROM t2); --- | c| --- | c| | a| | b| | b| | c| | d| | e| | a| | b| | b| ---可以看到UNION将t1c, a, b, b, c与t2d, e, a, b, b合并后去重最终得到 5 个不重复的值{c, a, b, d, e}而UNION ALL直接拼接两边的全部 10 行重复值全部保留。注意示例中整个查询被括号包裹这是为了明确集合操作的输入边界Flink SQL 语法上允许对两个SELECT语句用括号包裹后再做集合运算。从执行层看UNION ALL直接对应一个 Union 节点。在 CommonExecUnion.java 中它通过 DataStream API 的UnionTransformation将多个输入流合并为一个流见createExecutionTransformation中返回new UnionTransformation(inputTransforms)的实现并派生出 StreamExecUnion.java 与 BatchExecUnion.java 两种运行时实现。而UNION去重版则等价于 Union 去重聚合会被优化器进一步改写为聚合运算。INTERSECT取两表交集INTERSECT和INTERSECT ALL返回同时出现在两张表中的行。INTERSECT只保留不重复的行INTERSECT ALL不去除重复并遵循两边各出现几次就返回几次的匹配语义Flink SQL (SELECT s FROM t1) INTERSECT (SELECT s FROM t2); --- | s| --- | a| | b| --- Flink SQL (SELECT s FROM t1) INTERSECT ALL (SELECT s FROM t2); --- | s| --- | a| | b| | b| ---两个视图中共同出现的值只有a和b因此去重版交集结果为{a, b}。而INTERSECT ALL需要考虑出现次数b在t1中出现 2 次、在t2中出现 2 次取较小值 2因此结果中出现两行ba两边各出现 1 次结果为一行a。EXCEPT取两表差集EXCEPT和EXCEPT ALL返回出现在左表中但不在右表中的行。EXCEPT去重EXCEPT ALL按左表出现次数减去右表出现次数的剩余次数返回Flink SQL (SELECT s FROM t1) EXCEPT (SELECT s FROM t2); --- | s | --- | c | --- Flink SQL (SELECT s FROM t1) EXCEPT ALL (SELECT s FROM t2); --- | s | --- | c | | c | ---t1中独有的值是c出现 2 次去重版EXCEPT只返回一行cEXCEPT ALL中c在左表出现 2 次、右表出现 0 次剩余 2 次因此返回两行c。b虽然两边都有但右表出现次数2不少于左表2剩余次数为 0所以不出现在结果中。提示Flink SQL 使用EXCEPT关键字表示差集与 PostgreSQL 一致部分数据库使用MINUS二者语义相同但 Flink 语法层面采用EXCEPT。IN判断表达式是否存在于子查询结果中IN返回 true当且仅当左侧表达式存在于给定子查询的结果表中。语法要求子查询结果表必须只有一列且该列的数据类型必须与左侧表达式一致SELECT user, amount FROM Orders WHERE product IN ( SELECT product FROM NewProducts )该查询的含义是筛选出product出现在NewProducts表中的所有订单行等价于基于product键的半连接semi join过滤。优化器行为Flink 优化器会把IN条件重写为 join group 操作具体表现为 SEMI JOIN 配合聚合去重防止右表重复键导致结果行被放大。流式查询注意由于需要为 join 与 group 维护状态流式查询计算结果所需的状态大小会随输入中的 distinct 行数无限增长。缓解手段是为查询配置一个合适的状态存活时间State TTL设置table.exec.state.ttl防止状态无限膨胀但要注意设置 TTL 可能影响查询结果的正确性过期的状态被清理后迟到的数据将无法正确参与匹配完整参数说明参见 Query Configuration 与 流式概念文档。EXISTS判断子查询是否至少返回一行EXISTS返回 true当且仅当子查询至少返回一行SELECT user, amount FROM Orders WHERE product EXISTS ( SELECT product FROM NewProducts )EXISTS的语义是子查询非空即命中与IN的区别在于它不要求左侧表达式与子查询列做等值匹配只关心子查询结果是否有行。但 Flink SQL 对EXISTS的支持有一个前提仅当该操作可以被重写为 join group 操作时才支持即不能在所有场景下都使用非等值关联或无法改写的场景可能不被接受。与IN相同优化器会把EXISTS重写为 join group 操作流式查询同样面临状态无限增长的问题需要借助table.exec.state.ttl等配置治理状态规模同时警惕 TTL 对结果正确性的影响。详细配置入口同样指向 Query Configuration。流式查询下的状态治理table.exec.state.ttl 与 STATE_TTL Hint针对IN、EXISTS被重写为 join group 后状态无限增长的问题Flink 提供了两个层面的状态 TTL 治理手段。1. 作业级配置table.exec.state.ttltable.exec.state.ttl是流处理模式下标签为Streaming的 Duration 类型参数默认值为0 ms语义为空闲状态即长时间未更新的状态最短保留时间状态在空闲时间达到该时长之前绝不会被清理达到之后会在某个时间点被清理。默认值 0 表示永不清除状态。官方配置描述还指出清理状态需要额外的簿记开销bookkeeping overhead因此该参数需要按需设置。可以通过以下任一方式配置SQL Client / SQL Gateway 会话级设置SET table.exec.state.ttl 1000;这一用法在 SQL Client 初始化脚本示例 中有明确体现SET table.exec.state.ttl 1000; -- optional: table programs idle state time。Java / Scala / Python Table API通过EnvironmentSettings传入Configuration或通过TableEnvironment#getConfig()获取的TableConfig设置底层 key-value 选项示例见 Configuration 文档Configuration configuration new Configuration(); configuration.setString(table.exec.state.ttl, 1 h); EnvironmentSettings settings EnvironmentSettings.newInstance() .inStreamingMode().withConfiguration(configuration).build(); TableEnvironment tEnv TableEnvironment.create(settings);configuration Configuration() configuration.set(table.exec.state.ttl, 1 h) settings EnvironmentSettings.new_instance() \ .in_streaming_mode() \ .with_configuration(configuration) \ .build() t_env TableEnvironment.create(settings)2. 算子级配置STATE_TTL Hint对于有状态计算的 Regular Join 与 Group Aggregation用户还可以使用STATE_TTLhint 指定算子级的空闲状态维持时间从而让特定算子使用与作业级table.exec.state.ttl不同的 TTL 值详见 Hints 文档-- 表名作为 hint key SELECT /* STATE_TTL(orders3d, lineitem1d) */ * FROM orders LEFT JOIN lineitem ON orders.o_orderkey lineitem.l_orderkey; -- 表别名作为 hint key一旦指定别名必须使用别名 SELECT /* STATE_TTL(o3d, l1d) */ * FROM orders o LEFT JOIN lineitem l ON o.o_orderkey l.l_orderkey; -- 级联 Join 分别为每个参与表设置 TTL SELECT /* STATE_TTL(o 3d, l 1d, c 10d) */ * FROM orders o LEFT OUTER JOIN lineitem l ON o.o_orderkey l.l_orderkey LEFT OUTER JOIN customers c ON o.o_custkey c.c_custkey;需要强调的是STATE_TTLhint 目前面向的是 Regular Join 与 Group Aggregation 这两类有状态算子而IN/EXISTS重写后的 join group 结构同样可以借助这一机制为对应算子设定 TTL。此外基于窗口的操作如 Window Join、Window Aggregation、Interval Join不依赖table.exec.state.ttl控制状态保留其状态 TTL 无法在算子级别配置参见 流式概念文档。3. 正确性与状态的权衡无论采用作业级还是算子级 TTL都必须明确TTL 是正确性换资源的权衡。设置过小的 TTL 会显著缩减状态规模、降低内存与磁盘压力但会使空闲时间超过 TTL的状态被清理导致后续迟到的数据无法与已被清理的历史状态正确关联从而产出错误结果。对于IN/EXISTS/INTERSECT/EXCEPT这类对历史数据敏感的集合与半连接操作建议结合业务数据的时间窗口特征谨慎选择 TTL仅在确认数据迟到范围有限时才使用。源码视角优化器如何重写集合操作理解集合操作的底层执行有助于预判查询性能与状态开销。当前仓库的 planner 代码清晰地展示了三类关键重写INTERSECT → SEMI JOIN Aggregate在 ReplaceIntersectWithSemiJoinRule.java 中优化器将去重版INTERSECT!intersect.all getInputs().size() 2即仅处理两个输入的 case重写为先对两输入按全部字段生成等值条件执行JoinRelType.SEMI半连接再在全部键上做aggregate(groupKey(...))去重。注释明确说明Planner rule that replaces distinct Intersect with a distinct Aggregate on a SEMI Join.这印证了文档中优化器将条件重写为 join 和 group 操作的说法——INTERSECT与IN/EXISTS最终都收敛到 join group 的执行形态因此它们的流式状态成本模型是一致的。UNION → UnionTransformation在 CommonExecUnion.java 中Union 运行时节点通过UnionTransformation把多个输入变换合并为单一输出流是UNION ALL最直接、成本最低的实现而去重版UNION则在其基础上叠加去重聚合。SET 配置语句的解析配置项如SET table.exec.state.ttl 1000在 SQL Client 中由 SetOperationParseStrategy.java 解析它用正则SET(\s(?key[^\s])\s*\s*((?quotedVal[^]*)|(?val[^;\s])))?\s*;?匹配SET key value或裸SET查看全部配置两种形式将带引号与不带引号的值统一转换为SetOperation(key, value)。这也解释了为什么 SQL 中配置值既可以用单引号包裹也可以直接书写。使用建议与注意事项去重成本UNION/INTERSECT/EXCEPT不带ALL需要额外的聚合去重开销高于对应的ALL变体。如果业务上确认输入已经无重复应优先使用ALL版本。列数一致性集合操作要求两侧查询的列数一致IN的子查询则必须只输出一列且类型与左侧表达式严格一致。括号分隔多个集合操作组合时建议用括号明确运算顺序与输入边界避免语义歧义。流式状态所有涉及 join group 重写的操作IN、EXISTS、去重版集合操作在流模式下都会积累状态务必结合table.exec.state.ttl或STATE_TTLhint 规划状态规模并接受由此带来的正确性边界。EXISTS 支持范围EXISTS仅在能够被重写为 join group 的场景下受支持编写查询时应确认关联形式可被优化器改写。小结集合操作是 Flink SQL 在 Batch 与 Streaming 双模式下原生支持的一类查询能力UNION [ALL]、INTERSECT [ALL]、EXCEPT [ALL]提供标准的并/交/差语义IN与EXISTS提供子查询存在性判断。理解不带 ALL 去重、带 ALL 保重的语义规律、掌握IN/EXISTS/INTERSECT被重写为 join group 的执行模型并正确配置table.exec.state.ttl或STATE_TTLhint 来治理流式状态是写出既正确又可控的集合查询的关键。更多相关配置与概念可继续阅读 Query Configuration、流式概念、Hints 与 SQL Client 文档。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink SQL 集合操作Set Operations完全指南UNION / INTERSECT / EXCEPT / IN / EXISTS 的语法、去重语义与流式计算原理Flink SQL 集合操作Set Operations完全指南UNION / INTERSECT / EXCEPT / IN / EXISTS 的语法、大数据流处理批处理数据工程Flink Hive Dialect 集合操作Set Operations完整指南UNION / INTERSECT / EXCEPT 语法与实战Flink Hive Dialect 集合操作Set Operations完整指南UNION / INTERSECT / EXCEPT 语法与实战 Set大数据流处理批处理数据工程Apache Spark SQL 集合运算Set Operators完全指南UNION / INTERSECT / EXCEPT 语法、去重语义与执行原理Apache Spark SQL 集合运算Set Operators完全指南UNION / INTERSECT / EXCEPT 语法、去重语义与执行原理大数据数据分析批处理流处理机器学习图计算上一篇Mobile Security Framework (MobSF) 自定义工作流自动化安全检测流程编排终极指南下一篇pytorch-grad-cam与模型监控生产环境解释性告警系统创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关推荐

Salt Proxy 执行模块:在 Minion 上自动部署与管理 salt-proxy 进程
Salt Proxy 执行模块:在 Minion 上自动部署与管理 salt-proxy 进程

运维配置管理后端 【免费下载链接】salt Software to automate the management and configuration of infrastructure and applications at scale. 项目地址: https://gitcode.com/gh_mirrors/sa/salt 点击查看 免费下载 导读 salt_proxy 是 Salt 提供的执行模块&… · 2026/9/24 14:53:49

AI Agent 场景应用:基于 CodeGuide 智能体脚手架打造 ai + draw.io 交互式绘图产品
AI Agent 场景应用:基于 CodeGuide 智能体脚手架打造 ai + draw.io 交互式绘图产品

文档教程后端 【免费下载链接】CodeGuide :books: 本代码库是作者小傅哥多年从事一线互联网 Java 开发的学习历程技术汇总,旨在为大家提供一个清晰详细的学习教程,侧重点更倾向编写Java核心内容。如果本仓库能为您提供帮助,请给予支持(关注、… · 2026/9/24 14:53:49

KuGouMusicApi Vercel 部署实战:Serverless免费托管酷狗音乐API,零成本上线
KuGouMusicApi Vercel 部署实战:Serverless免费托管酷狗音乐API,零成本上线

KuGouMusicApi Vercel 部署实战:Serverless免费托管酷狗音乐API,零成本上线 【免费下载链接】KuGouMusicApi 酷狗音乐 Node.js API service 项目地址: https://gitcode.com/gh_mirrors/ku/KuGouMusicApi KuGouMusicApi 是一款酷狗音乐 Node.js AP… · 2026/9/24 14:53:42

PHPStan 纯函数错误标识符 `pureFunction.parameterByRef` 全解析:引用参数与纯度契约冲突
PHPStan 纯函数错误标识符 `pureFunction.parameterByRef` 全解析:引用参数与纯度契约冲突

开发工具代码质量静态分析 【免费下载链接】phpstan PHP Static Analysis Tool - discover bugs in your code without running it! 项目地址: https://gitcode.com/gh_mirrors/ph/phpstan 点击查看 免费下载 导读 pureFunction.parameterByRef 是 PHPStan 在 2.0… · 2026/9/24 15:25:23

深度解析 Druid:一款 data-first 的 Rust 原生 UI 工具包(基于 0.8.3 源码)
深度解析 Druid:一款 data-first 的 Rust 原生 UI 工具包(基于 0.8.3 源码)

跨平台桌面应用UI组件 【免费下载链接】druid A data-first Rust-native UI design toolkit. 项目地址: https://gitcode.com/gh_mirrors/drui/druid 点击查看 免费下载 本篇以仓库根目录 README.md 为骨架,结合 druid 源码、druid-shell 与 druid-der… · 2026/9/24 15:25:23

DevilutionX(暗黑破坏神 1)GKD350h 移植版:安装、按键映射与已知问题全解析
DevilutionX(暗黑破坏神 1)GKD350h 移植版:安装、按键映射与已知问题全解析

游戏开发 【免费下载链接】DevilutionX Diablo build for modern operating systems 项目地址: https://gitcode.com/gh_mirrors/de/DevilutionX 点击查看 免费下载 DevilutionX 是《暗黑破坏神 1》面向现代操作系统的开源重制移植,本指南聚焦其在 GKD3… · 2026/9/24 15:25:23

命题专家答题方法哪里学 技巧答题赛道哪家更靠谱
命题专家答题方法哪里学 技巧答题赛道哪家更靠谱

不少初高中家长都会遇到这样的困惑:孩子课本上的知识都能看懂,日常学习也投入不少精力,但面对试卷作答时,很难把自身积累充分展现出来。刷了大量习题,学习状态却没有明显改观,这往往不是孩子不够用心&#… · 2026/9/24 15:25:23

勾芡香菇汤做法详解:一碗汤品从菜谱到 RAG 知识库的完整旅程
勾芡香菇汤做法详解:一碗汤品从菜谱到 RAG 知识库的完整旅程

教程人工智能大模型RAG 【免费下载链接】all-in-rag 🔍大模型应用开发实战一:RAG 技术全栈指南,在线阅读地址:https://datawhalechina.github.io/all-in-rag/ 项目地址: https://gitcode.com/datawhalechina/all-in-ra… · 2026/9/24 15:25:17

golang-jwt/jwt v4 版本演进全解析:从 jwt-go 迁移到 Go JWT 库的 4.0 时代
golang-jwt/jwt v4 版本演进全解析:从 jwt-go 迁移到 Go JWT 库的 4.0 时代

golang-jwt/jwt v4 版本演进全解析:从 jwt-go 迁移到 Go JWT 库的 4.0 时代 【免费下载链接】sliver Adversary Emulation Framework 项目地址: https://gitcode.com/gh_mirrors/sl/sliver 本文以当前仓库中 vendor/github.com/golang-jwt/jwt/v4/VERSION_HI… · 2026/9/24 15:25:17

基于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

了解更多?预约专属演示

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

企业微信二维码