0%

Flink SQL基于ANSI SQL标准语法,使用ROW_NUMBER的开窗聚合函数来解决分组TopN的问题, 语法:

1
2
3
4
5
6
7
8
9
10
11
SELECT *
FROM
(
SELECT *,
ROW_NUMBER() OVER (
[PARTITION BY col1[, col2..]]
ORDER BY col1 [asc|desc][, col2 [asc|desc]...]
) AS rownum
FROM table_name
) ta
WHERE rownum <= N [AND conditions]

参数说明:
ROW_NUMBER(): 是一个计算行号的OVER窗口函数,行号计算从1开始。
PARTITION BY col1[, col2..] : 指定分区的列,可以不指定。
ORDER BY col1 [asc|desc][, col2 [asc|desc]...]: 指定排序的列,可以多列不同排序方向

TopN 需要两层 query,子查询中使用ROW_NUMBER()开窗函数来为每条数据标上排名,排名的计算根据PARTITION BYORDER BY来指定分区列和排序列,也就是说每一条数据会计算其在所属分区中,根据排序列排序得到的排名。在外层查询中,对排名进行过滤,只取出排名小于 N 的,如 N=10,那么就是取 Top 10 的数据。如果没有指定PARTITION BY那么就是一个全局 TopN 的计算,所以 ROW_NUMBER 在使用上更为灵活。

SQL范例:

对全网商家根据行业按销售额排序,计算出每个行业销售额前十名的商家:

1
2
3
4
5
6
7
8
9
SELECT *  
FROM (
SELECT *,
ROW_NUMBER() OVER (
PARTITION BY category
ORDER BY sales DESC
) AS rownum
FROM shop_sales)
WHERE rownum <= 10

ROW_NUMBER 方式的 TopN 语法非常灵活,能满足全局 TopN 和分组 TopN 的需求

LogicalWindow 会对所有数据进行排名,也就是说每当到达一个数据,就要对历史数据进行重排序,并输出历史数据的新的排名,然后 LogicalCalc 节点会根据排名进行过滤。这在性能上是非常糟糕的,因为这无限放大了流量。而我们知道,最优的流式 TopN 的计算只需要维护一个 N 元素大小的小根堆,每当有数据到达时,只需要与堆顶元素比较,如果比堆顶元素还小,则直接丢弃;如果比堆顶元素大,则更新小根堆,并输出更新后的排行榜。也就是说我们不需要分为两个节点进行计算,不需要将所有数据进行排序,只需要在一个节点中就可以高效地完成计算。所以我们在查询优化器中加入了一条规则,在使用 TopN 语法时,将 LogicalWindow 和 LogicalCalc 合并成了 LogicalRank 节点。LogicalRank 在翻译成物理执行计划时,会使用一个经过特殊设计的 TopN 算子

TopN 算子的实现上主要有两个数据结构,一个是 TreeMap,另一个是 MapState。TreeMap 的作用类似于上文的小根堆,有序地存放了排名前 N 的元素。但是 TreeMap 是个内存数据结构,在 failover 后会丢失,无法保证数据的一致性。因此我们还有一个 MapState 结构,MapState 是 Flink 提供的状态接口,用来存储 TopN 的数据(保证数据不丢)。当有 failover 发生后,MapState 能保证状态的恢复,而 TreeMap 会从 MapState 中重新构造出来。我们并有没有把顺序也存到状态中去,因为顺序是可以在恢复时重构的。因为每一次状态的读写操作都会涉及到序列化/反序列化,往往是性能的瓶颈,所以 TreeMap 的主要作用是降低了对 MapState 状态的读写操作。对大部分数据来说都是与 TreeMap 进行交互,不需要对 MapState 进行读写的,全是内存操作,所以 TopN 的性能是非常高的。

TopN 算子的主要处理流程是,每当有数据到达时,会与 TreeMap 的最小的元素比较,如果比它小,那么该数据就不可能是 TopN 的一员,直接丢弃即可。如果比它大,那么就会先更新 TreeMap,同时更新 MapState 中的存的数据。最后输出更新后的排行榜。为了减少冗余数据的输出,我们只会输出排名发生变化的数据。例如原先的第7名上升到了第六名,那么只需要输出新的第六名和第七名即可

