0%

FlinkSQL-Count

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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
========================================
count field 执行计划详情:
EXPLAIN PLAN FOR SELECT COUNT(name) FROM iceberg_local_table
========================================
== Abstract Syntax Tree ==
LogicalAggregate(group=[{}], EXPR$0=[COUNT($0)])
+- LogicalProject(name=[$1])
+- LogicalTableScan(table=[[default_catalog, default_database, iceberg_local_table]])

== Optimized Physical Plan ==
GroupAggregate(select=[COUNT(name) AS EXPR$0])
+- Exchange(distribution=[single])
+- TableSourceScan(table=[[default_catalog, default_database, iceberg_local_table, project=[name]]], fields=[name])

== Optimized Execution Plan ==
GroupAggregate(select=[COUNT(name) AS EXPR$0])
+- Exchange(distribution=[single])
+- TableSourceScan(table=[[default_catalog, default_database, iceberg_local_table, project=[name]]], fields=[name])



========================================
count 1 执行计划详情:
EXPLAIN PLAN FOR SELECT COUNT(1) FROM iceberg_local_table
========================================
== Abstract Syntax Tree ==
LogicalAggregate(group=[{}], EXPR$0=[COUNT()])
+- LogicalProject($f0=[0])
+- LogicalTableScan(table=[[default_catalog, default_database, iceberg_local_table]])

== Optimized Physical Plan ==
GroupAggregate(select=[COUNT(*) AS EXPR$0])
+- Exchange(distribution=[single])
+- TableSourceScan(table=[[default_catalog, default_database, iceberg_local_table, project=[]]], fields=[])

== Optimized Execution Plan ==
GroupAggregate(select=[COUNT(*) AS EXPR$0])
+- Exchange(distribution=[single])
+- TableSourceScan(table=[[default_catalog, default_database, iceberg_local_table, project=[]]], fields=[])



========================================
count star 执行计划详情:
EXPLAIN PLAN FOR SELECT COUNT(*) FROM iceberg_local_table
========================================
== Abstract Syntax Tree ==
LogicalAggregate(group=[{}], EXPR$0=[COUNT()])
+- LogicalProject($f0=[0])
+- LogicalTableScan(table=[[default_catalog, default_database, iceberg_local_table]])

== Optimized Physical Plan ==
GroupAggregate(select=[COUNT(*) AS EXPR$0])
+- Exchange(distribution=[single])
+- TableSourceScan(table=[[default_catalog, default_database, iceberg_local_table, project=[]]], fields=[])

== Optimized Execution Plan ==
GroupAggregate(select=[COUNT(*) AS EXPR$0])
+- Exchange(distribution=[single])
+- TableSourceScan(table=[[default_catalog, default_database, iceberg_local_table, project=[]]], fields=[])

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
2
3
4
5
6
@Override
public Expression[] accumulateExpressions() {
return new Expression[] {
/* count = */ ifThenElse(isNull(operand(0)), count, plus(count, literal(1L)))
};
}

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
2
3
4
@Override
public Expression[] accumulateExpressions() {
return new Expression[] {/* count1 = */ plus(count1, literal(1L))};
}

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 关键代码位置总结

  1. 聚合函数实现

    • 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
  2. count(*)判断逻辑

    • FlinkRelBuilder.isCountStarAgg(): flink-table-planner/src/main/java/org/apache/flink/table/planner/calcite/FlinkRelBuilder.java
  3. 聚合下推接口

    • SupportsAggregatePushDown: flink-table-common/src/main/java/org/apache/flink/table/connector/source/abilities/SupportsAggregatePushDown.java
  4. 聚合分割规则

    • SplitAggregateRule: flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/rules/logical/SplitAggregateRule.scala

1.6 元数据信息

