上一篇讲了 Table API 的整体原理和表达式体系。这篇深入 Table API 中最基础、最常用的两个操作——查询(select 投影)与过滤(filter/where)。
几乎每一个 Flink SQL 作业都会用到 select 和 filter,但很多同学只是简单地写select(...).filter(...),对背后的原理、优化机制、性能影响知之甚少。实际上,select 和 filter 是性能优化的关键抓手——投影裁剪可以减少 50% 以上的数据读取量,谓词下推可以将过滤推到数据源端执行,条件重排和短路求值可以减少大量无效计算。
这篇从本质定义、四大操作详解、执行顺序对比、六步内部执行流程、四大优化原理、表达式计算与短路求值、完整代码实现、六个常见坑、八条最佳实践,把查询与过滤一次性讲透。
一、查询与过滤本质定义
查询与过滤的本质一句话概括:select(投影)负责列的选择和计算,filter(过滤)负责行的筛选,两者配合从原始表中提取出需要的数据子集,是关系代数中最基础的两个操作。
下面这张图把查询与过滤的本质定义、四大核心操作、select 与 filter 执行顺序对比放在一起展示。
从关系代数的角度看,select 对应投影操作(π),filter 对应选择操作(σ)。这两个操作是所有复杂查询的基础——聚合、连接、窗口等高级操作都建立在投影和过滤之上。
1.1 执行顺序
在逻辑执行计划中,filter 通常先于 select 执行(先过滤行,再计算列),这样可以减少需要计算的数据量。优化器会自动进行谓词下推(Predicate Pushdown)和投影裁剪(Projection Pushdown),将过滤和投影尽可能推到数据源端执行,减少数据传输和计算量。
即使你写的是table.select(...).filter(...)(select 在前),优化器也会自动将 filter 提到 select 之前执行。这是查询优化器的基本能力,不需要手动调整顺序。
1.2 无状态算子
Projection 和 Filter 都是无状态算子,不需要维护状态,逐行处理:
- Filter 算子:对每一行计算条件表达式,返回 true 的行向下游传递,false 的行丢弃
- Projection 算子:对每一行计算所有输出列的表达式,生成新的行向下游传递
因为无状态,这两个算子的性能非常高,也不会有状态膨胀的问题(distinct 除外,distinct 需要状态)。
二、四大操作详解
2.1 select 投影(Projection)
作用:选择列、计算表达式、重命名列。
方法:
select($("a"), $("b")):选择列select($("a").plus($("b")).as("sum")):表达式计算 + 重命名select($("*")):选择所有列(不推荐)
原理:生成 Projection 算子,对每一行输入计算输出列。表达式在编译期解析为 RexNode,运行时通过 CodeGen 生成专门的 Java 代码,直接操作 BinaryRowData 的字段,性能接近手写 Java 代码。
优化:投影裁剪(Projection Pushdown)将不需要的列裁剪掉,推到数据源端。列式存储(Parquet/ORC)可以只读取需要的列,大幅减少 IO 和反序列化开销。
2.2 filter/where 过滤(Filter)
作用:根据条件表达式筛选行,只保留满足条件的行。
方法:
filter($("a").isGreater(10))where($("a").isGreater(10).and($("b").isEqual("x")))
filter 和 where 完全等价,where 是 SQL 风格的别名。
原理:生成 Filter 算子,对每一行计算条件表达式,返回 true 的行向下游传递,false 的行丢弃。条件表达式同样通过 CodeGen 生成专门的 Java 代码。
优化:谓词下推(Predicate Pushdown)将过滤条件推到数据源端,支持谓词下推的连接器(JDBC/Hive/Parquet)可以在数据源端执行过滤,大幅减少数据传输。
2.3 distinct 去重(Deduplication)
作用:去除重复行,只保留唯一行。
方法:table.distinct()
原理:
- 批处理:基于排序或哈希去重,一次性处理所有数据
- 流处理:基于状态去重,维护一个已见行的集合(状态),每来一行检查是否已存在,不存在则输出并加入状态,存在则丢弃
注意:流处理 distinct 需要状态管理,状态会持续增长,需要设置状态 TTL(table.exec.state.ttl)防止 OOM。distinct 是全局去重,不是分组去重(分组去重用 groupBy)。
2.4 limit/offset 分页(Pagination)
作用:限制返回行数,实现分页查询。
方法:
table.limit(10):取前 10 行table.orderBy($("a").desc()).fetch(10).offset(20):排序后分页
原理:
- 批处理:全局排序后取前 N 行,需要将所有数据拉到一个节点排序(全局排序开销大)
- 流处理:limit 不直接支持无界流(因为流是无限的,无法确定"前 N 行"),需要配合窗口或 ROW_NUMBER() 使用
注意:limit 是批处理操作,流处理中使用需要谨慎。offset 会跳过前 N 行,大数据量下 offset 较大时性能差(需要扫描并丢弃前 N 行)。
三、select 与 filter 执行顺序对比
| 对比维度 | select 投影 | filter 过滤 |
|---|---|---|
| 关系代数 | 投影操作 π | 选择操作 σ |
| 操作对象 | 列(垂直裁剪) | 行(水平筛选) |
| 逻辑执行顺序 | 通常在 filter 之后执行 | 通常在 select 之前执行 |
| 优化策略 | 投影裁剪(Projection Pushdown) | 谓词下推(Predicate Pushdown) |
| 状态需求 | 无状态(逐行计算) | 无状态(逐行判断) |
| 流处理支持 | 完全支持(无界流逐行投影) | 完全支持(无界流逐行过滤) |
| 性能影响 | 减少列数减少序列化/网络传输,表达式计算增加 CPU | 过滤率越高下游数据量越小,条件计算增加 CPU |
| 最佳实践 | 只选择需要的列,避免 select(*) | 过滤条件尽早执行,利用谓词下推,避免复杂 UDF |
核心结论:filter 先于 select 执行是优化器的默认行为,不需要手动调整。但在写代码时,应该有意识地将过滤条件尽早写出来(在聚合、连接之前),让优化器有更多优化空间。
四、查询与过滤内部执行流程
理解了本质和操作分类,下面深入内部执行流程,看 select 和 filter 的方法调用到底是怎么变成执行代码的。
下面这张图把六步执行流程、四大优化原理、表达式计算与短路求值放在一起展示。
4.1 六步执行流程
一条 select/filter 链式调用从输入到执行,完整的流程分为六步:
- 方法调用:用户调用 select/filter/distinct 等方法
- 构建 Operation:每个方法构建一个 QueryOperation 节点
- RelNode 转换:Operation 翻译为 Calcite LogicalProject/LogicalFilter
- 优化器优化:谓词下推/投影裁剪/常量折叠/条件重排
- 代码生成:CodeGen 生成 Projection/Filter 算子的 Java 代码
- 逐行执行:算子逐行处理:Filter 判断条件 → Projection 计算列
4.2 关键理解
select 和 filter 在逻辑执行计划中分别对应LogicalProject和LogicalFilter两个 RelNode 节点。优化器会对这两个节点进行多种优化。
代码生成(CodeGen)是性能的关键:Flink 会为 Projection 和 Filter 算子生成专门的 Java 代码,而不是用解释执行。生成的代码直接操作二进制行数据(BinaryRowData),避免了对象创建和虚方法调用,性能非常高。这也是 Flink Table API/SQL 性能优秀的关键原因之一。
比如$("amount").times(lit(0.8)).as("discount")会被翻译为类似out.setDecimal(0, in.getDecimal(2).multiply(0.8))这样的直接字段操作代码,没有任何抽象层开销。
五、四大优化原理
5.1 谓词下推(Predicate Pushdown)
原理:将 Filter 条件尽可能推到数据源端执行,减少从数据源读取的数据量。
效果:支持谓词下推的连接器(JDBC/Hive/Parquet)可以在数据源端执行过滤,只读取满足条件的数据,大幅减少网络传输和计算量。
示例:SELECT * FROM orders WHERE amount > 100,谓词 amount > 100 下推到 MySQL,MySQL 只返回 amount > 100 的行。
限制:不是所有谓词都能下推,包含 UDF、复杂表达式的谓词可能无法下推。用explain()查看执行计划,确认谓词是否下推。
5.2 投影裁剪(Projection Pushdown)
原理:将 Projection(列选择)推到数据源端,只读取需要的列,不需要的列直接裁剪掉。
效果:列式存储(Parquet/ORC)可以只读取需要的列,大幅减少 IO 和反序列化开销。行式存储(JSON/CSV)也可以在读取后立即裁剪不需要的列。
示例:SELECT user_id, amount FROM orders,只读取 user_id 和 amount 两列,其他列(order_id/status/order_time)直接裁剪。
最佳实践:永远不要用select(*)或SELECT *,只选择需要的列。这是最简单也最有效的性能优化手段。
5.3 常量折叠(Constant Folding)
原理:在优化阶段直接计算只包含常量的表达式,不需要在运行时逐行计算。
效果:减少运行时计算量,表达式结果在编译期就确定了。
示例:SELECT amount * 2, 1 + 2 FROM orders,常量表达式 1 + 2 在优化期直接计算为 3,运行时不需要计算。
5.4 条件重排(Condition Reordering)
原理:优化器会根据条件的选择性和计算成本重新排列 AND 条件的顺序,将高选择性(过滤率高)、低成本的条件放在前面,利用短路求值减少计算量。
效果:高选择性条件先执行,大部分行在第一个条件就被过滤掉,不需要计算后续的复杂条件。
注意:优化器的条件重排是基于统计信息的,如果没有统计信息(如流处理),优化器可能不会重排,此时需要手动将高选择性条件放在前面。
六、表达式计算与短路求值
6.1 表达式编译(RexNode → CodeGen)
Table API 中的 Expression(如$("a").plus($("b")))在编译期被解析为 Calcite 的 RexNode(行表达式节点),然后在代码生成阶段被翻译为专门的 Java 代码。生成的代码直接操作 BinaryRowData 的字段,避免了表达式解释执行的开销,性能接近手写 Java 代码。
6.2 短路求值(Short-circuit Evaluation)
原理:AND 条件中,如果前面的条件返回 false,后面的条件不会执行(因为整个表达式已经确定为 false)。OR 条件中,如果前面的条件返回 true,后面的条件不会执行。
效果:减少不必要的条件计算,尤其是包含 UDF 或复杂计算的条件。
最佳实践:将高选择性、低成本的条件放在 AND 的前面,将低选择性、高成本的条件(如 UDF)放在后面。这样大部分行在第一个条件就被过滤掉,不需要执行昂贵的 UDF。
6.3 NULL 处理(三值逻辑)
SQL 采用三值逻辑(TRUE/FALSE/UNKNOWN),任何与 NULL 的比较结果都是 UNKNOWN,UNKNOWN 在 WHERE 中被视为 false(行被过滤掉)。
常见坑:$("a").isEqual(lit(null))永远返回 UNKNOWN(不会匹配任何行),判断 NULL 必须用$("a").isNull()。
如果需要保留 NULL 行,用WHERE a > 10 OR a IS NULL。
6.4 过滤率与性能
过滤率(Filter Selectivity)= 过滤后行数 / 原始行数。过滤率越低(过滤掉的行越多),下游数据量越小,整体性能越好。
Filter 是减少数据量最有效的算子,一个高过滤率的 Filter 可以大幅减少下游聚合、连接、窗口等算子的计算量和状态大小。
最佳实践:
- 尽早执行 Filter(在聚合/连接之前)
- 利用谓词下推将 Filter 推到数据源端
- 避免在 Filter 中使用复杂 UDF(UDF 无法谓词下推,且计算成本高)
- 用
explain()确认 Filter 的位置和谓词下推情况
七、完整代码实现
下面是一个完整的查询与过滤使用示例,包含创建环境、DDL 建表、select 投影、filter 过滤、组合查询、explain 检查、输出执行。
下面这张图把完整代码实现、六个常见坑、八条最佳实践放在一起展示。
importorg.apache.flink.streaming.api.environment.StreamExecutionEnvironment;importorg.apache.flink.table.api.Table;importorg.apache.flink.table.api.TableResult;importorg.apache.flink.table.api.bridge.java.StreamTableEnvironment;importjava.time.ZoneId;importstaticorg.apache.flink.table.api.Expressions.*;publicclassTableQueryExample{publicstaticvoidmain(String[]args)throwsException{// 1. 创建执行环境StreamExecutionEnvironmentenv=StreamExecutionEnvironment.getExecutionEnvironment();env.enableCheckpointing(60000);StreamTableEnvironmenttableEnv=StreamTableEnvironment.create(env);tableEnv.getConfig().setLocalTimeZone(ZoneId.of("Asia/Shanghai"));// 2. DDL创建Kafka源表tableEnv.executeSql(""" CREATE TABLE orders ( order_id BIGINT, user_id BIGINT, amount DECIMAL(10,2), status STRING, city STRING, order_time TIMESTAMP(3), WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'orders', 'properties.bootstrap.servers' = 'localhost:9092', 'format' = 'json' ) """);// 3. DDL创建MySQL结果表tableEnv.executeSql(""" CREATE TABLE high_value_orders ( order_id BIGINT, user_id BIGINT, discount_amount DECIMAL(10,2), city STRING, PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:mysql://localhost:3306/flink_db', 'table-name' = 'high_value_orders', 'username' = 'root', 'password' = '123456' ) """);// 4. select投影:只选择需要的列,计算折扣金额Tableprojected=tableEnv.from("orders").select($("order_id"),$("user_id"),$("amount"),$("amount").times(lit(0.8)).as("discount_amount"),$("status"),$("city"));// 5. filter过滤:高价值订单(金额>500 且 已支付 且 一线城市)// 注意:高选择性条件放在前面,利用短路求值减少计算Tablefiltered=projected.filter($("status").isEqual(lit("PAID"))// 高选择性,先过滤.and($("city").in(lit("北京"),lit("上海"),lit("广州"),lit("深圳"))).and($("amount").isGreater(lit(500)))// 数值比较,成本低.and($("order_id").isNotNull())// NULL检查,放在最后);// 6. 再次select:只保留输出需要的列(投影裁剪)Tableresult=filtered.select($("order_id"),$("user_id"),$("discount_amount"),$("city"));// 7. 查看执行计划(确认谓词下推和投影裁剪)Stringplan=result.explain();System.out.println("=== 执行计划 ===");System.out.println(plan);// 8. 写入结果表并执行TableResultexecResult=result.executeInsert("high_value_orders");// 9. 等待作业完成execResult.await();}}代码关键点:
- 静态导入 Expressions:
import static org.apache.flink.table.api.Expressions.*,这是 Table API 的第一步。 - select 投影:只选择需要的列,计算折扣金额并用
as()重命名。避免select(*)。 - filter 过滤:高选择性条件(status = ‘PAID’)放在最前面,利用短路求值减少计算。城市 IN 条件次之,数值比较再次之,NULL 检查放在最后。
- 再次 select 投影裁剪:过滤后再次 select,只保留输出需要的列,让优化器进行投影裁剪。
- explain 检查:执行
result.explain()查看执行计划,确认谓词下推和投影裁剪是否生效。 - executeInsert 执行:写入结果表并触发执行,
result.await()等待作业完成。
八、六个常见坑
8.1 坑一:SELECT * 导致投影裁剪失效
现象:查询性能差,数据源读取了大量不需要的列,网络传输和反序列化开销大。
根因:用select($("*"))或SELECT *选择所有列,优化器无法进行投影裁剪,必须读取所有列。
解决方案:永远只选择需要的列,明确列出列名。这是最简单也最有效的性能优化手段,尤其是列式存储(Parquet/ORC)可以只读取需要的列。
8.2 坑二:Filter 中使用 UDF 导致谓词无法下推
现象:查询性能差,数据源读取了全量数据,过滤在 Flink 端执行。
根因:Filter 条件中包含自定义函数(UDF),优化器无法将包含 UDF 的谓词下推到数据源端,必须在 Flink 端读取全量数据后再过滤。
解决方案:
- 尽量用内置函数替代 UDF(内置函数可能支持下推)
- 将 UDF 条件和可下推条件分开,可下推条件放在前面(会被下推),UDF 条件放在后面(在 Flink 端执行)
- 用
explain()确认谓词下推情况
8.3 坑三:NULL 比较用 isEqual 导致结果错误
现象:查询结果不符合预期,NULL 行没有被正确匹配或过滤。
根因:SQL 采用三值逻辑,任何与 NULL 的比较结果都是 UNKNOWN(在 WHERE 中视为 false)。$("a").isEqual(lit(null))永远返回 UNKNOWN,不会匹配任何行。
解决方案:判断 NULL 必须用$("a").isNull()或$("a").isNotNull()。如果需要保留 NULL 行,用$("a").isGreater(10).or($("a").isNull())。
8.4 坑四:Filter 条件顺序不当导致性能差
现象:查询性能差,大量行执行了昂贵的 UDF 或复杂计算后才被过滤掉。
根因:AND 条件中,低选择性、高成本的条件(如 UDF)放在前面,高选择性、低成本的条件放在后面。虽然优化器可能会重排条件,但流处理中没有统计信息时优化器不会重排。
解决方案:手动将高选择性、低成本的条件放在 AND 的前面(如 status = ‘PAID’、city IN (…)),将低选择性、高成本的条件(如 UDF、复杂计算)放在后面。利用短路求值,大部分行在第一个条件就被过滤掉。
8.5 坑五:流处理 distinct 状态无限增长
现象:作业运行一段时间后 Checkpoint 越来越大,TaskManager OOM,作业失败。
根因:流处理 distinct 基于状态去重,维护一个已见行的集合。如果 key 基数大(如用户 ID、订单 ID),状态会持续增长,不会自动清理。
解决方案:
- 设置状态 TTL(
table.exec.state.ttl),过期的状态会被清理 - 考虑用 groupBy 替代 distinct(groupBy 可以配合窗口限制状态)
- 如果是基于事件时间的去重,用 ROW_NUMBER() 去重(可以配合窗口限制状态)
8.6 坑六:流处理 limit/offset 不支持或结果不确定
现象:流处理中使用 limit 报错,或者结果不确定、不符合预期。
根因:limit 是批处理操作,需要全局排序后取前 N 行。流处理是无界流,无法确定"前 N 行"(因为数据持续到来)。直接在无界流上使用 limit 会报错或结果不确定。
解决方案:
- 流处理中不要直接使用 limit,需要配合窗口使用(先窗口聚合,再在窗口内取 TopN)
- 用 ROW_NUMBER() OVER (PARTITION BY … ORDER BY …) 实现 TopN(流处理标准做法)
- 如果是批处理模式,可以正常使用 limit/offset
九、八条最佳实践 Checklist
上线前逐条检查:
永远只选择需要的列:不要用 select(*) 或 SELECT *,明确列出需要的列名,让优化器进行投影裁剪,减少数据读取和传输。
过滤条件尽早执行:在聚合、连接、窗口之前执行 Filter,减少下游数据量。高过滤率的 Filter 是减少数据量最有效的手段。
利用谓词下推:尽量用内置函数和简单比较条件,避免在 Filter 中使用 UDF。用 explain() 确认谓词是否下推到数据源端。
条件顺序手动优化:AND 条件中将高选择性、低成本的条件放在前面,低选择性、高成本的条件(UDF)放在后面,利用短路求值减少计算。
NULL 判断用 isNull/isNotNull:不要用 isEqual(lit(null)) 判断 NULL,SQL 三值逻辑下结果永远是 UNKNOWN。需要保留 NULL 行时用 OR isNull()。
流处理 distinct 设置状态 TTL:流处理 distinct 基于状态去重,必须设置 table.exec.state.ttl 防止状态无限增长导致 OOM。
流处理 TopN 用 ROW_NUMBER:不要在无界流上直接用 limit,用 ROW_NUMBER() OVER (PARTITION BY … ORDER BY …) 实现流处理 TopN。
上线前 explain() 检查执行计划:对复杂查询执行 explain() 检查执行计划,确认谓词下推、投影裁剪、Filter 位置、算子顺序是否符合预期。
十、总结与下一篇预告
Table 查询与过滤原理及代码实现要点回顾:
第一,本质:select(投影)负责列的选择和计算,filter(过滤)负责行的筛选,两者配合从原始表中提取出需要的数据子集,是关系代数中最基础的两个操作(π投影和σ选择)。
第二,四大操作:select 投影(列选择/表达式计算/重命名,投影裁剪优化)、filter 过滤(行筛选/条件判断/短路求值,谓词下推优化)、distinct 去重(批处理排序/哈希,流处理基于状态需设置 TTL)、limit 分页(批处理全局排序取 N 行,流处理需配合窗口或 ROW_NUMBER)。
第三,执行顺序:filter 先于 select 执行是优化器的默认行为(先过滤行再计算列),不需要手动调整。Projection 和 Filter 都是无状态算子,逐行处理,性能高。
第四,六步内部执行流程:方法调用 → 构建 Operation → RelNode 转换(LogicalProject/LogicalFilter)→ 优化器优化 → 代码生成(CodeGen 生成专门 Java 代码直接操作 BinaryRowData)→ 逐行执行。
第五,四大优化原理:谓词下推(Filter 推到数据源端,JDBC/Hive/Parquet 支持)、投影裁剪(只读取需要的列,列式存储效果显著)、常量折叠(常量表达式编译期计算)、条件重排(高选择性低成本条件在前,利用短路求值)。
第六,表达式计算与短路求值:表达式通过 RexNode → CodeGen 编译为直接操作 BinaryRowData 的 Java 代码;AND 条件短路求值(前 false 不执行后条件);NULL 三值逻辑(isEqual(null) 永远 UNKNOWN,用 isNull 判断);过滤率越低下游数据量越小,Filter 是减少数据量最有效算子。
第七,完整代码实现:创建 StreamEnv → 创建 StreamTableEnv → DDL 建表(Kafka 源/MySQL 结果)→ select 投影(只选需要的列 + 计算折扣金额)→ filter 过滤(高选择性条件在前 + 城市 IN + 金额 > 500 + isNotNull)→ 再次 select 投影裁剪 → explain 检查执行计划 → executeInsert 执行 → await 等待完成。
第八,六个常见坑:SELECT * 导致投影裁剪失效、Filter 中 UDF 导致谓词无法下推、NULL 比较用 isEqual 导致结果错误、Filter 条件顺序不当导致性能差、流处理 distinct 状态无限增长、流处理 limit 不支持或结果不确定。
第九,八条最佳实践:永远只选择需要的列、过滤条件尽早执行、利用谓词下推、条件顺序手动优化、NULL 判断用 isNull/isNotNull、流处理 distinct 设置状态 TTL、流处理 TopN 用 ROW_NUMBER、上线前 explain() 检查执行计划。
查询与过滤是 Flink SQL 中最基础也最重要的操作。理解了投影裁剪、谓词下推、短路求值等优化原理,就能写出高性能的查询代码,避免常见的性能陷阱。一个高过滤率的 Filter + 精准的投影裁剪,可以将查询性能提升数倍甚至数十倍。