嵌套 TopN 解决热点问题

TopN 的计算与 GroupBy 的计算类似,如果数据存在倾斜,则会有计算热点的现象。比如全局 TopN,那么所有的数据只能汇集到一个节点进行 TopN 的计算,那么计算能力就会受限于单台机器,无法做到水平扩展。解决思路与 GroupBy 是类似的,就是使用嵌套 TopN,或者说两层 TopN。在原先的 TopN 前面,再加一层 TopN,用于分散热点。例如,计算全网排名前十的商铺,会导致单点的数据热点,那么可以先加一层分组 TopN,组的划分规则是根据店铺 ID 哈希取模后分成128组(并发的倍数)。第二层 TopN 与原先的写法一样,没有 PARTITION BY。第一层会计算出每一组的 TopN,而后在第二层中进行合并汇总,得到最终的全网前十。第二层虽然仍是单点,但是大量的计算量由第一层分担了,而第一层是可以水平扩展的。使用嵌套 TopN 的优化写法如下所示:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
CREATE VIEW tmp_topn AS  
SELECT *
FROM (
SELECT *,
ROW_NUMBER() OVER (PARTITION BY HASH_CODE(shop_id)%128 ORDER BY sales DESC) AS rownum
FROM shop_sales)
WHERE rownum <= 10

SELECT *
FROM (
SELECT shop_id, shop_name, sales,
ROW_NUMBER() OVER (ORDER BY sales DESC) AS rownum
FROM tmp_topn)
WHERE rownum <= 10

Flink中如何实现TopN?

windowall是一个全局并发为1的操作,所有的数据只能汇集到一个节点进行 TopN 的计算,那么计算能力就会受限于单台机器,容易产生数据热点问题


[参考文献]

  1. Flink SQL高效Top-N方案的实现原理
  2. Flink SQL 功能解密系列-流式 TopN 挑战与实现–阿里天池
  3. 基于Flink快速开发实时TopN

Flink 去重

DataStream去重-timer

DataStream-two-phase-aggregation

SQL-count-distinct

SQL-group-by-latest

基数去重

HypeLogLog基数: 一种基于概率的基数估计算法,用于高效统计大规模数据集中不重复元素的数量(基数),其核心思想是通过哈希函数和分桶策略,以极低的内存消耗(通常仅需几KB)实现近似计算。

  1. 菜鸟供应链实时技术架构演进

  2. 趣头条实战-基于Flink+ClickHouse构建实时数据平台

  3. ApacheFlink新场景-OLAP引擎

  4. 说说Flink DataStream的八种物理分区逻辑

  5. State Processor API:如何读取,写入和修改 Flink 应用程序的状态

  6. Flink滑动窗口原理与细粒度滑动窗口的性能问题

  7. 基于Flink快速开发实时TopN

  8. 使用 Apache Flink 开发实时 ETL

  9. Flink Source/Sink探究与实践:RocketMQ数据写入HBase

  10. Spark/Flink广播实现作业配置动态更新

  11. Flink全链路延迟的测量方式

  12. Flink原理-Flink中的数据抽象及数据交换过程

  13. Flink SQL Window源码全解析

  14. Flink DataStream维度表Join的简单方案

  15. Apache Flink的内存管理

  16. Flink1.9整合Kafka实战

  17. Apache Flink在小米的发展和应用

  18. 基于Kafka+Flink+Redis的电商大屏实时计算案例

  19. Flink实战-贝壳找房基于Flink的实时平台建设

  20. 用Flink取代Spark Streaming!知乎实时数仓架构演进

  21. Flink实时数仓-美团点评实战

  22. 来将可留姓名?Flink最强学习资源合集!

  23. 数据不撒谎,Flink-Kafka性能压测全记录!

  24. 菜鸟在物流场景中基于Flink的流计算实践

  25. 基于Flink构建实时数据仓库

  26. Flink/Spark 如何实现动态更新作业配置

  27. Apache Flink结合Apache Kafka实现端到端的一致性语义

  28. Flink在滴滴出行的应用与实践

  29. 批流统一计算引擎的动力源泉—Flink Shuffle机制的重构与优化

  30. 腾讯基于Flink的实时流计算平台演进之路

  31. Flink进阶-Flink CEP(复杂事件处理)

  32. Flink基于EventTime和WaterMark处理乱序事件和晚到的数据

  33. Flink 最锋利的武器:Flink SQL 入门和实战

  34. Flink Back Pressure

  35. 最火的实时计算框架Flink和下一代分布式消息队列Pulsar的批流融合

  36. Flink Exactly-Once 投递实现浅析

  37. Flink 网络传输优化技术

  38. Apache Flink 在快手的应用与实践

  39. 基于Flink构建实时数据仓库

  40. 菜鸟在物流场景中基于Flink的流计算实践

  41. 数据不撒谎,Flink-Kafka性能压测全记录!

  42. 腾讯基于 Flink 的实时流计算平台演进之路

  43. Flink CEP在哈啰出行的应用

  44. 基于Flink和Kafka构建批流一体的数据集成平台

  45. 补齐短板 | 看阿里巴巴如何玩转Flink+Hive

  46. Flink实时数仓 | 美团点评实战

  47. 用Flink取代Spark Streaming!知乎实时数仓架构演进

  48. Flink CheckPoint奇技淫巧 | 原理和在生产中的应用

  49. 使用 Apache Flink 开发实时ETL

  50. 使用 Kubernetes 部署 Flink 应用

  51. Flink 实战 | 贝壳找房基于Flink的实时平台建设