1
2
3
4
5
6
7
8
9
10
11
12
ll  metadata
total 60K
-rw-r--r-- 1 root root 11K Dec 22 21:42 3031cf5a-37b0-4f97-8568-5af892f48244-m0.avro
-rw-r--r-- 1 root root 11K Dec 22 21:40 56d0c3ef-8f73-44c8-8d20-553bbdbbd200-m0.avro
-rw-r--r-- 1 root root 4.5K Dec 22 21:42 snap-7569084149956477108-1-3031cf5a-37b0-4f97-8568-5af892f48244.avro
-rw-r--r-- 1 root root 4.4K Dec 22 21:40 snap-7985011370354105799-1-56d0c3ef-8f73-44c8-8d20-553bbdbbd200.avro
-rw-r--r-- 1 root root 1.8K Dec 22 21:27 v1.metadata.json
-rw-r--r-- 1 root root 2.7K Dec 22 21:40 v2.metadata.json
-rw-r--r-- 1 root root 4.1K Dec 22 21:42 v3.metadata.json
-rw-r--r-- 1 root root 1 Dec 22 21:42 version-hint.text
#
cat v3.metadata.json
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
// metadata.json
{
"format-version": 2,
"table-uuid": "2bf9cc09-1099-454a-aae7-c13fe14f7f20",
"location": "file:///tmp/iceberg_local_warehouse/default_database/iceberg_local_table",
"last-sequence-number": 2,
"last-updated-ms": 1766410956369,
"last-column-id": 5,
"current-schema-id": 0,
"schemas": [
{
"type": "struct",
"schema-id": 0,
"fields": [
{
"id": 1,
"name": "id",
"required": false,
"type": "long"
},
{
"id": 2,
"name": "name",
"required": false,
"type": "string"
},
{
"id": 3,
"name": "category",
"required": false,
"type": "string"
},
{
"id": 4,
"name": "create_time",
"required": false,
"type": "timestamp"
},
{
"id": 5,
"name": "creator",
"required": false,
"type": "string"
}
]
}
],
"default-spec-id": 0,
"partition-specs": [
{
"spec-id": 0,
"fields": []
}
],
"last-partition-id": 999,
"default-sort-order-id": 0,
"sort-orders": [
{
"order-id": 0,
"fields": []
}
],
"current-index-id": 0,
"indexes": [
{
"index-id": 0,
"index-name": "empty",
"index-type": "",
"index-properties": {},
"fields": []
}
],
"properties": {
"connector": "iceberg",
"write.parquet.compression-codec": "zstd",
"catalog-type": "hadoop",
"keep-modified-paths": "file:///tmp/iceberg_local_warehouse/default_database/iceberg_local_table",
"catalog-name": "demo_catalog",
"warehouse": "file:///tmp/iceberg_local_warehouse"
},
"current-snapshot-id": 7569084149956477108,
"refs": {
"main": {
"snapshot-id": 7569084149956477108,
"type": "branch"
}
},
"snapshots": [
{
"sequence-number": 1,
"snapshot-id": 7985011370354105799,
"timestamp-ms": 1766410855548,
"summary": {
"operation": "append",
"flink.checkpoint-timestamp": "1766410855378",
"flink.operator-id": "a0e952a7da09b602686e75133408a70e",
"engine": "flink-1.15",
"oceanus.job-id": "8da0afd77543eb566e79c5dc0656bf1c",
"flink.job-id": "8da0afd77543eb566e79c5dc0656bf1c",
"jobType": "streaming",
"flink.max-committed-checkpoint-id": "9223372036854775807",
"added-data-files": "48",
"added-small-data-files": "48",
"added-records": "100",
"added-files-size": "82584",
"changed-partition-count": "1",
"user": "root",
"jobId": "8da0afd77543eb566e79c5dc0656bf1c",
"jobName": "insert-into_default_catalog.default_database.iceberg_local_table",
"total-records": "100",
"total-files-size": "82584",
"total-data-files": "48",
"total-delete-files": "0",
"total-position-deletes": "0",
"total-equality-deletes": "0",
"engine-version": "1.15-tq-0.1.42",
"engine-name": "flink",
"iceberg-version": "Apache Iceberg 1.8.2-6-tq (commit 1f7547520767fe0d4537361e6c812962fa750d83)"
},
"manifest-list": "file:/tmp/iceberg_local_warehouse/default_database/iceberg_local_table/metadata/snap-7985011370354105799-1-56d0c3ef-8f73-44c8-8d20-553bbdbbd200.avro",
"schema-id": 0
},
{
"sequence-number": 2,
"snapshot-id": 7569084149956477108,
"parent-snapshot-id": 7985011370354105799,
"timestamp-ms": 1766410956369,
"summary": {
"operation": "append",
"flink.checkpoint-timestamp": "1766410956164",
"flink.operator-id": "a0e952a7da09b602686e75133408a70e",
"engine": "flink-1.15",
"oceanus.job-id": "b95b354367e1fde2dcc30438f62e8e75",
"flink.job-id": "b95b354367e1fde2dcc30438f62e8e75",
"jobType": "streaming",
"flink.max-committed-checkpoint-id": "9223372036854775807",
"added-data-files": "48",
"added-small-data-files": "48",
"added-records": "100",
"added-files-size": "82585",
"changed-partition-count": "1",
"user": "root",
"jobId": "b95b354367e1fde2dcc30438f62e8e75",
"jobName": "insert-into_default_catalog.default_database.iceberg_local_table",
"total-records": "200",
"total-files-size": "165169",
"total-data-files": "96",
"total-delete-files": "0",
"total-position-deletes": "0",
"total-equality-deletes": "0",
"engine-version": "1.15-tq-0.1.42",
"engine-name": "flink",
"iceberg-version": "Apache Iceberg 1.8.2-6-tq (commit 1f7547520767fe0d4537361e6c812962fa750d83)"
},
"manifest-list": "file:/tmp/iceberg_local_warehouse/default_database/iceberg_local_table/metadata/snap-7569084149956477108-1-3031cf5a-37b0-4f97-8568-5af892f48244.avro",
"schema-id": 0
}
],
"statistics": [],
"partition-statistics": [],
"snapshot-log": [
{
"timestamp-ms": 1766410855548,
"snapshot-id": 7985011370354105799
},
{
"timestamp-ms": 1766410956369,
"snapshot-id": 7569084149956477108
}
],
"metadata-log": [
{
"timestamp-ms": 1766410022055,
"metadata-file": "file:/tmp/iceberg_local_warehouse/default_database/iceberg_local_table/metadata/v1.metadata.json"
},
{
"timestamp-ms": 1766410855548,
"metadata-file": "file:/tmp/iceberg_local_warehouse/default_database/iceberg_local_table/metadata/v2.metadata.json"
}
]
}
1
2
3
#

