0%

一个SQL了解Flink运行原理: Flink SQL window TopN

Flink SQL支持增量窗口TopN

假设定义数据源:

字段 类型 备注
category String 股票分类
field String 股票id
ftime String 字符串类型时间,格式为yyyyMMddHHmmss

样例数据:

category field ftime
k stock_id-1 20220930164620
k stock_id-0 20220930164625
k stock_id-1 20220930164625
k stock_id-0 20220930164625
k stock_id-1 20220930164625
k stock_id-0 20220930164625
k stock_id-0 20220930164625
k stock_id-0 20220930164625
k stock_id-0 20220930164625
k stock_id-0 20220930164625
k stock_id-1 20220930164630
k stock_id-1 20220930164630
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
select 
fixedTime(stat_time, 'HH:mm:ss') as time,
category,
field,
pv,
rk
from (
SELECT category,
field,
stat_time,
pv,
row_number() over (partition by category,INCREMENT(stat_time, INTERVAL '5' SECOND, INTERVAL '5' SECOND, Time '00:00:00', true) order by pv desc) as rk
FROM (SELECT category,
field,
INCREMENT_SEGMENT_ROWTIME(rtime, INTERVAL '1' DAY, INTERVAL '5' SECOND, Time '16:00:00', true) as stat_time,
count(1) as pv
FROM PageViews
GROUP BY category,
field,
INCREMENT(rtime, INTERVAL '1' DAY, INTERVAL '5' SECOND, Time '16:00:00', true)
)
)
where rk <= 10

计算结果:

row_kind time category field pv rk
true 16:49:54 k stock_id-1 2 1
true 16:49:54 k stock_id-0 1 2
true 16:49:59 k stock_id-1 8 1
true 16:49:59 k stock_id-0 5 2
true 16:50:04 k stock_id-1 12 1
true 16:50:04 k stock_id-0 11 2
true 16:50:09 k stock_id-0 19 1
true 16:50:09 k stock_id-1 14 2
true 16:50:14 k stock_id-0 23 1
true 16:50:14 k stock_id-1 20 2
true 16:50:19 k stock_id-1 27 1
true 16:50:19 k stock_id-0 26 2
1
2
3
4
5
6
7
8
9
Plan after converting SqlNode to RelNode:

LogicalProject(EXPR$0=[FixedTime($2, _UTF-16LE'HH:mm:ss')], category=[$0], field=[$1], pv=[$3], rk=[$4])
LogicalFilter(condition=[<=($4, 10)])
LogicalProject(category=[$0], field=[$1], stat_time=[$2], pv=[$3], rk=[ROW_NUMBER() OVER (PARTITION BY $0, INCREMENT($2, 5000:INTERVAL SECOND, 5000:INTERVAL SECOND, 00:00:00, true) ORDER BY $3 DESC)])
LogicalProject(category=[$0], field=[$1], stat_time=[INCREMENT_SEGMENT_ROWTIME($2)], pv=[$3])
LogicalAggregate(group=[{0, 1, 2}], pv=[COUNT()])
LogicalProject(category=[$0], field=[$1], rtime=[INCREMENT($2, 86400000:INTERVAL DAY, 5000:INTERVAL SECOND, 16:00:00, true)], $f3=[1])
FlinkLogicalDataStreamScan(id=[2], fields=[category, field, rtime])

Source 里面有三个字段: category, field, rtime
其中的rtime字段是内置字段rtime,来标示事件时间。通过tableEnv.registerDataStream("PageViews", stockClickStream, "category, field, rtime.rowtime"); 指定。

1
2
3
4
5
6
7
8
9
10
org.apache.calcite.plan.RelOptPlanner org.apache.calcite.plan.hep.HepPlanner.dumpGraph(HepPlanner.java:1027)- 
Breadth-first from root: {
rel#20:HepRelVertex#20 = rel#19:LogicalProject.NONE(input=HepRelVertex#18,exprs=[_UTF-16LE'Top10', FixedTime($2, _UTF-16LE'yyyyMMdd-HHmmss'), $0, $1, $3, $4]), rowcount=5.0, cumulative cost={241.25 rows, 630.0 cpu, 3600.0 io}
rel#18:HepRelVertex#18 = rel#17:LogicalFilter.NONE(input=HepRelVertex#16,condition=<=($4, 10)), rowcount=5.0, cumulative cost={236.25 rows, 600.0 cpu, 3600.0 io}
rel#16:HepRelVertex#16 = rel#15:LogicalProject.NONE(input=HepRelVertex#14,inputs=0..3,exprs=[ROW_NUMBER() OVER (PARTITION BY $0, INCREMENT($2, 5000:INTERVAL SECOND, 5000:INTERVAL SECOND, 00:00:00, true) ORDER BY $3 DESC)]), rowcount=10.0, cumulative cost={231.25 rows, 590.0 cpu, 3600.0 io}
rel#14:HepRelVertex#14 = rel#13:LogicalProject.NONE(input=HepRelVertex#12,inputs=0..1,exprs=[INCREMENT_SEGMENT_ROWTIME($2), $3]), rowcount=10.0, cumulative cost={221.25 rows, 540.0 cpu, 3600.0 io}
rel#12:HepRelVertex#12 = rel#11:LogicalAggregate.NONE(input=HepRelVertex#10,group={0, 1, 2},pv=COUNT()), rowcount=10.0, cumulative cost={211.25 rows, 500.0 cpu, 3600.0 io}
rel#10:HepRelVertex#10 = rel#9:LogicalProject.NONE(input=HepRelVertex#8,inputs=0..1,exprs=[INCREMENT($2, 86400000:INTERVAL DAY, 5000:INTERVAL SECOND, 16:00:00, true), 1]), rowcount=100.0, cumulative cost={200.0 rows, 500.0 cpu, 3600.0 io}
rel#8:HepRelVertex#8 = rel#1:FlinkLogicalDataStreamScan.LOGICAL(id=2,fields=page, event, rtime), rowcount=100.0, cumulative cost={100.0 rows, 100.0 cpu, 3600.0 io}
}

cumulative cost: 累计代价