资讯动态

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

发布时间:2026/9/24 14:53:48 来源:尧图企业网站定制
大数据流处理批处理数据工程【免费下载链接】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),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价