SQL

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
-- 创建订单表连接器  
CREATE TABLE orders
(
order_id STRING,
user_id STRING,
amount DOUBLE,
order_time TIMESTAMP(3),
ptime as PROCTIME(),
rtime as CAST(order_time AS TIMESTAMP_LTZ(3)),
WATERMARK FOR rtime AS rtime - INTERVAL '5' SECOND
) WITH (
'connector' = 'filesystem',
'path' = 'file://${order.file.path}',
'format' = 'csv',
'csv.ignore-parse-errors' = 'true',
'csv.allow-comments' = 'true'
);
--
CREATE TABLE console1
(
user_id STRING,
date_date STRING,
date_time STRING,
msg_time STRING,
total_amount DOUBLE,
order_times BIGINT
)
with (
'connector' = 'console',
'connector.name' = 'csl1'
);
-- 统计一天内,各用户的交易总额和支付次数,10分钟一次输出
INSERT INTO console1
SELECT user_id
, DATE_FORMAT(window_start, 'yyyyMMdd') AS data_date
, DATE_FORMAT(window_start, 'HH:mm:ss') AS data_time
, DATE_FORMAT(window_end, 'yyyy-MM-dd HH:mm:ss') AS msg_time
, SUM(amount) AS total_amount
, count(distinct order_id) AS video_exp_cnt
FROM
TABLE(
CUMULATE(
TABLE orders
, DESCRIPTOR(
rtime
)
, INTERVAL '10' MINUTES
, INTERVAL '1' DAYS
, INTERVAL '0' HOURS
, TRUE
)
)
GROUP BY user_id
, window_start
, window_end

RelNodes

