1 Flink SQL 查询iceberg count语义
1.1 结论:
| 查询 | 方式 | 是否读取数据 | 能否识别数据列 |
|---|---|---|---|
| count(字段) | metadata | Y | Y |
| count(distinct 字段) | datafile | Y | Y |
| count(*)/count(1) | metadata | Y | N |
Iceberg通过元数据优化:
- 如果表启用了 count-based metadata optimization(例如通过
TableScan+useSnapshotId+ manifest-level row count),Flink 可能直接从 manifest 获取总行数。 - 但 Flink 当前(截至 Flink 1.18 / Iceberg 1.4)对这类优化支持有限,多数情况下仍会触发 data file 扫描,除非手动启用特定优化或使用 Iceberg 的
metadata table查询。
执行计划:
1 | ======================================== |
select count(field)的projection,可以拿到对应的字段
● Flink与Iceberg的聚合下推问题
୦ 目前Flink与Iceberg在聚合下推方面配合不佳,Flink的Source Function仅支持filter和project下推,不支持聚合下推。
୦ Spark支持聚合下推优化,但Flink在1.15版本中仅支持project和filter下推,高版本(1.18+)实现了manifest和snapshot级别的count优化。
୦ 当前Flink在project为空时会扫描全表数据,但未实际读取数据,导致性能问题。
1.2 AggFunction
1.2.1 count(字段)
- 实现类:
CountAggFunction - 行为:只统计非null值
- 代码位置:
/Users/averyzhang/workspace/flink/flink-1.15/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/functions/aggfunctions/CountAggFunction.java
- 关键逻辑:
1 |
|
1.2.2 count(1)和count(*)
- 实现类:
Count1AggFunction - 行为:统计所有记录,包括null值
- 代码位置:
/Users/averyzhang/workspace/flink/flink-1.15/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/functions/aggfunctions/Count1AggFunction.java - 关键逻辑:
1 |
|
1.3 算子转换
1.3.1 count(字段)
- 算子:
CountAggFunction - 优化规则:在
SplitAggregateRule中,partial阶段使用COUNT,final阶段使用SUM0
1.3.2 count(1)和count(*)
- 算子:
Count1AggFunction - 判断逻辑:在
FlinkRelBuilder.isCountStarAgg()方法中判断是否为count(*)
1.4 iceberg表的具体行为
对于iceberg表,如果实现了SupportsAggregatePushDown接口
(代码位置:/Users/averyzhang/workspace/flink/flink-1.15/flink-table/flink-table-common/src/main/java/org/apache/flink/table/connector/source/abilities/SupportsAggregatePushDown.java):
1.4.1 count(字段)
- 如果字段有null值,需要扫描数据来判断
- 如果字段没有null值约束,可能可以利用iceberg的列统计信息进行优化
1.4.2 count(1)和count(*)
- 最佳优化场景:可以直接利用iceberg表的文件级别统计信息(如row_count)
- 性能优势:可能不需要扫描实际数据,直接从元数据获取计数
1.5 关键代码位置总结
聚合函数实现:
- CountAggFunction:
flink-table-planner/src/main/java/org/apache/flink/table/planner/functions/aggfunctions/CountAggFunction.java - Count1AggFunction:
flink-table-planner/src/main/java/org/apache/flink/table/planner/functions/aggfunctions/Count1AggFunction.java
- CountAggFunction:
count(*)判断逻辑:FlinkRelBuilder.isCountStarAgg():flink-table-planner/src/main/java/org/apache/flink/table/planner/calcite/FlinkRelBuilder.java
聚合下推接口:
SupportsAggregatePushDown:flink-table-common/src/main/java/org/apache/flink/table/connector/source/abilities/SupportsAggregatePushDown.java
聚合分割规则:
SplitAggregateRule:flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/rules/logical/SplitAggregateRule.scala
1.6 元数据信息
1 | ll metadata |
1 | // metadata.json |
1 | # |
每个数据文件,都会对应一行 avro 记录:
1 | // avro record of data file |
Iceberg文件级别元数据
୦ Iceberg的datafile元数据中包含record count、字段大小等信息,Flink写入时已记录这些信息。
୦ 未来可探讨利用manifest和文件级别元数据优化聚合下推功能。
2 Flink SQL执行与iceberg关系
| 名称 | 动作 | 输入 | 输出 | 关键类 | Iceberg调用 |
|---|---|---|---|---|---|
| Parsing | 词法解析与语法解析,构建SqlNode树 | SQL 字符串 | SqlNode 树 |
SqlParserorg.apache.flink.table.planner.parse.CalciteParser#parseSqlList |
|
| Validation | 语义验证 结合Catalog元数据(表、列、函数定义) |
SqlNode + Catalog |
验证后的 SqlNode内部绑定: 1. 表的标识 2. 列的类型信息 3. 函数调用合法性确认 |
FlinkSqlValidatororg.apache.flink.table.planner.calcite.FlinkPlannerImpl#validate |
|
| Conversion | 逻辑计划生成: 关系代数转换 1. 支持时间属性字段 2. 处理Watermark 3. 区分Append/Update/Retract流模式 |
SqlNode |
逻辑 RelNode 树 |
FlinkSqlToRelConvertersqlToOperationConverter#convertValidatedSqlNodeFlinkPlannerImpl#rel |
创建org.apache.iceberg.flink.source.IcebergTableSource |
| Optimization | HepPlanner基于规则做逻辑优化 VolcanoPlanner代价模型做物理优化 |
逻辑 RelNode |
优化后 RelNode(Physical RelNode) |
org.apache.flink.table.planner.delegation.PlannerBase#optimizeorg.apache.flink.table.planner.plan.optimize.program.FlinkChainedProgramRelOptPlanner, RelOptRule,FlinkLogicalOptRuleSetConvention,Trait |
IcebergTableSource#applyProjectionPushProjectIntoTableSourceScanRule |
| CodeGen | 转换为DataStream API | RelNode |
JobGraph |
RelNode#translateToPlanExecNode.translateToPlan()translateToExecNodeGraph |
创建InputFormatSourceFunctionIcebergTableSource#createDataStream |
2.1 优化
优化的入口: org.apache.flink.table.planner.plan.optimize.StreamCommonSubGraphBasedOptimizer#doOptimize
2.2 FlinkOptimizeProgram
用于优化Stream Table Plan的Programs序列: org.apache.flink.table.planner.plan.optimize.program.FlinkStreamProgram
| FlinkChainedProgram | 作用 |
|---|---|
| subquery_rewrite | |
| temporal_join_rewrite | |
| decorrelate | |
| default_rewrite | |
| predicate_pushdown | |
| join_reorder | |
| project_rewrite | |
| logical | |
| logical_rewrite | |
| time_indicator | |
| physical | |
| physical_rewrite |
2.3 优化规则
org.apache.flink.table.planner.plan.rules.FlinkStreamRuleSets
org/apache/flink/table/planner/plan/rules/logicorg/apache/flink/table/planner/plan/rules/physical
优化规则:本质是关系表达式的等价变换
Rules 的本质是一种模式匹配与替换机制。
- 输入: 一棵逻辑算子树(由
RelNode组成,如LogicalProject,LogicalFilter)。 - 匹配(Matching): 规则定义了一个
RelOptRuleOperand(操作数),用来描述它感兴趣的算子结构。例如,“一个 Filter 紧跟在另一个 Filter 之后”。 - 变换(Transformation): 当匹配成功时,调用
onMatch方法,将旧的算子树替换为等价但更优的新算子树。
IcebergTableSource 支持了 SupportsProjectionPushDown, SupportsFilterPushDown, SupportsLimitPushDown 三种下推
- 投影下推 (Projection PushDown)
方法: applyProjection(int[][] projectFields)
功能: 只读取需要的列,减少数据传输
实现: 通过projectedFields数组记录需要投影的列索引
2. 过滤条件下推 (Filter PushDown)
方法: applyFilters(List<ResolvedExpression> flinkFilters)
功能: 将过滤条件推送到数据源层执行
实现: 将Flink表达式转换为Iceberg表达式,返回接受的过滤器列表
3. 限制下推 (Limit PushDown)
方法: applyLimit(long newLimit)
功能: 限制返回的数据条数
实现: 设置limit字段值,在数据源层进行限制
这三种pushdown功能最终在getScanRuntimeProvider方法中整合,通过FlinkSource.forRowData()构建数据流时应用相应的投影、过滤和限制条件,从而在数据源层面实现优化,减少网络传输和计算开销。
2.3.1 投影下推(Projection-PushDown)规则
org.apache.flink.table.planner.plan.rules.logical.PushProjectIntoTableSourceScanRule
- 减少I/O开销: 只读取查询需要的列数据
- 降低网络传输: 减少不必要的数据传输
- 提高处理效率: 在数据源层面进行过滤
- 支持复杂类型: 处理嵌套和变体数据结构的投影
graph LR
A[获取Project和Scan] --> B[检查嵌套投影支持]
B --> C[提取引用字段]
C --> D[处理变体字段]
D --> E[构建投影模式]
E --> F[执行投影下推]
F --> G[创建新的表源和扫描]
G --> H[重写投影表达式]
H --> I[返回优化结果]
flowchart TB subgraph ProjectionPushDown direction LR subgraph matches direction LR a1 --> a2 end subgraph onMatch direction LR a4 --> b3 end end
2.3.2 过滤下推
2.3.3 Limit下推
2.4 Iceberg读取数据
核心类: org.apache.iceberg.flink.source.RowDataFileScanTaskReader
- 初始化阶段
创建分区常量映射:从FileScanTask中提取分区信息,创建常量映射
创建删除过滤器:初始化FlinkDeleteFilter处理数据删除逻辑 - 数据读取阶段
根据文件格式选择对应的Reader:
Parquet Reader:使用FlinkParquetReaders.buildReader
Avro Reader:使用FlinkAvroReader
ORC Reader:使用FlinkOrcReader - 过滤处理阶段
基础数据过滤:应用task.residual()表达式过滤或rowFilter过滤
删除过滤:处理行删除和位置删除逻辑 - 投影转换阶段
投影转换:当requiredSchema与projectedSchema不一致时,应用RowDataProjection进行字段投影
2.5 代办
- org.apache.iceberg.flink.source.RowDataFileScanTaskReader#newParquetIterable 中Parquet读取数据的split操作?
- FunctionGenerator.scala 用于Function的CodeGen