java -jar avro-tools.java tojson snap-7569084149956477108-1-3031cf5a-37b0-4f97-8568-5af892f48244.avro > local_iceberg_schema.avro.json

每个数据文件,都会对应一行 avro 记录:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
// avro record of data file
{
"status": 1,
"snapshot_id": {
"long": 7569084149956477108
},
"sequence_number": null,
"file_sequence_number": null,
"min_sequence_number": null,
"data_file": {
"content": 0,
"file_path": "file:/tmp/iceberg_local_warehouse/default_database/iceberg_local_table/data/00010-0-30b0e113-e15a-484a-a4db-9107ecbc3034-00001.parquet",
"file_format": "PARQUET",
"partition": {},
"record_count": 2,
"file_size_in_bytes": 1705,
"column_sizes": {
"array": [
{
"key": 1,
"value": 57
},
{
"key": 2,
"value": 62
},
{
"key": 3,
"value": 83
},
{
"key": 4,
"value": 57
},
{
"key": 5,
"value": 84
}
]
},
"value_counts": {
"array": [
{
"key": 1,
"value": 2
},
{
"key": 2,
"value": 2
},
{
"key": 3,
"value": 2
},
{
"key": 4,
"value": 2
},
{
"key": 5,
"value": 2
}
]
},
"null_value_counts": {
"array": [
{
"key": 1,
"value": 0
},
{
"key": 2,
"value": 0
},
{
"key": 3,
"value": 0
},
{
"key": 4,
"value": 0
},
{
"key": 5,
"value": 0
}
]
},
"nan_value_counts": {
"array": []
},
"lower_bounds": {
"array": [
{
"key": 1,
"value": "\u000B\u0000\u0000\u0000\u0000\u0000\u0000\u0000"
},
{
"key": 2,
"value": "产品_11"
},
{
"key": 3,
"value": "图书"
},
{
"key": 4,
"value": "èˆþM‘F\u0006\u0000"
},
{
"key": 5,
"value": "赵六"
}
]
},
"upper_bounds": {
"array": [
{
"key": 1,
"value": ";\u0000\u0000\u0000\u0000\u0000\u0000\u0000"
},
{
"key": 2,
"value": "产品_59"
},
{
"key": 3,
"value": "图书"
},
{
"key": 4,
"value": "Я\rN‘F\u0006\u0000"
},
{
"key": 5,
"value": "赵六"
}
]
},
"key_metadata": null,
"split_offsets": {
"array": [
4
]
},
"equality_ids": null,
"sort_order_id": {
"int": 0
},
"referenced_data_file": null,
"schema_id": {
"int": 0
},
"index_metadata": {
"map": {}
},
"extra_metadata": {
"map": {
"write.parquet.compression-codec": "ZSTD"
}
}
}
}