1
2
3
4
5
6
7
8
Sink(table=[default_catalog.default_database.console1], fields=[user_id, data_date, data_time, msg_time, total_amount, video_exp_cnt])
+- Calc(select=[user_id, DATE_FORMAT(window_start, _UTF-16LE'yyyyMMdd') AS data_date, DATE_FORMAT(window_start, _UTF-16LE'HH:mm:ss') AS data_time, DATE_FORMAT(window_end, _UTF-16LE'yyyy-MM-dd HH:mm:ss') AS msg_time, total_amount, video_exp_cnt])
+- GlobalWindowAggregate(groupBy=[user_id], window=[CUMULATE(slice_end=[$slice_end], max_size=[86400000 ms], step=[10 min], offset=[0 ms], allowLazyTriggering=[true])], select=[user_id, SUM(sum$0) AS total_amount, COUNT(distinct$0 count$1) AS video_exp_cnt, start('w$) AS window_start, end('w$) AS window_end])
+- Exchange(distribution=[hash[user_id]])
+- LocalWindowAggregate(groupBy=[user_id], window=[CUMULATE(time_col=[order_time], max_size=[86400000 ms], step=[10 min], offset=[0 ms], allowLazyTriggering=[true])], select=[user_id, SUM(amount) AS sum$0, COUNT(distinct$0 order_id) AS count$1, DISTINCT(order_id) AS distinct$0, slice_end('w$) AS $slice_end])
+- Calc(select=[user_id, amount, order_id, order_time])
+- WatermarkAssigner(rowtime=[order_time], watermark=[-(order_time, 5000:INTERVAL SECOND)])
+- TableSourceScan(table=[[default_catalog, default_database, orders]], fields=[order_id, user_id, amount, order_time])

Transformations

过程

Transformation Description
SourceTransformation [1]:TableSourceScan(table=[[default_catalog, default_database, orders]], fields=[order_id, user_id, amount, order_time])
OneInputTransformation [2]:Calc(select=[order_id, user_id, amount, order_time, PROCTIME() AS ptime, CAST(order_time AS TIMESTAMP_WITH_LOCAL_TIME_ZONE(3)) AS rtime])
OneInputTransformation [3]:WatermarkAssigner(rowtime=[rtime], watermark=[(rtime - 5000:INTERVAL SECOND)])
OneInputTransformation [4]:Calc(select=[user_id, amount, order_id, ptime])
PartitionTransformation [5]:Exchange(distribution=[hash[user_id]])
OneInputTransformation [6]:WindowAggregate(groupBy=[user_id], window=[CUMULATE(time_col=[ptime], max_size=[86400000 ms], step=[10 min], offset=[0 ms], allowLazyTriggering=[true])], select=[user_id, SUM(amount) AS total_amount, COUNT(DISTINCT order_id) AS video_exp_cnt, start('w$) AS window_start, end('w$) AS window_end])
OneInputTransformation [7]:Calc(select=[user_id, DATE_FORMAT(window_start, 'yyyyMMdd') AS data_date, DATE_FORMAT(window_start, 'HH:mm:ss') AS data_time, DATE_FORMAT(window_end, 'yyyy-MM-dd HH:mm:ss') AS msg_time, total_amount, video_exp_cnt])
LegacySinkTransformation

new Source API

  • 问题1: SQL -> Operation -> Transformation

Codegens

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
Source: orders[1] -> WatermarkAssigner[2] -> Calc[3] -> LocalWindowAggregate[4] (7/8)#0 - The compiled code for StreamExecCalc$11:   
/* 1 */
/* 2 */ public class StreamExecCalc$11 extends org.apache.flink.table.runtime.operators.TableStreamOperator
/* 3 */ implements org.apache.flink.streaming.api.operators.OneInputStreamOperator {
/* 4 */
/* 5 */ private final Object[] references;
/* 6 */ private transient org.apache.flink.table.runtime.typeutils.StringDataSerializer typeSerializer$5;
/* 7 */ org.apache.flink.table.data.BoxedWrapperRowData out = new org.apache.flink.table.data.BoxedWrapperRowData(4);
/* 8 */ private final org.apache.flink.streaming.runtime.streamrecord.StreamRecord outElement = new org.apache.flink.streaming.runtime.streamrecord.StreamRecord(null);
/* 9 */
/* 10 */ public StreamExecCalc$11(
/* 11 */ Object[] references,
/* 12 */ org.apache.flink.streaming.runtime.tasks.StreamTask task,
/* 13 */ org.apache.flink.streaming.api.graph.StreamConfig config,
/* 14 */ org.apache.flink.streaming.api.operators.Output output,
/* 15 */ org.apache.flink.streaming.runtime.tasks.ProcessingTimeService processingTimeService) throws Exception {
/* 16 */ this.references = references;
/* 17 */ typeSerializer$5 = (((org.apache.flink.table.runtime.typeutils.StringDataSerializer) references[0]));
/* 18 */ this.setup(task, config, output);
/* 19 */ if (this instanceof org.apache.flink.streaming.api.operators.AbstractStreamOperator) {
/* 20 */ ((org.apache.flink.streaming.api.operators.AbstractStreamOperator) this)
/* 21 */ .setProcessingTimeService(processingTimeService);
/* 22 */ }
/* 23 */ }
/* 24 */
/* 25 */ @Override
/* 26 */ public void open() throws Exception {
/* 27 */ super.open();
/* 28 */
/* 29 */ }
/* 30 */
/* 31 */ @Override
/* 32 */ public void processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord element) throws Exception {
/* 33 */ org.apache.flink.table.data.RowData in1 = (org.apache.flink.table.data.RowData) element.getValue();
/* 34 */
/* 35 */ org.apache.flink.table.data.binary.BinaryStringData field$4;
/* 36 */ boolean isNull$4;
/* 37 */ org.apache.flink.table.data.binary.BinaryStringData field$6;
/* 38 */ double field$7;
/* 39 */ boolean isNull$7;
/* 40 */ org.apache.flink.table.data.binary.BinaryStringData field$8;
/* 41 */ boolean isNull$8;
/* 42 */ org.apache.flink.table.data.binary.BinaryStringData field$9;
/* 43 */ org.apache.flink.table.data.TimestampData field$10;
/* 44 */ boolean isNull$10;
/* 45 */
/* 46 */
/* 47 */
/* 48 */ isNull$8 = in1.isNullAt(0);
/* 49 */ field$8 = org.apache.flink.table.data.binary.BinaryStringData.EMPTY_UTF8;
/* 50 */ if (!isNull$8) {
/* 51 */ field$8 = ((org.apache.flink.table.data.binary.BinaryStringData) in1.getString(0));
/* 52 */ }
/* 53 */ field$9 = field$8;
/* 54 */ if (!isNull$8) {
/* 55 */ field$9 = (org.apache.flink.table.data.binary.BinaryStringData) (typeSerializer$5.copy(field$9));
/* 56 */ }
/* 57 */
/* 58 */ isNull$7 = in1.isNullAt(2);
/* 59 */ field$7 = -1.0d;
/* 60 */ if (!isNull$7) {
/* 61 */ field$7 = in1.getDouble(2);
/* 62 */ }
/* 63 */
/* 64 */ isNull$4 = in1.isNullAt(1);
/* 65 */ field$4 = org.apache.flink.table.data.binary.BinaryStringData.EMPTY_UTF8;
/* 66 */ if (!isNull$4) {
/* 67 */ field$4 = ((org.apache.flink.table.data.binary.BinaryStringData) in1.getString(1));
/* 68 */ }
/* 69 */ field$6 = field$4;
/* 70 */ if (!isNull$4) {
/* 71 */ field$6 = (org.apache.flink.table.data.binary.BinaryStringData) (typeSerializer$5.copy(field$6));
/* 72 */ }
/* 73 */
/* 74 */ isNull$10 = in1.isNullAt(3);
/* 75 */ field$10 = null;
/* 76 */ if (!isNull$10) {
/* 77 */ field$10 = in1.getTimestamp(3, 3);
/* 78 */ }
/* 79 */
/* 80 */ out.setRowKind(in1.getRowKind());
/* 81 */
/* 82 */
/* 83 */
/* 84 */ /*+ [CODE HINT START] ROW_FIELD_0 */
/* 85 */
/* 86 */ if (isNull$4) {
/* 87 */ out.setNullAt(0);
/* 88 */ } else {
/* 89 */ out.setNonPrimitiveValue(0, field$6);
/* 90 */ }
/* 91 */ /*+ [CODE HINT END] ROW_FIELD_0 */
/* 92 */
/* 93 */
/* 94 */ /*+ [CODE HINT START] ROW_FIELD_1 */
/* 95 */
/* 96 */ if (isNull$7) {
/* 97 */ out.setNullAt(1);
/* 98 */ } else {
/* 99 */ out.setDouble(1, field$7);
/* 100 */ }
/* 101 */ /*+ [CODE HINT END] ROW_FIELD_1 */
/* 102 */
/* 103 */
/* 104 */ /*+ [CODE HINT START] ROW_FIELD_2 */
/* 105 */
/* 106 */ if (isNull$8) {
/* 107 */ out.setNullAt(2);
/* 108 */ } else {
/* 109 */ out.setNonPrimitiveValue(2, field$9);
/* 110 */ }
/* 111 */ /*+ [CODE HINT END] ROW_FIELD_2 */
/* 112 */
/* 113 */
/* 114 */ /*+ [CODE HINT START] ROW_FIELD_3 */
/* 115 */
/* 116 */ if (isNull$10) {
/* 117 */ out.setNullAt(3);
/* 118 */ } else {
/* 119 */ out.setNonPrimitiveValue(3, field$10);
/* 120 */ }
/* 121 */ /*+ [CODE HINT END] ROW_FIELD_3 */
/* 122 */
/* 123 */
/* 124 */ output.collect(outElement.replace(out));
/* 125 */
/* 126 */
/* 127 */ }
/* 128 */
/* 129 */
/* 130 */
/* 131 */ @Override
/* 132 */ public void close() throws Exception {
/* 133 */ super.close();
/* 134 */
/* 135 */ }
/* 136 */
/* 137 */
/* 138 */ }
/* 139 */

Load module

1
LOAD MODULE hive WITH ('hive-version' = '3.1.2')

org.apache.flink.table.api.internal.TableEnvironmentImpl#executeInternal(org.apache.flink.table.operations.Operation)

org.apache.flink.table.api.internal.TableEnvironmentImpl#loadModule(org.apache.flink.table.operations.LoadModuleOperation)

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

date created: 2025-12-23 09:47
tags:

  • 时间
  • Timestmap
  • Flink
  • 时区
    category: Bigdata
    date updated: 2026-08-17 20:24
    title: 各引擎timestamp格式名称对齐

1 各引擎timestamp格式名称对齐

1.1.1 Timestmap

Flink Iceberg Spark MySQL
timestmap 默认就是毫秒精度 timestamp timestampntz DATETIME(3)
timestampltz timestamptz timestamp

从DLA拉取元数据时,默认精度为MILLIS

1.1.2 时间精度

类型 精度 单位
TIMESTAMP(0) 0 位小数
TIMESTAMP(1) 1 位小数 0.1 秒(100ms)
TIMESTAMP(2) 2 位小数 10ms
TIMESTAMP(3) 3 位小数 毫秒(ms)
TIMESTAMP(6) 6 位小数 微秒(μs)
TIMESTAMP(9) 9 位小数 纳秒(ns)

1.1.3 避免隐式转换

1
2
3
4
5
6
7
8
9
10
11
12
13
CREATE TEMPORARY VIEW SalesDataSource AS
SELECT
*
FROM
(
VALUES
(1766073600000),
(1766080800000)
) AS T(testt);
insert into mysql_table
select
CAST(FROM_UNIXTIME(testt/1000) AS DATE) as day_time
from SalesDataSource

MySQL中day_time的字段类型是DATEDATETIME,输入的结果总是1970

  • testtBIGINT(毫秒时间戳)
  • FROM_UNIXTIME(testt/1000) 返回的是 TIMESTAMP
  • CAST(... AS DATE) 会变成 Flink 的 DATE 类型
  • 如果 MySQL 字段是 DATETIMETIMESTAMP,写入时 Flink JDBC connector 会把 DATE 转换为 0,显示成 1970-01-01

解决方案: 避免隐式转换

1
2
3
4
-- Flink的Timestamp 对应的类型就是 DateTime(3), 无需强制转换
FROM_UNIXTIME(testt/1000) AS day_time
-- 显示Cast成TIMESTAMP(3)也可以
CAST(FROM_UNIXTIME(testt/1000) AS TIMESTAMP(3)) AS day_time
Flink SQL Connector/Java MySQL 类型
DATE java.sql.Date DATE 只有年月日,精度到天
TIMESTAMP java.sql.Timestamp DATETIME或TIMESTAMP 精度到毫秒

Flink SQL中使用CAST(... AS DATE)写入MySQL DATETIME, JDBC Sink会把毫秒部分丢掉,甚至映射失败,默认值会变成0->1970-01-01