☰
Flink SQL 之 Table 查询与过滤原理及代码实现:从投影裁剪到谓词下推与短路求值
2026/9/27 11:28:20 网站建设 项目流程

上一篇讲了 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 链式调用从输入到执行,完整的流程分为六步:

  1. 方法调用:用户调用 select/filter/distinct 等方法
  2. 构建 Operation:每个方法构建一个 QueryOperation 节点
  3. RelNode 转换:Operation 翻译为 Calcite LogicalProject/LogicalFilter
  4. 优化器优化:谓词下推/投影裁剪/常量折叠/条件重排
  5. 代码生成:CodeGen 生成 Projection/Filter 算子的 Java 代码
  6. 逐行执行:算子逐行处理: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 可以大幅减少下游聚合、连接、窗口等算子的计算量和状态大小。

最佳实践:

  1. 尽早执行 Filter(在聚合/连接之前)
  2. 利用谓词下推将 Filter 推到数据源端
  3. 避免在 Filter 中使用复杂 UDF(UDF 无法谓词下推,且计算成本高)
  4. 用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();}}

代码关键点:

  1. 静态导入 Expressions:import static org.apache.flink.table.api.Expressions.*,这是 Table API 的第一步。
  2. select 投影:只选择需要的列,计算折扣金额并用as()重命名。避免select(*)。
  3. filter 过滤:高选择性条件(status = ‘PAID’)放在最前面,利用短路求值减少计算。城市 IN 条件次之,数值比较再次之,NULL 检查放在最后。
  4. 再次 select 投影裁剪:过滤后再次 select,只保留输出需要的列,让优化器进行投影裁剪。
  5. explain 检查:执行result.explain()查看执行计划,确认谓词下推和投影裁剪是否生效。
  6. executeInsert 执行:写入结果表并触发执行,result.await()等待作业完成。

八、六个常见坑

8.1 坑一:SELECT * 导致投影裁剪失效

现象:查询性能差,数据源读取了大量不需要的列,网络传输和反序列化开销大。

根因:用select($("*"))或SELECT *选择所有列,优化器无法进行投影裁剪,必须读取所有列。

解决方案:永远只选择需要的列,明确列出列名。这是最简单也最有效的性能优化手段,尤其是列式存储(Parquet/ORC)可以只读取需要的列。

8.2 坑二:Filter 中使用 UDF 导致谓词无法下推

现象:查询性能差,数据源读取了全量数据,过滤在 Flink 端执行。

根因:Filter 条件中包含自定义函数(UDF),优化器无法将包含 UDF 的谓词下推到数据源端,必须在 Flink 端读取全量数据后再过滤。

解决方案:

  1. 尽量用内置函数替代 UDF(内置函数可能支持下推)
  2. 将 UDF 条件和可下推条件分开,可下推条件放在前面(会被下推),UDF 条件放在后面(在 Flink 端执行)
  3. 用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),状态会持续增长,不会自动清理。

解决方案:

  1. 设置状态 TTL(table.exec.state.ttl),过期的状态会被清理
  2. 考虑用 groupBy 替代 distinct(groupBy 可以配合窗口限制状态)
  3. 如果是基于事件时间的去重,用 ROW_NUMBER() 去重(可以配合窗口限制状态)

8.6 坑六:流处理 limit/offset 不支持或结果不确定

现象:流处理中使用 limit 报错,或者结果不确定、不符合预期。

根因:limit 是批处理操作,需要全局排序后取前 N 行。流处理是无界流,无法确定"前 N 行"(因为数据持续到来)。直接在无界流上使用 limit 会报错或结果不确定。

解决方案:

  1. 流处理中不要直接使用 limit,需要配合窗口使用(先窗口聚合,再在窗口内取 TopN)
  2. 用 ROW_NUMBER() OVER (PARTITION BY … ORDER BY …) 实现 TopN(流处理标准做法)
  3. 如果是批处理模式,可以正常使用 limit/offset

九、八条最佳实践 Checklist

上线前逐条检查:

  1. 永远只选择需要的列:不要用 select(*) 或 SELECT *,明确列出需要的列名,让优化器进行投影裁剪,减少数据读取和传输。

  2. 过滤条件尽早执行:在聚合、连接、窗口之前执行 Filter,减少下游数据量。高过滤率的 Filter 是减少数据量最有效的手段。

  3. 利用谓词下推:尽量用内置函数和简单比较条件,避免在 Filter 中使用 UDF。用 explain() 确认谓词是否下推到数据源端。

  4. 条件顺序手动优化:AND 条件中将高选择性、低成本的条件放在前面,低选择性、高成本的条件(UDF)放在后面,利用短路求值减少计算。

  5. NULL 判断用 isNull/isNotNull:不要用 isEqual(lit(null)) 判断 NULL,SQL 三值逻辑下结果永远是 UNKNOWN。需要保留 NULL 行时用 OR isNull()。

  6. 流处理 distinct 设置状态 TTL:流处理 distinct 基于状态去重,必须设置 table.exec.state.ttl 防止状态无限增长导致 OOM。

  7. 流处理 TopN 用 ROW_NUMBER:不要在无界流上直接用 limit,用 ROW_NUMBER() OVER (PARTITION BY … ORDER BY …) 实现流处理 TopN。

  8. 上线前 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 + 精准的投影裁剪,可以将查询性能提升数倍甚至数十倍。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询