Iceberg文件级别元数据

୦ Iceberg的datafile元数据中包含record count、字段大小等信息,Flink写入时已记录这些信息。

୦ 未来可探讨利用manifest和文件级别元数据优化聚合下推功能。

2 Flink SQL执行与iceberg关系

名称 动作 输入 输出 关键类 Iceberg调用
Parsing 词法解析与语法解析,构建SqlNode树 SQL 字符串 SqlNode 树 SqlParser
org.apache.flink.table.planner.parse.CalciteParser#parseSqlList
Validation 语义验证
结合Catalog元数据(表、列、函数定义)
SqlNode + Catalog 验证后的 SqlNode
内部绑定:
1. 表的标识
2. 列的类型信息
3. 函数调用合法性确认
FlinkSqlValidator
org.apache.flink.table.planner.calcite.FlinkPlannerImpl#validate
Conversion 逻辑计划生成: 关系代数转换
1. 支持时间属性字段
2. 处理Watermark
3. 区分Append/Update/Retract流模式
SqlNode 逻辑 RelNode 树 FlinkSqlToRelConverter
sqlToOperationConverter#convertValidatedSqlNode
FlinkPlannerImpl#rel
创建org.apache.iceberg.flink.source.IcebergTableSource
Optimization HepPlanner基于规则做逻辑优化
VolcanoPlanner代价模型做物理优化
逻辑 RelNode 优化后 RelNode(Physical RelNode) org.apache.flink.table.planner.delegation.PlannerBase#optimize
org.apache.flink.table.planner.plan.optimize.program.FlinkChainedProgram
RelOptPlanner
RelOptRule,
FlinkLogicalOptRuleSet
ConventionTrait
IcebergTableSource#applyProjection
PushProjectIntoTableSourceScanRule
CodeGen 转换为DataStream API RelNode JobGraph RelNode#translateToPlan
ExecNode.translateToPlan()
translateToExecNodeGraph
创建InputFormatSourceFunction
IcebergTableSource#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/logic
org/apache/flink/table/planner/plan/rules/physical

优化规则:本质是关系表达式的等价变换

Rules 的本质是一种模式匹配与替换机制

  • 输入: 一棵逻辑算子树(由 RelNode 组成,如 LogicalProject, LogicalFilter)。
  • 匹配(Matching): 规则定义了一个 RelOptRuleOperand(操作数),用来描述它感兴趣的算子结构。例如,“一个 Filter 紧跟在另一个 Filter 之后”。
  • 变换(Transformation): 当匹配成功时,调用 onMatch 方法,将旧的算子树替换为等价但更优的新算子树。

IcebergTableSource 支持了 SupportsProjectionPushDown, SupportsFilterPushDown, SupportsLimitPushDown 三种下推

  1. 投影下推 (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

  1. 初始化阶段
    创建分区常量映射:从FileScanTask中提取分区信息,创建常量映射
    创建删除过滤器:初始化FlinkDeleteFilter处理数据删除逻辑
  2. 数据读取阶段
    根据文件格式选择对应的Reader:
    Parquet Reader:使用FlinkParquetReaders.buildReader
    Avro Reader:使用FlinkAvroReader
    ORC Reader:使用FlinkOrcReader
  3. 过滤处理阶段
    基础数据过滤:应用task.residual()表达式过滤或rowFilter过滤
    删除过滤:处理行删除和位置删除逻辑
  4. 投影转换阶段
    投影转换:当requiredSchema与projectedSchema不一致时,应用RowDataProjection进行字段投影

2.5 代办

  • org.apache.iceberg.flink.source.RowDataFileScanTaskReader#newParquetIterable 中Parquet读取数据的split操作?
  • FunctionGen‎erator.scala 用于Function的CodeGen