0%

flink窗口函数包含滚动窗口、滑动窗口、会话窗口和OVER窗口

滚动窗口

滚动窗口(TUMBLE)将每个元素分配到一个指定大小的窗口中。通常,滚动窗口有一个固定的大小,并且不会出现重叠。例如,如果指定了一个5分钟大小的滚动窗口,无限流的数据会根据时间划分为[0:00 - 0:05)[0:05, 0:10)[0:10, 0:15)等窗口。下图展示了一个30秒的滚动窗口。

使用标识函数选出窗口的起始时间或者结束时间,窗口的时间属性用于下级Window的聚合。

窗口标识函数 返回类型 描述
TUMBLE_START(time-attr, size-interval) TIMESTAMP 返回窗口的起始时间(包含边界)。例如[00:10, 00:15) 窗口,返回00:10
TUMBLE_END(time-attr, size-interval) TIMESTAMP 返回窗口的结束时间(包含边界)。例如[00:00, 00:15]窗口,返回00:15
TUMBLE_ROWTIME(time-attr, size-interval) TIMESTAMP(rowtime-attr) 返回窗口的结束时间(不包含边界)。例如[00:00, 00:15]窗口,返回00:14:59.999 。返回值是一个rowtime attribute,即可以基于该字段做时间属性的操作,例如,级联窗口只能用在基于Event Time的Window上
TUMBLE_PROCTIME(time-attr, size-interval) TIMESTAMP(rowtime-attr) 返回窗口的结束时间(不包含边界)。例如[00:00, 00:15]窗口,返回00:14:59.999。返回值是一个proctime attribute,即可以基于该字段做时间属性的操作,例如,级联窗口只能用在基于Processing Time的Window上

TUMBLE window示例

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
import org.apache.flink.api.common.typeinfo.TypeHint;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.java.tuple.Tuple3;
import org.apache.flink.streaming.api.TimeCharacteristic;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.timestamps.AscendingTimestampExtractor;
import org.apache.flink.table.api.EnvironmentSettings;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;


import java.sql.Timestamp;
import java.util.Arrays;

public class TumbleWindowExample {

public static void main(String[] args) throws Exception {

/**
* 1 注册环境
*/
EnvironmentSettings mySetting = EnvironmentSettings
.newInstance()
// .useOldPlanner()
.useBlinkPlanner()
.inStreamingMode()
.build();

// 获取 environment
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 指定系统时间概念为 event time
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);

StreamTableEnvironment tEnv = StreamTableEnvironment.create(env,mySetting);


// 初始数据
DataStream<Tuple3<Long, String,Integer>> log = env.fromCollection(Arrays.asList(
//时间 14:53:00
new Tuple3<>(1572591180_000L,"xiao_ming",300),
//时间 14:53:09
new Tuple3<>(1572591189_000L,"zhang_san",303),
//时间 14:53:12
new Tuple3<>(1572591192_000L, "xiao_li",204),
//时间 14:53:21
new Tuple3<>(1572591201_000L,"li_si", 208)
));

// 指定时间戳
SingleOutputStreamOperator<Tuple3<Long, String, Integer>> logWithTime = log.assignTimestampsAndWatermarks(new AscendingTimestampExtractor<Tuple3<Long, String, Integer>>() {

@Override
public long extractAscendingTimestamp(Tuple3<Long, String, Integer> element) {
return element.f0;
}
});

// 转换为 Table
Table logT = tEnv.fromDataStream(logWithTime, "t.rowtime, name, v");

Table result = tEnv.sqlQuery("SELECT TUMBLE_START(t, INTERVAL '10' SECOND) AS window_start," +
"TUMBLE_END(t, INTERVAL '10' SECOND) AS window_end, SUM(v) FROM "
+ logT + " GROUP BY TUMBLE(t, INTERVAL '10' SECOND)");

TypeInformation<Tuple3<Timestamp,Timestamp,Integer>> tpinf = new TypeHint<Tuple3<Timestamp,Timestamp,Integer>>(){}.getTypeInfo();
tEnv.toAppendStream(result, tpinf).print();

env.execute();
}


}

sql逻辑,每十秒钟聚合
执行结果:

1
2
3
(2019-11-01 06:53:00.0,2019-11-01 06:53:10.0,603)  
(2019-11-01 06:53:20.0,2019-11-01 06:53:30.0,208)
(2019-11-01 06:53:10.0,2019-11-01 06:53:20.0,204)

滑动窗口

滑动窗口(HOP),也被称作Sliding Window。不同于滚动窗口,滑动窗口的窗口可以重叠。

滑动窗口有两个参数:slide和size。slide为每次滑动的步长,size为窗口的大小。

  • slide < size,则窗口会重叠,每个元素会被分配到多个窗口。
  • slide = size,则等同于滚动窗口(TUMBLE)。
  • slide > size,则为跳跃窗口,窗口之间不重叠且有间隙。

通常,大部分元素符合多个窗口情景,窗口是重叠的。因此,滑动窗口在计算移动平均数(moving averages)时很实用。例如,计算过去5分钟数据的平均值,每10秒钟更新一次,可以设置slide为10秒,size为5分钟。下图为您展示间隔为30秒,窗口大小为1分钟的滑动窗口。

滑动窗口

使用滑动窗口标识函数选出窗口的起始时间或者结束时间,窗口的时间属性用于下级Window的聚合。

窗口标识函数 返回类型 描述
HOP_START(<time-attr>, <slide-interval>, <size-interval>) TIMESTAMP 返回窗口的起始时间(包含边界)。例如[00:10, 00:15) 窗口,返回00:10
HOP_END(<time-attr>, <slide-interval>, <size-interval>) TIMESTAMP 返回窗口的结束时间(包含边界)。例如[00:00, 00:15) 窗口,返回00:15
HOP_ROWTIME(<time-attr>, <slide-interval>, <size-interval>) TIMESTAMP(rowtime-attr) 返回窗口的结束时间(不包含边界)。例如[00:00, 00:15) 窗口,返回00:14:59.999。返回值是一个rowtime attribute,即可以基于该字段做时间类型的操作,只能用在基于event time的window上。
HOP_PROCTIME(<time-attr>, <slide-interval>, <size-interval>) TIMESTAMP(rowtime-attr) 返回窗口的结束时间(不包含边界)。例如[00:00, 00:15) 窗口,返回00:14:59.999 。返回值是一个proctime attribute

滑动窗口实例:
java代码同上,sql语句改为:

1
2
3
4
5
6
SELECT 
HOP_START(t, INTERVAL '5' SECOND, INTERVAL '10' SECOND) AS window_start,
HOP_END(t, INTERVAL '5' SECOND, INTERVAL '10' SECOND) AS window_end,
SUM(v)
FROM logT
GROUP BY HOP(t, INTERVAL '5' SECOND, INTERVAL '10' SECOND)

每间隔5秒统计10秒内的数据
sql结果如下:

1
2
3
4
5
6
(2019-11-01 06:53:15.0,2019-11-01 06:53:25.0,208)  
(2019-11-01 06:53:10.0,2019-11-01 06:53:20.0,204)
(2019-11-01 06:53:05.0,2019-11-01 06:53:15.0,507)
(2019-11-01 06:53:20.0,2019-11-01 06:53:30.0,208)
(2019-11-01 06:53:00.0,2019-11-01 06:53:10.0,603)
(2019-11-01 06:52:55.0,2019-11-01 06:53:05.0,300)

会话窗口

会话窗口(SESSION)通过Session活动来对元素进行分组。会话窗口与滚动窗口和滑动窗口相比,没有窗口重叠,没有固定窗口大小。相反,当它在一个固定的时间周期内不再收到元素,即会话断开时,这个窗口就会关闭。

会话窗口通过一个间隔时间(Gap)来配置,这个间隔定义了非活跃周期的长度。例如,一个表示鼠标点击活动的数据流可能具有长时间的空闲时间,并在两段空闲之间散布着高浓度的点击。 如果数据在指定的间隔(Gap)之后到达,则会开始一个新的窗口。

会话窗口示例如下图。每个Key由于不同的数据分布,形成了不同的Window。


使用标识函数选出窗口的起始时间或者结束时间,窗口的时间属性用于下级Window的聚合。

窗口标识函数 返回类型 描述
SESSION_START(<time-attr>, <gap-interval>) Timestamp 返回窗口的起始时间(包含边界)。如[00:10, 00:15) 的窗口,返回 00:10 ,即为此会话窗口内第一条记录的时间。
SESSION_END(<time-attr>, <gap-interval>) Timestamp 返回窗口的结束时间(包含边界)。如[00:00, 00:15) 的窗口,返回 00:15,即为此会话窗口内最后一条记录的时间+<gap-interval>
SESSION_ROWTIME(<time-attr>, <gap-interval>) Timestamp(rowtime-attr) 返回窗口的结束时间(不包含边界)。如 [00:00, 00:15) 的窗口,返回00:14:59.999 。返回值是一个rowtime attribute,也就是可以基于该字段进行时间类型的操作。该参数只能用于基于event time的window 。
SESSION_PROCTIME(<time-attr>, <gap-interval>) Timestamp(rowtime-attr) 返回窗口的结束时间(不包含边界)。如 [00:00, 00:15) 的窗口,返回 00:14:59.999 。返回值是一个 proctime attribute,也就是可以基于该字段进行时间类型的操作。该参数只能用于基于processing time的window 。

会话窗口实例:
java代码同上
sql语句如下:
每隔5秒聚合

1
2
3
4
5
6
SELECT 
SESSION_START(t, INTERVAL '5' SECOND) AS window_start,
SESSION_END(t, INTERVAL '5' SECOND) AS window_end,
SUM(v)
FROM logT
GROUP BY SESSION(t, INTERVAL '5' SECOND)

sql结果:

1
2
3
(2019-11-01 06:53:21.0,2019-11-01 06:53:26.0,208)  
(2019-11-01 06:53:00.0,2019-11-01 06:53:05.0,300)
(2019-11-01 06:53:09.0,2019-11-01 06:53:17.0,507)

OVER窗口

OVER窗口(OVER Window)是传统数据库的标准开窗,不同于Group By Window,OVER窗口中每1个元素都对应1个窗口。窗口内的元素是当前元素往前多少个或往前多长时间的元素集合,因此流数据元素分布在多个窗口中。

在应用OVER窗口的流式数据中,每1个元素都对应1个OVER窗口。每1个元素都触发1次数据计算,每个触发计算的元素所确定的行,都是该元素所在窗口的最后1行。在实时计算的底层实现中,OVER窗口的数据进行全局统一管理(数据只存储1份),逻辑上为每1个元素维护1个OVER窗口,为每1个元素进行窗口计算,完成计算后会清除过期的数据。

Flink SQL中对OVER窗口的定义遵循标准SQL的定义语法,传统OVER窗口没有对其进行更细粒度的窗口类型命名划分。按照计算行的定义方式,OVER Window可以分为以下两类:

  • ROWS OVER Window:每一行元素都被视为新的计算行,即每一行都是一个新的窗口。
  • RANGE OVER Window:具有相同时间值的所有元素行视为同一计算行,即具有相同时间值的所有行都是同一个窗口。

Rows OVER Window语义

窗口数据

ROWS OVER Window的每个元素都确定一个窗口。ROWS OVER Window分为Unbounded(无界流)和Bounded(有界流)两种情况。
Unbounded ROWS OVER Window数据示例如下图所示。

虽然上图所示窗口user1的w7、w8及user2的窗口w3、w4都是同一时刻到达,但它们仍然在不同的窗口,这一点与RANGE OVER Window不同。

Bounded ROWS OVER Window数据以3个元素(往前2个元素)的窗口为例,如下图所示。

虽然上图所示窗口user1的w5、w6及user2的窗口w1、w2都是同一时刻到达,但它们仍然在不同的窗口,这一点与RANGE OVER Window不同。

RANGE OVER Window语义

窗口数据

RANGE OVER Window所有具有共同元素值(元素时间戳)的元素行确定一个窗口,RANGE OVER Window分为Unbounded和Bounded的两种情况。
Unbounded RANGE OVER Window数据示例如下图所示。


上图所示窗口user1的w7、user2的窗口w3,两个元素同一时刻到达,属于相同的window,这一点与ROWS OVER Window不同。

Bounded RANGE OVER Window数据,以3秒中数据(INTERVAL '2' SECOND)的窗口为例,如下图所示。

上图所示窗口user1的w6、user2的窗口w3,元素都是同一时刻到达,属于相同的window,这一点与ROWS OVER Window不同。

OVER窗口实例:
java代码同上
初始数据如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
// 初始数据
DataStream<Tuple3<Long, String,Integer>> log = env.fromCollection(Arrays.asList(
//时间 14:53:00
new Tuple3<>(1572591180_000L,"xiao_ming",999),
//时间 14:53:09
new Tuple3<>(1572591189_000L,"zhang_san",303),
//时间 14:53:12
new Tuple3<>(1572591192_000L, "xiao_li",888),
//时间 14:53:21
new Tuple3<>(1572591201_000L,"li_si", 908),
//2019-11-01 14:53:31
new Tuple3<>(1572591211_000L,"li_si", 555),
//2019-11-01 14:53:41
new Tuple3<>(1572591221_000L,"zhang_san", 666),
//2019-11-01 14:53:51
new Tuple3<>(1572591231_000L,"xiao_ming", 777),
//2019-11-01 14:54:01
new Tuple3<>(1572591241_000L,"xiao_ming", 213),
//2019-11-01 14:54:11
new Tuple3<>(1572591251_000L,"zhang_san", 300),
//2019-11-01 14:54:21
new Tuple3<>(1572591261_000L,"li_si", 112)
));

ROWS over Windown sql语句如下:

1
2
3
4
5
6
7
8
9
SELECT 
name,
v,
MAX(v) OVER(
PARTITION BY name
ORDER BY t
ROWS BETWEEN 2 PRECEDING AND CURRENT ROW
)
FROM logT

sql结果如下:

1
2
3
4
5
6
7
8
9
10
11
(zhang_san,303,303)  
(xiao_li,888,888)
(li_si,908,908)
(xiao_ming,999,999)
(zhang_san,666,666)
(li_si,555,908)
(xiao_ming,777,999)
(li_si,112,908)
(zhang_san,300,666)
(xiao_ming,213,999)

RANGE OVER Window sql 语句如下:

1
2
3
4
5
6
7
8
9
SELECT 
name,
v,
MAX(v) OVER(
PARTITION BY name
ORDER BY t
RANGE BETWEEN INTERVAL '15' SECOND PRECEDING AND CURRENT ROW
)
FROM logT

sql结果如下:

1
2
3
4
5
6
7
8
9
10
(xiao_ming,999,999)  
(xiao_li,888,888)
(zhang_san,303,303)
(li_si,908,908)
(li_si,555,908)
(xiao_ming,777,777)
(zhang_san,666,666)
(li_si,112,112)
(xiao_ming,213,777)
(zhang_san,300,300)

本文的java代码来自:
https://github.com/CheckChe08...

底层实现

https://mp.weixin.qq.com/s/UkpkS_JiRGR0ibZKYechbg

概述

窗口是无限流上一种核心机制,可以流分割为有限大小的“窗口”,同时,在窗口内进行聚合,从而把源源不断产生的数据根据不同的条件划分成一段一段有边界的数据区间,使用户能够利用窗口功能实现很多复杂的统计分析需求。

Window分类

1、TimeWindow与CountWindow Flink Window可以是时间驱动的(TimeWindow),也可以是数据驱动的(CountWindow)。由于flink-planner-blink SQL中目前只支持TimeWindow相应的表达语句(TUMBLEHOPSESSION),因此,本文主要介绍TimeWindow SQL示例和逻辑,CountWindow感兴趣的读者可自行分析。

2、TimeWindow子类型 Flink TimeWindow有滑动窗口(HOP)、滚动窗口(TUMBLE)以及会话窗口(SESSION)三种,所选取的字段时间,可以是系统时间(PROCTIME)或事件时间(EVENT TIME)两种,接来下依次介绍。

Tumble Window(滚动窗口)

翻转窗口Assigner将每个元素分配给具有指定大小的窗口。翻转窗口的大小是固定的,且不会重叠。例如,指定一个大小为5分钟的翻滚窗口,并每5分钟启动一个新窗口,如下图所示:

图片

TUMBLE ROWTIME语法示例:

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
CREATE TABLE sessionOrderTableRowtime (
ctime TIMESTAMP,
categoryName VARCHAR,
shopName VARCHAR,
itemName VARCHAR,
userId VARCHAR,
price FLOAT,
action BIGINT,
WATERMARK FOR ctime AS withOffset(ctime, 1000),
proc AS PROCTIME()
) with (
`type` = 'kafka',
format = 'json',
updateMode = 'append',
`group.id` = 'groupId',
bootstrap.servers = 'xxxxx:9092',
version = '0.10',
`zookeeper.connect` = 'xxxxx:2181',
startingOffsets = 'latest',
topic = 'sessionsourceproctime'
);


CREATE TABLE popwindowsink (
countA BIGINT,
ctime_start TIMESTAMP,
ctime_end VARCHAR,
ctime_rowtime VARCHAR,
categoryName VARCHAR,
price_sum FLOAT
) with (
format = 'json',
updateMode = 'append',
bootstrap.servers = 'xxxxx:9092',
version = '0.10',
topic = 'sessionsinkproctime',
`type` = 'kafka'
);

INSERT INTO popwindowsink
(
SELECT
COUNT(*),
TUMBLE_START(ctime, INTERVAL '5' MINUTE),
DATE_FORMAT(TUMBLE_END(ctime, INTERVAL '5' MINUTE), 'yyyy-MM-dd-HH-mm-ss:SSS'), --将TUMBLE_END转为可视化的日期
DATE_FORMAT(TUMBLE_ROWTIME(ctime, INTERVAL '5' MINUTE), 'yyyy-MM-dd-HH-mm-ss:SSS'), --这里TUMBLE_ROWTIME为TUMBLE_END-1ms,一般用于后续窗口级联聚合
categoryName,
SUM(price)
FROM sessionOrderTableRowtime
GROUP BY TUMBLE(ctime, INTERVAL '5' MINUTE), categoryName
)

TUMBLEP ROCTIME语法示例:

1
2
3
4
5
6
7
8
9
10
INSERT INTO popwindowsink
(SELECT
COUNT(*),
TUMBLE_START(proc, INTERVAL '5' MINUTE),
DATE_FORMAT(TUMBLE_END(proc, INTERVAL '5' MINUTE), 'yyyy-MM-dd-HH-mm-ss:SSS'),
DATE_FORMAT(TUMBLE_PROCTIME(proc, INTERVAL '5' MINUTE), 'yyyy-MM-dd-HH-mm-ss:SSS'), --注意这里proc字段即Source DDL中指定的PROCTIME
categoryName,
SUM(price)
FROM sessionOrderTableRowtime
GROUP BY TUMBLE(proc, INTERVAL '5' MINUTE), categoryName)

ROWTIME与PROCTIME区别:

  • 在使用上:主要是填入的ctime、proc关键字的区别,这两个字段在Source DDL中指定方式不一样.
  • 在实现原理上:ROWTIME模式,根据ctime对应的值,去确定窗口的start、end;PROCTIME模式,在WindowOperator处理数据时,获取本地系统时间,去确定窗口的start、end.

由于生产系统中,主要使用ROWTIME来计算、聚合、统计,PROCTIME一般用于测试或对统计精度要求不高的场景,本文后续都主要以ROWTIME进行分析。

Hop Window(滑动窗口)

滑动窗口Assigner将元素分配给多个固定长度的窗口。类似于滚动窗口分配程序,窗口的大小由窗口大小参数配置。因此,如果滑动窗口小于窗口大小,则滑动窗口可以重叠。在这种情况下,元素被分配到多个窗口。其实,滚动窗口TUMBLE是滑动窗口的一个特例。例子,设置一个10分钟长度的窗口,以5分钟间隔滑动。这样,每5分钟就会出现一个窗口,其中包含最近10分钟内到达的事件,如下图:

图片

HOP ROWTIME语法示例:

1
2
3
4
5
6
7
8
9
10
INSERT INTO popwindowsink
(SELECT
COUNT(*),
HOP_START(ctime, INTERVAL '5' MINUTE, INTERVAL '10' MINUTE),
DATE_FORMAT(HOP_END(ctime, INTERVAL '5' MINUTE, INTERVAL '10' MINUTE), 'yyyy-MM-dd-HH-mm-ss:SSS'),
DATE_FORMAT(HOP_ROWTIME(ctime, INTERVAL '5' MINUTE, INTERVAL '10' MINUTE), 'yyyy-MM-dd-HH-mm-ss:SSS'), --注意这里ctime字段即Source DDL中指定的ROWTIME
categoryName,
SUM(price)
FROM sessionOrderTableRowtime
GROUP BY HOP(ctime, INTERVAL '5' MINUTE, INTERVAL '10' MINUTE), categoryName)

Session Window(会话窗口)

会话窗口Assigner根据活动会话对元素进行分组。与翻滚窗口和滑动窗口相比,会话窗口不会重叠,也没有固定的开始和结束时间。相反,会话窗口在一段时间内不接收元素时关闭,即,当一段不活跃的间隙发生时,当前会话关闭,随后的元素被分配给新的会话。

图片

SESSION ROWTIME语法示例:

1
2
3
4
5
6
7
8
9
10
INSERT INTO popwindowsink
(SELECT
COUNT(*),
SESSION_START(ctime, INTERVAL '5' MINUTE),
DATE_FORMAT(SESSION_END(ctime, INTERVAL '5' MINUTE, 'yyyy-MM-dd-HH-mm-ss:SSS'),
DATE_FORMAT(SESSION_ROWTIME(ctime, INTERVAL '5' MINUTE), 'yyyy-MM-dd-HH-mm-ss:SSS'), --注意这里ctime字段即Source DDL中指定的ROWTIME
categoryName,
SUM(price)
FROM sessionOrderTableRowtime
GROUP BY SESSION(ctime, INTERVAL '5' MINUTE), categoryName)

Window分类及整体流程

图片

上图内部流程分析:

应用层SQL:
1.1 window分类及配置,包括滑动、翻转、会话类型窗口
1.2 window时间类型配置,默认待字段名的EventTime,也可以通过PROCTIME()配置为ProcessingTime
Calcite解析引擎:
2.1 Calcite SQL解析,包括逻辑、优化、物理计划和算子绑定(#translateToPlanInternal),在本文特指StreamExecGroupWindowAggregateRule和StreamExecGroupWindowAggregate物理计划
WindowOperator算子创建相关:
3.1 StreamExecGroupWindowAggregate#createWindowOperator创建算子
3.2 WindowAssigner的创建,根据输入的数据,和窗口类型,生成多个窗口
3.3 processElement()真实处理数据,包括聚合运算,生成窗口,更新缓存,提交数据等功能
3.4 Trigger根据数据或时间,来决定窗口触发

创建WindowOperator算子

由于window语法主要是在group by语句中使用,calcite创建WindowOperator算子伴随着聚合策略的实现,包括聚合规则匹配(StreamExecGroupWindowAggregateRule),以及生成聚合physical算子StreamExecGroupWindowAggregate两个子流程:

图片

上图内部流程分析:

a. StreamExecGroupWindowAggregateRule会对window进行提前匹配,
生成的WindowEmitStrategy内部具有:是否为EventTime表标识、是否为SessionWindow、early fire和late fire配置、延迟毫秒数(窗口结束时间加上这个毫秒数即数据清理时间)
b. StreamExecGroupWindowAggregateRule会获取聚合逻辑计划中,window配置的时间字段,记录时间字段index信息,window的触发和清理都会用到这个时间
c. StreamExecGroupWindowAggregate入口即为translateToPlanInternal,它的实现方式与spark比较类似,会先循环调用child子节点translateToPlan方法,生成inputtranform信息作为输入
d.创建aggregateHandler是一个代码生成的过程,其生成的创建的class实现了accumulate、retract、merge、update方法,这个handler最后也传递给了WindowOperater,处理数据时,可以进行聚合、回撤并输出最新数据给下游
e. StreamExecGroupWindowAggregate与window相关的最后一步就是调用#createWindowOperator创建算子,其内部先创建了一个WindowOperatorBuilder,设置window类型、retract标识、trigger(window触发条件)、聚合函数句柄等,最后创建WindowOperator

WindowOperator处理数据图解

在上一小节,已经完成了WindowOperator参数的设定,并创建实例,接下来我们主要分析WindowOperator真实处理数据的流程(起点在WindowOperator#processElement方法):

图片

processElement处理数据流程:

a、 获取当前record具有的事件时间,如果是Processing Time模式,从时间服务Service里面获取时间即可
b、使用上一步获取的时间,接着调用windowFunction.assignWindow生成窗口,其内部实际上是调用各类型的WindowAssigner生成窗口,windowFunction有三大类,分别是Paned(滑动)、Merge(会话)、General(前两种以外的),WindowAssigner类型大致有5类,分别是Tumbling(翻转)、Sliding(滑动)、Session(会话)、CountTumbling 、CountSlide这几类,根据输入的一条数据和时间,可以生成1到多个窗口
c、接下来是遍历涉及的窗口进行聚合,包括从windowState获取聚合前值、使用句柄进行聚合、更新状态至windowState,将当前转态
d、上一步聚合完成后,就可以遍历窗口,使用TriggerContext(其实就是不同类型窗口Trigger触发器的代理),综合early fire、late fire、水印时间与窗口结束时间,综合判断是否触发窗口写出
e、如果TriggerContext判断出触发条件为true,则调用emitWindowResult写出,其内部有retract判断,更新当前state及previous state,写出数据等操作
f、如果TriggerContext判断出触发条件为false,则触发需要注册cleanupTimer,到达指定时间后,触发onEventTime或onProcessingTime
g、onEventTime或onProcessingTime功能十分类似,首先会触发emitWindowResult提交结果,另外会判断窗口结束时间+Lateness和当前时间是否相等,相等则表示可以清除窗口数据、当前state及previous state、窗口对应trigger。

WindowOperator源码调试

为了更直观的理解Window内部运行原理,这里我们引入一个Flink源码中已有的SQL Window测试用例,并进行了简单的修改(即修改为使用HOP滑动窗口)

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
classWindowJoinITCase{
@Test
def testRowTimeInnerJoinWithWindowAggregateOnFirstTime(): Unit = {
val sqlQuery =
"""
|SELECT t1.key, HOP_END(t1.rowtime, INTERVAL '4' SECOND, INTERVAL '20' SECOND), COUNT(t1.key)
|FROM T1 AS t1
|GROUP BY HOP(t1.rowtime, INTERVAL '4' SECOND, INTERVAL '20' SECOND), t1.key
|""".stripMargin

val data1 = new mutable.MutableList[(String, String, Long)]
data1.+=(("A", "L-1", 1000L))
data1.+=(("A", "L-2", 2000L))
data1.+=(("A", "L-3", 3000L))
//data1.+=(("B", "L-8", 2000L))
data1.+=(("B", "L-4", 4000L))
data1.+=(("C", "L-5", 2100L))
data1.+=(("A", "L-6", 10000L))
data1.+=(("A", "L-7", 13000L))

val t1 = env.fromCollection(data1)
.assignTimestampsAndWatermarks(new Row3WatermarkExtractor2)
.toTable(tEnv, 'key, 'id, 'rowtime)

tEnv.registerTable("T1", t1)

val sink = new TestingAppendSink
val t_r = tEnv.sqlQuery(sqlQuery)
val result = t_r.toAppendStream[Row]
result.addSink(sink)
env.execute()
}
}

1、StreamExecGroupWindowAggregate#createWindowOperator()创建算子

StreamExecGroupWindowAggregate#createWindowOperator()是创建WindowOperator算子的地方,对应的代码和注释:

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
class StreamExecGroupWindowAggregate{
private def createWindowOperator(
config: TableConfig,
aggsHandler: GeneratedNamespaceAggsHandleFunction[_],
recordEqualiser: GeneratedRecordEqualiser,
accTypes: Array[LogicalType],
windowPropertyTypes: Array[LogicalType],
aggValueTypes: Array[LogicalType],
inputFields: Seq[LogicalType],
timeIdx: Int): WindowOperator[_, _] = {

val builder = WindowOperatorBuilder
.builder()
.withInputFields(inputFields.toArray)
val timeZoneOffset = -config.getTimeZone.getOffset(Calendar.ZONE_OFFSET)

// 设置WindowOperatorBuilder,最后通过Builder创建WindowOperator
val newBuilder = window match {
case TumblingGroupWindow(_, timeField, size) //Tumble PROCTIME模式,内部设置Assiger
if isProctimeAttribute(timeField) && hasTimeIntervalType(size) =>
builder.tumble(toDuration(size), timeZoneOffset).withProcessingTime()

case TumblingGroupWindow(_, timeField, size) //Tumble ROWTIME模式,内部设置Assiger
if isRowtimeAttribute(timeField) && hasTimeIntervalType(size) =>
builder.tumble(toDuration(size), timeZoneOffset).withEventTime(timeIdx)

case SlidingGroupWindow(_, timeField, size, slide) //HOP PROCTIME模式,内部设置Assiger
if isProctimeAttribute(timeField) && hasTimeIntervalType(size) =>
builder.sliding(toDuration(size), toDuration(slide), timeZoneOffset)
.withProcessingTime()
.....
case SessionGroupWindow(_, timeField, gap)
if isRowtimeAttribute(timeField) =>
builder.session(toDuration(gap)).withEventTime(timeIdx)
}

// Retraction和Trigger设置
//默认是no retract和EventTime.afterEndOfWindow
if (emitStrategy.produceUpdates) {
// mark this operator will send retraction and set new trigger
newBuilder
.withSendRetraction()
.triggering(emitStrategy.getTrigger)
}

newBuilder
.aggregate(aggsHandler, recordEqualiser, accTypes, aggValueTypes, windowPropertyTypes)
.withAllowedLateness(Duration.ofMillis(emitStrategy.getAllowLateness))
.build()
}
}

2、WindowOperator#processElement()处理数据,注册Timer

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
class StreamExecGroupWindowAggregate{
private def createWindowOperator(
config: TableConfig,
aggsHandler: GeneratedNamespaceAggsHandleFunction[_],
recordEqualiser: GeneratedRecordEqualiser,
accTypes: Array[LogicalType],
windowPropertyTypes: Array[LogicalType],
aggValueTypes: Array[LogicalType],
inputFields: Seq[LogicalType],
timeIdx: Int): WindowOperator[_, _] = {

val builder = WindowOperatorBuilder
.builder()
.withInputFields(inputFields.toArray)
val timeZoneOffset = -config.getTimeZone.getOffset(Calendar.ZONE_OFFSET)

// 设置WindowOperatorBuilder,最后通过Builder创建WindowOperator
val newBuilder = window match {
case TumblingGroupWindow(_, timeField, size) //Tumble PROCTIME模式,内部设置Assiger
if isProctimeAttribute(timeField) && hasTimeIntervalType(size) =>
builder.tumble(toDuration(size), timeZoneOffset).withProcessingTime()

case TumblingGroupWindow(_, timeField, size) //Tumble ROWTIME模式,内部设置Assiger
if isRowtimeAttribute(timeField) && hasTimeIntervalType(size) =>
builder.tumble(toDuration(size), timeZoneOffset).withEventTime(timeIdx)

case SlidingGroupWindow(_, timeField, size, slide) //HOP PROCTIME模式,内部设置Assiger
if isProctimeAttribute(timeField) && hasTimeIntervalType(size) =>
builder.sliding(toDuration(size), toDuration(slide), timeZoneOffset)
.withProcessingTime()
.....
case SessionGroupWindow(_, timeField, gap)
if isRowtimeAttribute(timeField) =>
builder.session(toDuration(gap)).withEventTime(timeIdx)
}

// Retraction和Trigger设置
//默认是no retract和EventTime.afterEndOfWindow
if (emitStrategy.produceUpdates) {
// mark this operator will send retraction and set new trigger
newBuilder
.withSendRetraction()
.triggering(emitStrategy.getTrigger)
}

newBuilder
.aggregate(aggsHandler, recordEqualiser, accTypes, aggValueTypes, windowPropertyTypes)
.withAllowedLateness(Duration.ofMillis(emitStrategy.getAllowLateness))
.build()
}
}

运行数据:

图片

3、Timer触发 I、InternalTimerServiceImpl#advanceWatermark()

WindowOperator#onEventTime()的调用前,可以先看其上层调用:InternalTimerServiceImpl#advanceWatermark()

图片

当获取的watermark为9999L时,把eventTimeTimerQueue队列中所有小于这个值的timer poll出来,调用WindowOperator.onEnventTime(timer)

II、WindwOperator#onEventTime()

WindwOperator#onEventTime()方法比较清晰,主要是window的触发和window的清理两段逻辑:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
public class WindowOperator{
publicvoidonEventTime(InternalTimer<K, W> timer) throws Exception {
setCurrentKey(timer.getKey());

triggerContext.window = timer.getNamespace();
if (triggerContext.onEventTime(timer.getTimestamp())) {
// fire
emitWindowResult(triggerContext.window);
}

if (windowAssigner.isEventTime()) {
windowFunction.cleanWindowIfNeeded(triggerContext.window, timer.getTimestamp());
}
}
}

III、emitWindowResult()提交结果

#emitWindowResult()重点关注下其第一行代码:BaseRow aggResult = windowFunction.getWindowAggregationResult(window); 这个表示根据具体的TimeWindow{start=4000, end=24000},去获取聚合数据,如果是滑动窗口,需要将4000, 8000 ,12000,16000 , 20000, 24000这几段affect窗口里面的聚合值合并起来,内部逻辑:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
public classPanedWindowProcessFunction{
public BaseRow getWindowAggregationResult(W window) throws Exception {
Iterable<W> panes = windowAssigner.splitIntoPanes(window);
BaseRow acc = windowAggregator.createAccumulators();
// null namespace means use heap data views
windowAggregator.setAccumulators(null, acc);
for (W pane : panes) {
BaseRow paneAcc = ctx.getWindowAccumulators(pane);
if (paneAcc != null) {
windowAggregator.merge(pane, paneAcc);
}
}
return windowAggregator.getValue(window);
}
}

图片

Emit(Trigger)触发器

  • 配置方式指定Trigger:Flink1.9.0目前支持通过TableConifg配置earlyFireInterval、lateFireInterval毫秒数,来指定窗口结束之前、窗口结束之后的触发策略(默认是watermark超过窗口结束后触发一次),策略的解析在WindowEmitStrategy,在StreamExecGroupWindowAggregateRule就会创建和解析这个策略
  • SQL方式指定Trigger:Flink1.9.0代码中calcite部分已有SqlEmit相关的实现,后续可以支持SQL 语句(INSERT INTO)中配置EMIT触发器

本文Emit和Trigger都是触发器这一个概念,只是使用的方式不一样

1、Emit策略 Emit 策略是指在Flink SQL 中,query的输出策略(如能忍受的延迟)可能在不同的场景有不同的需求,而这部分需求,传统的 ANSI SQL 并没有对应的语法支持。比如用户需求:1小时的时间窗口,窗口触发之前希望每分钟都能看到最新的结果,窗口触发之后希望不丢失迟到一天内的数据。针对这类需求,抽象出了EMIT语法,并扩展到了SQL语法。

2、用途 EMIT语法的用途目前总结起来主要提供了:控制延迟、数据精确性,两方面的功能。

  • 控制延迟。针对大窗口,设置窗口触发之前的EMIT输出频率,减少用户看到结果的延迟(WITH| WITHOUT DELAY)。
  • 数据精确性。不丢弃窗口触发之后的迟到的数据,修正输出结果(minIdleStateRetentionTime,在WindowEmitStrategy中生成allowLateness)。

在选择EMIT策略时,还需要与处理开销进行权衡。因为越低的输出延迟、越高的数据精确性,都会带来越高的计算开销。

3、语法 EMIT 语法是用来定义输出的策略,即是定义在输出(INSERT INTO)上的动作。当未配置时,保持原有默认行为,即 window 只在 watermark 触发时 EMIT 一个结果。

语法:INSERT INTO tableName query EMIT strategy [, strategy]*

strategy ::= {WITH DELAY timeInterval | WITHOUT DELAY} [BEFORE WATERMARK |AFTER WATERMARK]

timeInterval ::=‘string’ timeUnit

WITH DELAY:声明能忍受的结果延迟,即按指定 interval 进行间隔输出。WITHOUT DELAY:声明不忍受延迟,即每来一条数据就进行输出。BEFORE WATERMARK:窗口结束之前的策略配置,即watermark 触发之前。AFTER WATERMARK:窗口结束之后的策略配置,即watermark 触发之后。注:

  • 其中 strategy可以定义多个,同时定义before和after的策略。但不能同时定义两个 before 或 两个after 的策略。
  • 若配置了AFTER WATERMARK 策略,需要显式地在TableConfig中配置minIdleStateRetentionTime标识能忍受的最大迟到时间。
  • minIdleStateRetentionTime在window中只影响窗口何时清除,不直接影响窗口何时触发, 例如配置为3600000,最多容忍1小时的迟到数据,超过这个时间的数据会直接丢弃

4、示例 如果我们已经有一个TUMBLE(ctime, INTERVAL ‘1’ HOUR)的窗口,tumble_window 的输出是需要等到一小时结束才能看到结果,我们希望能尽早能看到窗口的结果(即使是不完整的结果)。例如,我们希望每分钟看到最新的窗口结果:INSERT INTO result SELECT * FROM tumble_window EMIT WITH DELAY ‘1’ MINUTE BEFORE WATERMARK – 窗口结束之前,每隔1分钟输出一次更新结果

tumble_window 会忽略并丢弃窗口结束后到达的数据,而这部分数据对我们来说很重要,希望能统计进最终的结果里。而且我们知道我们的迟到数据不会太多,且迟到时间不会超过一天以上,并且希望收到迟到的数据立刻就更新结果:INSERT INTO result SELECT * FROM tumble_window EMIT WITH DELAY ‘1’ MINUTE BEFORE WATERMARK, WITHOUT DELAY AFTER WATERMARK –窗口结束之后,每条到达的数据都输出

tEnv.getConfig.setIdleStateRetentionTime(Time.days(1), Time.days(2))//min、max,只有Time.days(1)这个参数直接对window生效

补充一下WITH DELAY ‘1’这种配置的周期触发策略(即DELAY大于0),最后都是由ProcessingTime系统时间触发:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
class WindowEmitStrategy{
private def createTriggerFromInterval(
enableDelayEmit: Boolean,
interval: Long): Option[Trigger[TimeWindow]] = {
if (!enableDelayEmit) {
None
} else {
if (interval > 0) {
// 系统时间触发,小于wm的所有timer都执行onProcessingTime()
Some(ProcessingTimeTriggers.every(Duration.ofMillis(interval)))
} else {
// 为0则每条都触发
Some(ElementTriggers.every())
}
}
}
}

5、Trigger类和结构关系 在源码中,Window Trigger的实现子类有10个左右,需要结合上一个小节的EMIT SQL能更容易理清他们之间的关系,这里简单介绍下:

图片

  • AfterEndOfWindow:这个就是没配置任何EMIT策略时,默认的EvenTime、ProcTime

  • Window触发策略(即窗口结束后触发一次)

  • EveryElement:即delay=0,在processElement()时直接触发,无论是在窗口结束之前或者窗口结束之后都触发,且不再注册timer

  • AfterEndOfWindowNoLate:对应EMIT WITHOUT DELAY AFTER WATERMARK,窗口结束之前不输出,窗口结束之后无延迟输出

  • AfterFirstElementPeriodic:对应WITH DELAY ‘1’ MINUTE BEFORE| AFTER WATERMARK,即按系统时间周期执行,由ProcessingTime系统时间周期触发

Sink的三种模式

Flink table的三种sink模式

Sink有INSERT、UPDATE 和 DELETE 三类,Table的sink模式有append、upsert和retract三种

Sink模式 Insert Update Delete 支持的存储
Append(追加) 支持 不支持 不支持
Upsert(重复时更新) 支持 支持 支持 KV: HBase JDBC
Retract(允许撤销) 支持 支持 支持

Upsert模式与Retract模式的区别

Upsert模式需要唯一的key来传递更新消息,外部连接器需要明确知道这个唯一key的属性

Upsert模式和Retract模式

消息: 都为(Boolean,Row)二元组, 第一个元素代表操作类型:

模式 操作类型 插入 更新 删除
Upsert模式 true 为 UPSERT消息(不存在则INSERT, 存在则UPDATE)
false 为 DELETE消息
upsert消息 upsert消息 delete消息
Retract模式 true为添加消息
false为撤回消息
添加消息 已更新行(上一行)的撤回消息
更新行(新行)的添加消息
撤回消息

Append模式-窗口聚合中的应用

在实时聚合统计中,聚合统计的结果输出是由 Trigger 决定的,而 Append-Only 则意味着对于每个窗口实例(Pane,窗格)Trigger 只能触发一次,则就导致无法在迟到数据到达时再刷新结果。

通常来说,我们可以给 Watermark 设置一个较大的延迟容忍阈值来避免这种刷新(再有迟到数据则丢弃),但代价是却会引入较大的延迟。

Upsert模式

支持 Append-Only 的操作和在有主键的前提下的 Update 和 Delete 操作.

重复时更新

Upsert 模式依赖业务主键来实现输出结果的更新和删除,因此非常适合 KV 数据库,比如 HBase、JDBC 的 TableSink 都使用了这种方式。

Upsert 模式是目前来说比较实用的模式,因为大部分业务都会提供原子或复合类型的主键,而在支持 KV 的存储系统也非常多,但要注意的是不要变更主键,具体原因会在下一节谈到。

Retract模式

允许撤销

举个例子,假设我们将电商订单按照承运快递公司进行分类计数,有如下的结果表。

公司 订单数
中通 2
圆通 1
顺丰 3

那么如果原本一单为中通的快递,后续更新为用顺丰发货,对于 Upsert 模式会产生 (true, (顺丰, 4)) 这样一条 changelog,但中通的订单数没有被修正。相比之下,Retract 模式产出 (false, (中通, 1))(true, (顺丰, 1)) 两条数据,则可以正确地更新数据。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
select /*+ MAPJOIN(hj3,hj1,hj0)*/ 20210324 as fdate, TO_CHAR(SYSTIMESTAMP(), 'yyyymmddhh24miss') as fetl_time,hj2.fuin as fuin,COUNT(DISTINCT(SUBSTR(CAST(if (((hj2.fevent_time between 20210318000000 and 20210324235959) and (hj2.fscn in ('1'))) , hj2.fevent_time, null) AS STRING),1,8))) as flabel_25942,   COUNT(DISTINCT(SUBSTR(CAST(if (((hj2.fevent_time between 20210318000000 and 20210324235959)  and (hj2.fscn in ('1'))) , hj2.fevent_time, null) AS STRING),1,8))) as flabel_25970, COUNT(DISTINCT(SUBSTR(CAST(if (((hj2.furl_orig = '/mb/v5/fund/list/steady.shtml')  or (hj2.furl_orig = '/mb/v4/fundlist/fund_all.shtml')    or (hj2.furl_orig = '/mb/v4/fundlist/fund_all_v5.shtml'))  and ((hj2.fevent_time between 20210318000000 and 20210324235959)) , hj2.fevent_time, null) AS STRING),1,8))) as flabel_25972

from dml_base::dml_evt_lct_cft_label_factory_mta_access_dd hj2

left join dim_base::dim_lct_hq_beacon_report_manage hj1 on if(hj2.furl_orig is null
or trim(hj2.furl_orig) = '', concat('null_', floor(rand() * 10000)), hj2.furl_orig) = hj1.fevent_code

left join dim_base::dim_lct_hq_virtual_event_info hj3 on if(hj2.furl_orig is null
or trim(hj2.furl_orig) = '', concat('null_', floor(rand() * 10000)), hj2.furl_orig) = hj3.fevent_code
left join dim_base::dim_prd_lct_cft_fund_type_conf hj0 on if(hj2.fspid_fundcode is null
or trim(hj2.fspid_fundcode) = '', concat('null_', floor(rand() * 10000)), hj2.fspid_fundcode) = hj0.fspid_fundcode

where hj2.fuin is not null

and trim(hj2.fuin) <> ''

and hj2.fdate >= 20210318

and hj2.fdate <= 20210324

group by hj2.fuin
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
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
ABSTRACT SYNTAX TREE:
(TOK_QUERY (TOK_FROM (TOK_LEFTOUTERJOIN (TOK_LEFTOUTERJOIN (TOK_LEFTOUTERJOIN (TOK_TABREF (TOK_TAB dml_evt_lct_cft_label_factory_mta_access_dd dml_base) hj2) (TOK_TABREF (TOK_TAB dim_lct_hq_beacon_report_manage dim_base) hj1) (= (TOK_FUNCTION if (or (TOK_FUNCTION TOK_ISNULL (. (TOK_TABLE_OR_COL hj2) furl_orig)) (= (TOK_FUNCTION trim (. (TOK_TABLE_OR_COL hj2) furl_orig)) '')) (TOK_FUNCTION concat 'null_' (TOK_FUNCTION floor (* (TOK_FUNCTION rand) 10000))) (. (TOK_TABLE_OR_COL hj2) furl_orig)) (. (TOK_TABLE_OR_COL hj1) fevent_code))) (TOK_TABREF (TOK_TAB dim_lct_hq_virtual_event_info dim_base) hj3) (= (TOK_FUNCTION if (or (TOK_FUNCTION TOK_ISNULL (. (TOK_TABLE_OR_COL hj2) furl_orig)) (= (TOK_FUNCTION trim (. (TOK_TABLE_OR_COL hj2) furl_orig)) '')) (TOK_FUNCTION concat 'null_' (TOK_FUNCTION floor (* (TOK_FUNCTION rand) 10000))) (. (TOK_TABLE_OR_COL hj2) furl_orig)) (. (TOK_TABLE_OR_COL hj3) fevent_code))) (TOK_TABREF (TOK_TAB dim_prd_lct_cft_fund_type_conf dim_base) hj0) (= (TOK_FUNCTION if (or (TOK_FUNCTION TOK_ISNULL (. (TOK_TABLE_OR_COL hj2) fspid_fundcode)) (= (TOK_FUNCTION trim (. (TOK_TABLE_OR_COL hj2) fspid_fundcode)) '')) (TOK_FUNCTION concat 'null_' (TOK_FUNCTION floor (* (TOK_FUNCTION rand) 10000))) (. (TOK_TABLE_OR_COL hj2) fspid_fundcode)) (. (TOK_TABLE_OR_COL hj0) fspid_fundcode)))) (TOK_INSERT (TOK_DESTINATION (TOK_DIR TOK_TMP_FILE)) (TOK_SELECT (TOK_HINTLIST (TOK_HINT TOK_MAPJOIN (TOK_HINTARGLIST (TOK_TABLE_OR_COL hj3) (TOK_TABLE_OR_COL hj1) (TOK_TABLE_OR_COL hj0)))) (TOK_SELEXPR 20210324 fdate) (TOK_SELEXPR (TOK_FUNCTION TO_CHAR (TOK_FUNCTION SYSTIMESTAMP) 'yyyymmddhh24miss') fetl_time) (TOK_SELEXPR (. (TOK_TABLE_OR_COL hj2) fuin) fuin) (TOK_SELEXPR (TOK_FUNCTIONDI COUNT (TOK_FUNCTION SUBSTR (TOK_FUNCTION TOK_STRING (TOK_FUNCTION if (and (and (>= (. (TOK_TABLE_OR_COL hj2) fevent_time) 20210318000000) (<= (. (TOK_TABLE_OR_COL hj2) fevent_time) 20210324235959)) (in (. (TOK_TABLE_OR_COL hj2) fscn) '1')) (. (TOK_TABLE_OR_COL hj2) fevent_time) TOK_NULL)) 1 8)) flabel_25942) (TOK_SELEXPR (TOK_FUNCTIONDI COUNT (TOK_FUNCTION SUBSTR (TOK_FUNCTION TOK_STRING (TOK_FUNCTION if (and (and (>= (. (TOK_TABLE_OR_COL hj2) fevent_time) 20210318000000) (<= (. (TOK_TABLE_OR_COL hj2) fevent_time) 20210324235959)) (in (. (TOK_TABLE_OR_COL hj2) fscn) '1')) (. (TOK_TABLE_OR_COL hj2) fevent_time) TOK_NULL)) 1 8)) flabel_25970) (TOK_SELEXPR (TOK_FUNCTIONDI COUNT (TOK_FUNCTION SUBSTR (TOK_FUNCTION TOK_STRING (TOK_FUNCTION if (and (or (or (= (. (TOK_TABLE_OR_COL hj2) furl_orig) '/mb/v5/fund/list/steady.shtml') (= (. (TOK_TABLE_OR_COL hj2) furl_orig) '/mb/v4/fundlist/fund_all.shtml')) (= (. (TOK_TABLE_OR_COL hj2) furl_orig) '/mb/v4/fundlist/fund_all_v5.shtml')) (and (>= (. (TOK_TABLE_OR_COL hj2) fevent_time) 20210318000000) (<= (. (TOK_TABLE_OR_COL hj2) fevent_time) 20210324235959))) (. (TOK_TABLE_OR_COL hj2) fevent_time) TOK_NULL)) 1 8)) flabel_25972)) (TOK_WHERE (and (and (and (TOK_FUNCTION TOK_ISNOTNULL (. (TOK_TABLE_OR_COL hj2) fuin)) (<> (TOK_FUNCTION trim (. (TOK_TABLE_OR_COL hj2) fuin)) '')) (>= (. (TOK_TABLE_OR_COL hj2) fdate) 20210318)) (<= (. (TOK_TABLE_OR_COL hj2) fdate) 20210324))) (TOK_GROUPBY (. (TOK_TABLE_OR_COL hj2) fuin))))

STAGE DEPENDENCIES:
Stage-1
type:root stage;
Stage-2
type:;depends on:Stage-1;
Stage-3
type:;depends on:Stage-2;
Stage-0
type:root stage;

STAGE PLANS:
Stage: Stage-1
Map Reduce
Alias -> Map Operator Tree:
dml_base/dml_evt_lct_cft_label_factory_mta_access_dd#hj2
Operator: TableScan
alias: dml_base/dml_evt_lct_cft_label_factory_mta_access_dd#hj2
Operator: Filter Operator
predicate:
expr: (((fuin is not null and (trim(fuin) <> '')) and (fdate >= 20210318)) and (fdate <= 20210324))
type: boolean
Operator: Common Join Operator
condition map:
Left Outer Join0 to 1
Left Outer Join0 to 2
condition expressions:
0 {fdate} {fuin} {furl_orig} {fevent_time} {fspid_fundcode} {fscn}
1
2
handleSkewJoin: false
keys:
0 [class com.tencent.tdw_udf_cloud.hive.udf.generic.GenericUDFIf(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFOPOr(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFOPNull(Column[furl_orig](), class org.apache.hadoop.hive.ql.udf.generic.GenericUDFOPEqual(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(Column[furl_orig](), Const string ()(), class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(Const string null_, class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge((), Const int 10000()()(), Column[furl_orig]()]
1 [Column[fevent_code]]
2 [Column[fevent_code]]
outputColumnNames: _col0, _col42, _col85, _col88, _col103, _col104
Position of Big Table: 0
Operator: File Output Operator
compressed: false
GlobalTableId: 0
table:
table descs
input format: org.apache.hadoop.mapred.SequenceFileInputFormat
output format: org.apache.hadoop.hive.ql.io.HiveSequenceFileOutputFormat
Local Work:
Map Reduce Local Work
Alias -> Map Local Tables:
dim_base/dim_lct_hq_beacon_report_manage#hj1
Fetch Operator
limit: -1
dim_base/dim_lct_hq_virtual_event_info#hj3
Fetch Operator
limit: -1
Alias -> Map Local Operator Tree:
dim_base/dim_lct_hq_beacon_report_manage#hj1
Operator: TableScan
alias: dim_base/dim_lct_hq_beacon_report_manage#hj1
Operator: Common Join Operator
condition map:
Left Outer Join0 to 1
Left Outer Join0 to 2
condition expressions:
0 {fdate} {fuin} {furl_orig} {fevent_time} {fspid_fundcode} {fscn}
1
2
handleSkewJoin: false
keys:
0 [class com.tencent.tdw_udf_cloud.hive.udf.generic.GenericUDFIf(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFOPOr(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFOPNull(Column[furl_orig](), class org.apache.hadoop.hive.ql.udf.generic.GenericUDFOPEqual(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(Column[furl_orig](), Const string ()(), class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(Const string null_, class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge((), Const int 10000()()(), Column[furl_orig]()]
1 [Column[fevent_code]]
2 [Column[fevent_code]]
outputColumnNames: _col0, _col42, _col85, _col88, _col103, _col104
Position of Big Table: 0
Operator: File Output Operator
compressed: false
GlobalTableId: 0
table:
table descs
input format: org.apache.hadoop.mapred.SequenceFileInputFormat
output format: org.apache.hadoop.hive.ql.io.HiveSequenceFileOutputFormat
dim_base/dim_lct_hq_virtual_event_info#hj3
Operator: TableScan
alias: dim_base/dim_lct_hq_virtual_event_info#hj3
Operator: Common Join Operator
condition map:
Left Outer Join0 to 1
Left Outer Join0 to 2
condition expressions:
0 {fdate} {fuin} {furl_orig} {fevent_time} {fspid_fundcode} {fscn}
1
2
handleSkewJoin: false
keys:
0 [class com.tencent.tdw_udf_cloud.hive.udf.generic.GenericUDFIf(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFOPOr(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFOPNull(Column[furl_orig](), class org.apache.hadoop.hive.ql.udf.generic.GenericUDFOPEqual(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(Column[furl_orig](), Const string ()(), class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(Const string null_, class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge((), Const int 10000()()(), Column[furl_orig]()]
1 [Column[fevent_code]]
2 [Column[fevent_code]]
outputColumnNames: _col0, _col42, _col85, _col88, _col103, _col104
Position of Big Table: 0
Operator: File Output Operator
compressed: false
GlobalTableId: 0
table:
table descs
input format: org.apache.hadoop.mapred.SequenceFileInputFormat
output format: org.apache.hadoop.hive.ql.io.HiveSequenceFileOutputFormat
Path -> Alias:
hdfs://<hive_hdfs_base_path>/warehouse/dml_base.db/dml_evt_lct_cft_label_factory_mta_access_dd/par_20210318 [dml_base/dml_evt_lct_cft_label_factory_mta_access_dd#hj2]
hdfs://<hive_hdfs_base_path>/warehouse/dml_base.db/dml_evt_lct_cft_label_factory_mta_access_dd/par_20210319 [dml_base/dml_evt_lct_cft_label_factory_mta_access_dd#hj2]
hdfs://<hive_hdfs_base_path>/warehouse/dml_base.db/dml_evt_lct_cft_label_factory_mta_access_dd/par_20210320 [dml_base/dml_evt_lct_cft_label_factory_mta_access_dd#hj2]
hdfs://<hive_hdfs_base_path>/warehouse/dml_base.db/dml_evt_lct_cft_label_factory_mta_access_dd/par_20210321 [dml_base/dml_evt_lct_cft_label_factory_mta_access_dd#hj2]
hdfs://<hive_hdfs_base_path>/warehouse/dml_base.db/dml_evt_lct_cft_label_factory_mta_access_dd/par_20210322 [dml_base/dml_evt_lct_cft_label_factory_mta_access_dd#hj2]
hdfs://<hive_hdfs_base_path>/warehouse/dml_base.db/dml_evt_lct_cft_label_factory_mta_access_dd/par_20210323 [dml_base/dml_evt_lct_cft_label_factory_mta_access_dd#hj2]
hdfs://<hive_hdfs_base_path>/warehouse/dml_base.db/dml_evt_lct_cft_label_factory_mta_access_dd/par_20210324 [dml_base/dml_evt_lct_cft_label_factory_mta_access_dd#hj2]

Stage: Stage-2
Map Reduce
Alias -> Map Operator Tree:
hdfs://<tmp_hive_hdfs_base_path>/20210325/<hdfs_user>_tdwadmin_20210325195415346_29895878_7_699597773/10002
Operator: Select Operator
expressions:
expr: _col0
type: bigint
expr: _col42
type: string
expr: _col85
type: string
expr: _col88
type: bigint
expr: _col103
type: string
expr: _col104
type: string
outputColumnNames: _col0, _col42, _col85, _col88, _col103, _col104
Operator: Common Join Operator
condition map:
Left Outer Join0 to 1
condition expressions:
0 {_col0} {_col42} {_col85} {_col88} {_col104}
1
handleSkewJoin: false
keys:
0 [class com.tencent.tdw_udf_cloud.hive.udf.generic.GenericUDFIf(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFOPOr(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFOPNull(Column[_col103](), class org.apache.hadoop.hive.ql.udf.generic.GenericUDFOPEqual(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(Column[_col103](), Const string ()(), class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(Const string null_, class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge((), Const int 10000()()(), Column[_col103]()]
1 [Column[fspid_fundcode]]
outputColumnNames: _col0, _col42, _col85, _col88, _col104
Position of Big Table: 0
Operator: File Output Operator
compressed: false
GlobalTableId: 0
table:
table descs
input format: org.apache.hadoop.mapred.SequenceFileInputFormat
output format: org.apache.hadoop.hive.ql.io.HiveSequenceFileOutputFormat
Local Work:
Map Reduce Local Work
Alias -> Map Local Tables:
dim_base/dim_prd_lct_cft_fund_type_conf#hj0
Fetch Operator
limit: -1
Alias -> Map Local Operator Tree:
dim_base/dim_prd_lct_cft_fund_type_conf#hj0
Operator: TableScan
alias: dim_base/dim_prd_lct_cft_fund_type_conf#hj0
Operator: Common Join Operator
condition map:
Left Outer Join0 to 1
condition expressions:
0 {_col0} {_col42} {_col85} {_col88} {_col104}
1
handleSkewJoin: false
keys:
0 [class com.tencent.tdw_udf_cloud.hive.udf.generic.GenericUDFIf(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFOPOr(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFOPNull(Column[_col103](), class org.apache.hadoop.hive.ql.udf.generic.GenericUDFOPEqual(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(Column[_col103](), Const string ()(), class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(Const string null_, class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge((), Const int 10000()()(), Column[_col103]()]
1 [Column[fspid_fundcode]]
outputColumnNames: _col0, _col42, _col85, _col88, _col104
Position of Big Table: 0
Operator: File Output Operator
compressed: false
GlobalTableId: 0
table:
table descs
input format: org.apache.hadoop.mapred.SequenceFileInputFormat
output format: org.apache.hadoop.hive.ql.io.HiveSequenceFileOutputFormat
Path -> Alias:
hdfs://<tmp_hive_hdfs_base_path>/20210325/<hdfs_user>_tdwadmin_20210325195415346_29895878_7_699597773/10002 [hdfs://<tmp_hive_hdfs_base_path>/20210325/<hdfs_user>_tdwadmin_20210325195415346_29895878_7_699597773/10002]

Stage: Stage-3
Map Reduce
Alias -> Map Operator Tree:
hdfs://<tmp_hive_hdfs_base_path>/20210325/<hdfs_user>_tdwadmin_20210325195415346_29895878_7_699597773/10002
Operator: Select Operator
expressions:
expr: _col0
type: bigint
expr: _col42
type: string
expr: _col85
type: string
expr: _col88
type: bigint
expr: _col103
type: string
expr: _col104
type: string
outputColumnNames: _col0, _col42, _col85, _col88, _col103, _col104
Operator: Common Join Operator
condition map:
Left Outer Join0 to 1
condition expressions:
0 {_col0} {_col42} {_col85} {_col88} {_col104}
1
handleSkewJoin: false
keys:
0 [class com.tencent.tdw_udf_cloud.hive.udf.generic.GenericUDFIf(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFOPOr(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFOPNull(Column[_col103](), class org.apache.hadoop.hive.ql.udf.generic.GenericUDFOPEqual(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(Column[_col103](), Const string ()(), class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(Const string null_, class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge(class org.apache.hadoop.hive.ql.udf.generic.GenericUDFBridge((), Const int 10000()()(), Column[_col103]()]
1 [Column[fspid_fundcode]]
outputColumnNames: _col0, _col42, _col85, _col88, _col104
Position of Big Table: 0
Operator: File Output Operator
compressed: false
GlobalTableId: 0
table:
table descs
input format: org.apache.hadoop.mapred.SequenceFileInputFormat
output format: org.apache.hadoop.hive.ql.io.HiveSequenceFileOutputFormat
Path -> Alias:
hdfs://<tmp_hive_hdfs_base_path>/20210325/<hdfs_user>_tdwadmin_20210325195415346_29895878_7_699597773/10002 [hdfs://<tmp_hive_hdfs_base_path>/20210325/<hdfs_user>_tdwadmin_20210325195415346_29895878_7_699597773/10002]
Reduce Operator Tree:
Operator: Group By Operator
aggregations:
expr: count(DISTINCT KEY._col1:1._col0)
expr: count(DISTINCT KEY._col1:2._col0)
keys:
expr: KEY._col0
type: string
mode: mergepartial
outputColumnNames: _col0, _col1, _col2
UseNewGroupBy: true
Operator: Select Operator
expressions:
expr: 20210324
type: int
expr: (systimestamp to_char 'yyyymmddhh24miss')
type: string
expr: _col0
type: string
expr: _col1
type: bigint
expr: _col1
type: bigint
expr: _col2
type: bigint
outputColumnNames: _col0, _col1, _col2, _col3, _col4, _col5
Operator: File Output Operator
compressed: false
GlobalTableId: 0
table:
table descs
input format: org.apache.hadoop.mapred.TextInputFormat
output format: org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat

Stage: Stage-0
Fetch Operator
limit: 100000000

error

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
submitSql has error: 
org.apache.spark.sql.AnalysisException: nondeterministic expressions are only allowed in
Project, Filter, Aggregate or Window, found:
((IF(((hj2.`furl_orig` IS NULL) OR (com.tencent.tdw_udf_cloud.hive.udf.UDFTrim(hj2.`furl_orig`) = '')), com.tencent.tdw_udf_cloud.hive.udf.UDFConcat('null_', com.tencent.tdw_udf_cloud.hive.udf.UDFFloor((org.apache.hadoop.hive.ql.udf.UDFRand() * CAST(10000 AS DOUBLE)))), hj2.`furl_orig`)) = hj1.`fevent_code`)
in operator Join LeftOuter, (if ((isnull(furl_orig#92) || (HiveSimpleUDF#com.tencent.tdw_udf_cloud.hive.udf.UDFTrim(furl_orig#92) = ))) HiveSimpleUDF#com.tencent.tdw_udf_cloud.hive.udf.UDFConcat(null_,HiveSimpleUDF#com.tencent.tdw_udf_cloud.hive.udf.UDFFloor((HiveSimpleUDF#org.apache.hadoop.hive.ql.udf.UDFRand() * cast(10000 as double)))) else furl_orig#92 = fevent_code#132)
;;
Aggregate [fuin#49], [20210324 AS fdate#0, HiveSimpleUDF#com.tencent.tdw_udf_cloud.hive.udf.UDFToChar(HiveSimpleUDF#org.apache.hadoop.hive.ql.udf.UDFSysTimestamp(),yyyymmddhh24miss) AS fetl_time#1, fuin#49 AS fuin#2, count(distinct HiveSimpleUDF#com.tencent.tdw_udf_cloud.hive.udf.UDFSubstr(cast(if ((((fevent_time#95L >= 20210318000000) && (fevent_time#95L <= 20210324235959)) && fscn#111 IN (1))) fevent_time#95L else cast(null as bigint) as string),1,8)) AS flabel_25942#3L, count(distinct HiveSimpleUDF#com.tencent.tdw_udf_cloud.hive.udf.UDFSubstr(cast(if ((((fevent_time#95L >= 20210318000000) && (fevent_time#95L <= 20210324235959)) && fscn#111 IN (1))) fevent_time#95L else cast(null as bigint) as string),1,8)) AS flabel_25970#4L, count(distinct HiveSimpleUDF#com.tencent.tdw_udf_cloud.hive.udf.UDFSubstr(cast(if (((((furl_orig#92 = /mb/v5/fund/list/steady.shtml) || (furl_orig#92 = /mb/v4/fundlist/fund_all.shtml)) || (furl_orig#92 = /mb/v4/fundlist/fund_all_v5.shtml)) && ((fevent_time#95L >= 20210318000000) && (fevent_time#95L <= 20210324235959)))) fevent_time#95L else cast(null as bigint) as string),1,8)) AS flabel_25972#5L]
+- Filter ((isnotnull(fuin#49) && NOT (HiveSimpleUDF#com.tencent.tdw_udf_cloud.hive.udf.UDFTrim(fuin#49) = )) && ((fdate#7L >= cast(20210318 as bigint)) && (fdate#7L <= cast(20210324 as bigint))))
+- Join LeftOuter, (if ((isnull(fspid_fundcode#110) || (HiveSimpleUDF#com.tencent.tdw_udf_cloud.hive.udf.UDFTrim(fspid_fundcode#110) = ))) HiveSimpleUDF#com.tencent.tdw_udf_cloud.hive.udf.UDFConcat(null_,HiveSimpleUDF#com.tencent.tdw_udf_cloud.hive.udf.UDFFloor((HiveSimpleUDF#org.apache.hadoop.hive.ql.udf.UDFRand() * cast(10000 as double)))) else fspid_fundcode#110 = fspid_fundcode#209)
:- Join LeftOuter, (if ((isnull(furl_orig#92) || (HiveSimpleUDF#com.tencent.tdw_udf_cloud.hive.udf.UDFTrim(furl_orig#92) = ))) HiveSimpleUDF#com.tencent.tdw_udf_cloud.hive.udf.UDFConcat(null_,HiveSimpleUDF#com.tencent.tdw_udf_cloud.hive.udf.UDFFloor((HiveSimpleUDF#org.apache.hadoop.hive.ql.udf.UDFRand() * cast(10000 as double)))) else furl_orig#92 = fevent_code#201)
: :- Join LeftOuter, (if ((isnull(furl_orig#92) || (HiveSimpleUDF#com.tencent.tdw_udf_cloud.hive.udf.UDFTrim(furl_orig#92) = ))) HiveSimpleUDF#com.tencent.tdw_udf_cloud.hive.udf.UDFConcat(null_,HiveSimpleUDF#com.tencent.tdw_udf_cloud.hive.udf.UDFFloor((HiveSimpleUDF#org.apache.hadoop.hive.ql.udf.UDFRand() * cast(10000 as double)))) else furl_orig#92 = fevent_code#132)
: : :- SubqueryAlias hj2
: : : +- MetastoreRelation dml_base, dml_evt_lct_cft_label_factory_mta_access_dd
: : +- BroadcastHint
: : +- SubqueryAlias hj1
: : +- MetastoreRelation dim_base, dim_lct_hq_beacon_report_manage
: +- BroadcastHint
: +- SubqueryAlias hj3
: +- MetastoreRelation dim_base, dim_lct_hq_virtual_event_info
+- BroadcastHint
+- SubqueryAlias hj0
+- MetastoreRelation dim_base, dim_prd_lct_cft_fund_type_conf
1
java.lang.RuntimeException: Map local work failed at org.apache.hadoop.hive.ql.exec.ExecMapper.processOldMapLocalWork(ExecMapper.java:317) at org.apache.hadoop.hive.ql.exec.ExecMapper.map(ExecMapper.java:151) at org.apache.hadoop.mapred.MapRunner.run(MapRunner.java:54) at org.apache.hadoop.mapred.MapTask.runOldMapper(MapTask.java:453) at org.apache.hadoop.mapred.MapTask.run(MapTask.java:343) at org.apache.hadoop.mapred.YarnChild$2.run(YarnChild.java:175) at java.security.AccessController.doPrivileged(Native Method) at javax.security.auth.Subject.doAs(Subject.java:422) at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:2286) at org.apache.hadoop.mapred.YarnChild.main(YarnChild.java:169) Caused by: org.apache.hadoop.hive.ql.metadata.HiveException: java.io.IOException: 没有那个文件或目录 at org.apache.hadoop.hive.ql.exec.persistence.HashMapWrapper.getPersistentHash(HashMapWrapper.java:189) at org.apache.hadoop.hive.ql.exec.persistence.HashMapWrapper.put(HashMapWrapper.java:155) at org.apache.hadoop.hive.ql.exec.MapJoinOperator.process(MapJoinOperator.java:474) at org.apache.hadoop.hive.ql.exec.Operator.forward(Operator.java:471) at org.apache.hadoop.hive.ql.exec.TableScanOperator.process(TableScanOperator.java:37) at org.apache.hadoop.hive.ql.exec.ExecMapper.processOldMapLocalWork(ExecMapper.java:302) ... 9 more Caused by: java.io.IOException: 没有那个文件或目录 at java.io.UnixFileSystem.createFileExclusively(Native Method) at java.io.File.createTempFile(File.java:2024) at org.apache.hadoop.hive.ql.exec.persistence.HashMapWrapper.getPersistentHash(HashMapWrapper.java:176) ... 14 more	N/A

定位出错误原因是:强制指定了mapjoin,内存溢出了

Hive的自动join策略选择:

由于开启了hive.auto.convert.join,但是实际小表大小是hive.mapjoin.smalltable.filesize(默认25M,小表不会超过25M)。由于使用的是orc压缩,解压缩后可能大小到了250M,存放到内存大小可能就会超过1G。mapjoin的时候,hive orcfile 放到内存中会放大40倍
可以看到JVM Max Heap Size大小为:1013645312 (大约1G)

一句话总结

由于使用了hive.auto.convert.join,对小表进行广播,但是原表是orc的,存放到内存可能膨胀到大于localtask的堆内存大小,导致sql执行失败。

解决措施

方案一

调大localtask的内存,set hive.mapred.local.mem=XX ,默认1G,调大到4G

方案二

直接关表autojoin,将hive.auto.convert.join设置成false

flink join

Cogroup

CoGroupFunction中会返回所有数据,不管有没有匹配上

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
DataStream<Tuple3<Long, String, String>> input1 = ...;
input1 = input1.assignTimestampsAndWatermarks(new AscendingTimestampExtractor<Tuple3<Long, String, String>>() {

@Override
public long extractAscendingTimestamp(Tuple3<Long, String, String> arg0) {
return arg0.f0;
}

});

DataStream<Tuple2<Long, String>> input2 = ...;
input2 = input2.assignTimestampsAndWatermarks(new AscendingTimestampExtractor<Tuple2<Long, String>>() {
@Override
public long extractAscendingTimestamp(Tuple2<Long, String> stringStringTuple2) {
return stringStringTuple2.f0;
}
});

input1.coGroup(input2).where(new KeySelector<Tuple3<Long, String, String>, String>() {
@Override
public String getKey(Tuple3<Long, String, String> itemEntity) throws Exception {
return itemEntity.f1;
}
})
.equalTo(new KeySelector<Tuple2<Long, String>, String>() {
@Override
public String getKey(Tuple2<Long, String> value) throws Exception {
return value.f1;
}
})
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.apply(new CoGroupFunction<Tuple3<Long, String, String>, Tuple2<Long, String>, String>() {
@Override
public void coGroup(Iterable<Tuple3<Long, String, String>> first,
Iterable<Tuple2<Long, String>> second, Collector<String> collector) throws Exception {
StringBuilder buffer = new StringBuilder();
buffer.append("DataStream first:\n");
for (Tuple3<Long, String, String> value : first) {
buffer.append(value).append("\n");
}
buffer.append("DataStream second:\n");
for (Tuple2<Long, String> value : second) {
buffer.append(value.f0).append("=>").append(value.f1).append("\n");
}
collector.collect(buffer.toString());
}
})
.print();

上面的例子,左流有三个元素 Tuple3<String,String,String>,右流有两个元素Tuple2<String,String>
两个流第一个元素相互关联。分别指定两个流的事件时间字段。
两个流关联后,按照EventTime划分窗口。与单流类似。
不管元素是否可以关联上,都会输出

用户可以定义CoGroupFunction函数, 可以实现在窗口内,任意组合,如笛卡尔积

举例说明:

window Join

interval join

broadcast join

1 同步IO VS 异步IO

1.1 从操作系统的角度,异步IO为什么可以提升性能?

2 Flink-asyncIO

优点:写库不会阻塞性能更好

然后来看一下, Flink 中异步io主要分为两种

  一种是有序Ordered

  一种是无序UNordered

主要区别是往下游output的顺序(注意这里顺序不是写库的顺序既然都异步了写库的顺序自然是无法保证的),有序的会按接收的顺序继续往下游output发送,无序就是谁先处理完谁就先往下游发送

两张图了解这两种模式的实现

有序:record数据会通过异步线程写库,Emitter是一个守护进程,会不停的拉取queue头部的数据,如果头部的数据异步写库完成,Emitter将头数据往下游发送,如果头元素还没有异步写库完成,柱塞

无序:record数据会通过异步线程写库,这里有两个queue,一开始放在uncompleteedQueue,当哪个record异步写库成功后就直接放到completedQueue中,Emitter是一个守护进程,completedQueue只要有数据,会不停的拉取queue数据往下游发送

可以看到原理还是很简单的,两句话就总结完了,就是利用queue和java的异步线程,现在来看下源码

这里AsyncIO在Flink中被设计成operator中的一种,自然去OneInputStreamOperator的实现类中去找

于是来看一下AsyncWaitOperator.java

  

看到它的open方法(open方法会在taskmanager启动job的时候全部统一调用,可以翻一下以前的文章)

这里启动了一个守护线程Emitter,来看下线程具体做了什么

1处拉取数据,2处就是常规的将拉取到的数据往下游emit,Emitter拉取数据,这里先不讲因为分为有序的和无序的

这里已经知道了这个Emitter的作用是循环的拉取数据往下游发送

回到AsyncWaitOperator.java在它的open方法初始化了Emitter,那它是如何处理接收到的数据的呢,看它的ProcessElement()方法

1
![](_v_images/20201207210250637_60325951.png)

其实主要就是三个个方法

先是!!!将record封装成了一个包装类StreamRecordQueueEntry,主要是这个包装类的构造方法中,创建了一个CompleteableFuture(这个的complete方法其实会等到用户代码执行的时候用户自己决定什么时候完成)

1处主要就是讲元素加入到了对应的queue,这里也分为两种有序和无序的

这里也先不讲这两种模式加入数据的区别

接着2处就是调用用户的代码了,来看看官网的异步io的例子

给了一个Future作为参数,用户自己起了一个线程(这里思考一下就知道了为什么要新起一个异步线程去执行,因为如果不起线程的话,那processElement方法就柱塞了,无法异步了)去写库读库等,然后调用了这个参数的complete方法(也就是前面那个包装类中的CompleteableFuture)并且传入了一个结果

看下complete方法源码

这个resultFuture是每个record的包装类StreamRecordQueueEntry的其中一个属性是一个CompletableFuture

那现在就清楚了,用户代码在自己新起的线程中当自己的逻辑执行完以后会使这个异步线程结束,并输入一个结果

那这个干嘛用的呢

最开始的图中看到有序和无序实现原理,有序用一个queue,无序用两个queue分别就对应了

OrderedStreamElementQueue类中

UnorderedStreamElementQueue类中

回到前面有两个地方没有细讲,一是两种模式的Emitter是如何拉取数据的,二是两种模式下数据是如何加入OrderedStreamElementQueue的

有序模式:

1.先来看一下有序模式的,Emitter的数据拉取,和数据的加入

    其tryPut()方法

    

   

    onComplete**方法

       

1
onCompleteHandler方法

       

  这里比较绕,先将接收的数据加入queue中,然后onComplete()中当上一个异步线程getFuture() 其实就是每个元素包装类里面的那个CompletableFuture,当他结束时(会在用户方法用户调用complete时结束)异步调用传入的对象的 accept方法,accept方法中调用了onCompleteHandler()方法,onCompleteHandler方法中会判断queue是否为空,以及queue的头元素是否完成了用户的异步方法,当完成的时候,就会将headIsCompleted这个对象signalAll()唤醒

2.接着看有序模式Emitter的拉取数据

这里有序方式拉取数据的逻辑很清晰,如果为空或者头元素没有完成用户的异步方法,headIsCompleted这个对象会wait住(上面可以知道,当加入元素的到queue且头元素完成异步方法的时候会signalAll())然后将头数据返回,往下游发送

这样就实现了有序发送,因为Emitter只拉取头元素且已经完成用户异步方法的头元素

无序模式:

  这里和有序模式就大同小异了,只是变成了,接收数据后直接加入uncompletedQueue,当数据完成异步方法的时候就,放到completedQueue里面去并signalAll(),只要completedqueue里面有数据,Emitter就拉取往下发

这样就实现了无序模式,也就是异步写入谁先处理完就直接放到完成队列里面去,然后往下发,不用管接收数据的顺序

3 FlinkSQL通过异步IO优化

4 注意点

5 参考文献

  1. 阿里面试题:使用 Flink 的异步 IO 需要注意哪些细节?

10.12.0

reference

  • 在 DataStream API 上添加了高效的批执行模式的支持。这是批处理和流处理实现真正统一的运行时的一个重要里程碑。
  • 实现了基于Kubernetes的高可用性(HA)方案,作为生产环境中,ZooKeeper方案之外的另外一种选择。
  • 扩展了 Kafka SQL connector,使其可以在 upsert 模式下工作,并且支持在 SQL DDL 中处理 connector 的 metadata。现在,时态表 Join 可以完全用 SQL 来表示,不再依赖于 Table API 了。
  • PyFlink 中添加了对于 DataStream API 的支持,将 PyFlink 扩展到了更复杂的场景,比如需要对状态或者定时器 timer 进行细粒度控制的场景。除此之外,现在原生支持将 PyFlink 作业部署到 Kubernetes上。

DataStream API支持批量

可复用性:作业可以在流和批这两种执行模式之间自由地切换,而无需重写任何代码。因此,用户可以复用同一个作业,来处理实时数据和历史数据。

维护简单:统一的 API 意味着流和批可以共用同一组 connector,维护同一套代码,并能够轻松地实现流批混合执行,例如 backfilling 之类的场景。

Data Sink API

Sort-Merge Shuffle

SQL 中 支持 Temporal Table Join

Flink-temporal-table-join

Temporal Table记录了表历史上任何时间点所有的数据改动

ANSI-SQL 2011 Temporal Table示例

我们以一个DDL和一套DML示例说明Temporal Table的原理,DDL定义PK是可选的,下面的示例我们以不定义PK的为例进行说明:

  • DDL 示例
1
2
3
4
5
6
7
8
9
CREATE TABLE Emp
ENo INTEGER,
Sys_start TIMESTAMP(12) GENERATED
ALWAYS AS ROW START,
Sys_end TIMESTAMP(12) GENERATED
ALWAYS AS ROW END,
EName VARCHAR(30),
PERIOD FOR SYSTEM_TIME (Sys_start,Sys_end)
) WITH SYSTEM VERSIONING
  • DML 示例
    • INSERT
1
INSERT INTO Emp (ENo, EName) VALUES (22217, 'Joe')

说明: 其中Sys_Start和Sys_End是数据库系统默认填充的。

    • UPDATE
1
UPDATE Emp SET EName = 'Tom' WHERE ENo = 22217

说明: 假设是在 2012-02-03 10:00:00 执行的UPDATE,执行之后上一个值 "Joe" 的Sys_End值由 9999-12-31 23:59:59 变成了 2012-02-03 10:00:00 , 也就是下一个值 "Tom" 生效的开始时间。可见我们执行的是UPDATE但是数据库里面会存在两条数据,数据值和有效期不同,也就是版本不同 。

  • DELETE (假设执行DELETE之前的表内容如下)

DELETE FROM Emp WHERE ENo = 22217

说明: 假设我们是在 2012-06-01 00:00:00 执行的DELETE,则Sys_End值由 9999-12-31 23:59:59 变成了 2012-06-01 00:00:00 , 也就是在执行DELETE时候没有真正的删除符合条件的行,而是系统将符合条件的行的Sys_end修改为执行DELETE的事物时间。标识数据的有效期到DELETE执行那一刻为止。

  • SELECT
1
2
SELECT ENo,EName,Sys_Start,Sys_End FROM Emp
FOR SYSTEM_TIME AS OF TIMESTAMP '2011-01-02 00:00:00'

说明: 这个查询会返回所有 Sys_start <= 2011-01-02 00:00:00 并且 Sys_end > 2011-01-02 00:00:00 的记录。

Apache Flink® SQL Training

Flink SQL

  • Flink SQL/Table 如何转化为Flink graph?
  • Blink 进行了什么优化?

本节将主要从 SQL/Table API 如何转化为真正的 Job Graph 的流程开始,让大家对 Blink Planner 有一个比较清晰的认识,希望对大家阅读 Blink 代码,或者使用 Blink 方面有所帮助。然后介绍 Blink Planner 的改进及优化。

从上图可以很清楚的看到,解析的过程涉及到了三层:Table API/SQL,Blink Planner,Runtime,下面将对主要的步骤进行讲解:

Table API&SQL 解析验证:在 Flink 1.9 中,Table API 进行了大量的重构,引入了一套新的 Operation,这套 Operation 主要是用来描述任务的 Logic Tree。

当 SQL 传输进来后,首先会去做 SQL 解析,SQL 解析完成之后,会得到 SqlNode Tree(抽象语法树),然后会紧接着去做 Validator(验证),验证时会去访问 FunctionManger 和 CatalogManger,FunctionManger 主要是查询用户定义的 UDF,以及检查 UDF 是否合法,CatalogManger 主要是检查这个 Table 或者 Database 是否存在,如果验证都通过,就会生成一个 Operation DAG(有向无环图)。

从这一步可以看出,Table API 和 SQL 在 Flink 中最终都会转化为统一的结构,即 Operation DAG。

生成RelNode:Operation DAG 会被转化为 RelNode(关系表达式) DAG。

优化:优化器会对 RelNode 做各种优化,优化器的输入是各种优化的规则,以及各种统计信息。当前,在 Blink Planner 里面,绝大部分的优化规则,Stream 和 Batch 是共享的。差异在于,对 Batch 而言,它没有 state 的概念,而对于 Stream 而言,它是不支持 sort 的,所以目前 Blink Planner 中,还是运行了两套独立的规则集(Rule Set),然后定义了两套独立的 Physical Rel:BatchPhysical Rel 和 StreamPhysical Rel。优化器优化的结果,就是具体的 Physical Rel DAG。

转化:得到 Physical Rel Dag 后,继续会转化为 ExecNode,通过名字可以看出,ExecNode 已经属于执行层的概念了,但是这个执行层是 Blink 的执行层,在 ExecNode 中,会进行大量的 CodeGen 的操作,还有非 Code 的 Operator 操作,最后,将 ExecNode 转化为 Transformation DAG。

**生成可执行 Job Graph:**得到 Transformation DAG 后,最终会被转化成 Job Graph,完成 SQL 或者 Table API 的解析。

Blink Planner 功能方面改进主要包含如下几个方面:

  • 更完整的 SQL 语法支持:例如,IN,EXISTS,NOT EXISTS,子查询,完整的 Over 语句,Group Sets 等。而且已经跑通了所有的 TPCH,TPCDS 这两个测试集,性能还非常不错。
  • 提供了更丰富,高效的算子。
  • 提供了非常完善的 cost 模型,同时能够对接 Catalog 中的统计信息,使 cost 根据统计信息得到更优的执行计划。
  • 支持 join reorder。
  • shuffle service:对 Batch 而言,Blink Planner 还支持 shuffle service,这对 Batch 作业的稳定性有非常大的帮助,如果遇到 Batch 作业失败,通过 shuffle service 能够很快的进行恢复。

性能方面,主要包括以下部分:

  • 分段优化。

  • Sub-Plan Reuse。

  • 更丰富的优化 Rule:共一百多个 Rule ,并且绝大多数 Rule 是 Stream 和 Batch 共享的。

  • 更高效的数据结构 BinaryRow:能够节省序列化和反序列化的操作。

  • mini-batch 支持(仅 Stream):节省 state 的访问的操作。

  • 节省多余的 Shuffle 和 Sort(Batch 模式):两个算子之间,如果已经按 A 做 Shuffle,紧接着他下的下游也是需要按 A Shuffle 的数据,那中间的这一层 Shuffle,就可以省略,这样就可以省很多网络的开销,Sort 的情况也是类似。Sort 和 Shuffle 如果在整个计算里面是占大头,对整个性能是有很大的提升的。

深入性能优化及实践

本节中,将使用具体的示例进行讲解,让你深入理解 Blink Planner 性能优化的设计。

分段优化

示例 5

1
2
3
create view MyView as select word, count(1) as freq from SourceTable group by word; insert into SinkTable1 select \* from MyView where freq >10;

insert into SinkTable2 select count(word) as freq2, freq from MyView group by freq;

上面的这几个 SQL,转化为 RelNode DAG,大致图形如下:

图5 示例5 RelNode DAG

如果是使用 Flink Planner,经过优化层后,会生成如下执行层的 DAG:

图6 示例 5 Flink Planner DAG

可以看到,Flink Planner 只是简单的从 Sink 出发,反向的遍历到 Source,从而形成两个独立的执行链路,从上图也可以清楚的看到,Scan 和第一层 Aggregate 是有重复计算的。

在 Blink Planner 中,经过优化层之后,会生成如下执行层的 DAG:

图7 示例 5 Blink Planner DAG

Blink Planner 不是在每次调用 insert into 的时候就开始优化,而是先将所有的 insert into 操作缓存起来,等到执行前才进行优化,这样就可以看到完整的执行图,可以知道哪些部分是重复计算的。Blink Planner 通过寻找可以优化的最大公共子图,找到这些重复计算的部分。经过优化后,Blink Planner 会将最大公共子图的部分当做一个临时表,供其他部分直接使用。

这样,上面的图可以分为三部分,最大公共子图部分(临时表),临时表与 Filter 和 SinkTable1 优化,临时表与第二个 Aggregate 和 SinkTable 2 优化。

Blink Planner 其实是通过声明的 View 找到最大公共子图的,因此在开发过程中,如果需要复用某段逻辑,就将其定义为 View,这样就可以充分利用 Blink Planner 的分段优化功能,减少重复计算。

当然,当前的优化也不是最完美的,因为提前对图进行了切割,可能会导致一些优化丢失,今后会持续地对这部分算法进行改进。

总结一下,Blink Planner 的分段优化,其实解的是多 Sink 优化问题(DAG 优化),单 Sink 不是分段优化关心的问题,单 Sink 可以在所有节点上优化,不需要分段。

Sub-Plan Reuse

示例 6

1
2
3
4
5
6
7
insert into SinkTabl

select freq from (select word, count(1) as freq from SourceTable group by word) t where word like 'T%'

union all

select count(word) as freq2 from (select word, count(1) as freq from SourceTable group by word) t group by freq;

这个示例的 SQL 和分段优化的 SQL 其实是类似的,不同的是,没有将结果 Sink 到两个 Table 里面,而是将结果 Union 起来,Sink 到一个结果表里面。

下面看一下转化为 RelNode 的 DAG 图:

图 8 示例 6 RelNode DAG

从上图可以看出,Scan 和第一层的 Aggregate 也是有重复计算的,Blink Planner 其实也会将其找出来,变成下面的图:

图9 示例 6 Blink Planner DAG

Sub-Plan 优化的启用,有两个相关的配置:

  • table.optimizer.reuse-sub-plan-enabled (默认开启)

  • table.optimizer.reuse-source-enabled(默认开启)

这两个配置,默认都是开启的,用户可以根据自己的需求进行关闭。这里主要说明一下 table.optimizer.reuse-source-enabled 这个参数。在 Batch 模式下,join 操作可能会导致死锁,具体场景是在执行 hash-join 或者 nested-loop-join 时一定是先读 build 端,然后再读 probe 端,如果启用 reuse-source-enabled,当数据源是同一个 Source 的时候,Source 的数据会同时发送给 build 和 probe 端。这时候,build 端的数据将不会被消费,导致 join 操作无法完成,整个 join 就被卡住了。

为了解决死锁问题,Blink Planner 会先将 probe 端的数据落盘,这样 build 端读数据的操作才会正常,等 build 端的数据全部读完之后,再从磁盘中拉取 probe 端的数据,从而解决死锁问题。但是,落盘会有额外的开销,会多一次写的操作;有时候,读两次 Source 的开销,可能比一次写的操作更快,这时候,可以关闭 reuse-source,性能会更好。当然,如果读两次 Source 的开销,远大于一次落盘的开销,可以保持 reuse-source 开启。需要说明的是,Stream 模式是不存在死锁问题的,因为 Stream 模式 join 不会有选边的问题。

总结而言,sub-plan reuse 解的问题是优化结果的子图复用问题,它和分段优化类似,但他们是一个互补的过程。

注:Hash Join:对于两张待 join 的表 t1, t2。选取其中的一张表按照 join 条件给的列建立hash 表。然后扫描另外一张表,一行一行去建好的 hash 表判断是否有对应相等的行来完成 join 操作,这个操作称之为 probe (探测)。前一张表叫做 build 表,后一张表的叫做 probe 表。

Agg 分类优化

Blink 中的 Aggregate 操作是非常丰富的:

  • group agg,例如:select count(a) from t group by b

  • over agg,例如:select count(a) over (partition by b order by c) from t

  • window agg,例如:select count(a) from t group by tumble(ts, interval ‘10’ second), b

  • table agg ,例如:tEnv.scan(‘t’).groupBy(‘a’).flatAggregate(flatAggFunc(‘b’ as (‘c’, ‘d’)))

下面主要对 Group Agg 优化进行讲解,主要是两类优化。

1. Local/Global Agg 优化

Local/Global Agg 主要是为了减少网络 Shuffle。要运用 Local/Global 的优化,必要条件如下:

  • Aggregate 的所有 Agg Function 都是 mergeable 的,每个 Aggregate 需要实现 merge 方法,例如 SUM,COUNT,AVG,这些都是可以分多阶段完成,最终将结果合并;但是求中位数,计算 95% 这种类似的问题,无法拆分为多阶段,因此,无法运用 Local/Global 的优化。

  • table.optimizer.agg-phase-strategy 设置为 AUTO 或者 TWO_PHASE。

  • Stream 模式下,mini-batch 开启 ;Batch 模式下 AUTO 会根据 cost 模型加上统计数据,选择是否进行 Local/Global 优化。

示例 7

select count(*) from t group by color

没有优化的情况下,下面的这个 Aggregate 会产生 10 次的 Shuffle 操作。

图 10 示例 7 未做优化的 Count 操作

使用 Local/Global 优化后,会转化为下面的操作,会在本地先进行聚合,然后再进行 Shuffle 操作,整个 Shuffle 的数据剩下 6 条。在 Stream 模式下,Blink 其实会以 mini-batch 的维度对结果进行预聚合,然后将结果发送给 Global Agg 进行汇总。

图 11 示例 7 经过 Local/Global 优化的 Count 操作

2. Distinct Agg 优化

Distinct Agg 进行优化,主要是对 SQL 语句进行改写,达到优化的目的。但 Batch 模式和 Stream 模式解决的问题是不同的:

  • Batch 模式下的 Distinct Agg,需要先做 Distinct,再做 Agg,逻辑上需要两步才能实现,直接实现 Distinct Agg 开销太大。

  • Stream 模式下,主要是解决热点问题,因为 Stream 需要将所有的输入数据放在 State 里面,如果数据有热点,State 操作会很频繁,这将影响性能。

Batch 模式

第一层,求 distinct 的值和非 distinct agg function 的值,第二层求 distinct agg function 的值

示例 8

select color, count(distinct id), count(*) from t group by color

手工改写成:

1
2
3
4
5
6
7
8
9
10
11
12
13
select color, count(id), min(cnt) from (

select color, id, count(*) filter (where $e=2) as cnt from (

select color, id, 1 as $e from t --for distinct id

union all

select color, null as id, 2 as $e from t -- for count(\*)

) group by color, id, $e

) group by color

转化的逻辑过程,如下图所示:

图 12 示例 8 Batch 模式 Distinct 改写逻辑

Stream 模式

Stream 模式的启用有一些必要条件:

  • 必须是支持的 agg function:avg/count/min/max/sum/first_value/concat_agg/single_value;

  • table.optimizer.distinct-agg.split.enabled(默认关闭)

示例 9

select color, count(distinct id), count(*) from t group by color

手工改写成:

select color, sum(dcnt), sum(cnt) from (

select color, count(distinct id) as dcnt, count(*) as cnt from t

group by color, mod(hash_code(id), 1024)

) group by color

改写前,逻辑图大概如下:

图 13 示例 9 Stream 模式未优化 Distinct

改写后,逻辑图就会变为下面这样,热点数据被打散到多个中间节点上。

图14 示例 9 Stream 模式优化 Distinct

需要注意的是,示例 5 的 SQL 中 mod(hash_code(id),1024)中的这个 1024 为打散的维度,这个值建议设置大一些,设置太小产生的效果可能不好。

总结

本文首先对新的 TableEnvironment 的整体设计进行了介绍,并且列举了各种模式下TableEnvironment 的选择,然后通过具体的示例,展示了各种模式下代码的写法,以及需要注意的事项。

在新的 Catalog 和 DDL 部分,对 Catalog 的整体设计、DDL 的使用部分也都以实例进行拆分讲解。最后,对 Blink Planner 解析 SQL/Table API 的流程、Blink Planner 的改进以及优化的原理进行了讲解,希望对大家探索和使用 Flink SQL 有所帮助。

SQL解析工具

hive使用了antlr3实现了自己的HQL,
Flink使用Apache Calcite,
而Calcite的解析器是使用JavaCC实现的,
Spark2.x以后采用了antlr4实现自己的解析器,
Presto也是使用antlr4。

Window

在Apache Flink中有2种类型的Window,一种是OverWindow,即传统数据库的标准开窗,每一个元素都对应一个窗口。一种是GroupWindow,目前在SQL中GroupWindow都是基于时间进行窗口划分的。

Over Window

Apache Flink中对OVER Window的定义遵循标准SQL的定义语法。
按ROWS和RANGE分类是传统数据库的标准分类方法,在Apache Flink中还可以根据时间类型(ProcTime/EventTime)和窗口的有限和无限(Bounded/UnBounded)进行分类,共计8种类型。为了避免大家对过细分类造成困扰,我们按照确定当前行的不同方式将OVER Window分成两大类进行介绍,如下:

  • ROWS OVER Window - 每一行元素都视为新的计算行,即,每一行都是一个新的窗口。
  • RANGE OVER Window - 具有相同时间值的所有元素行视为同一计算行,即,具有相同时间值的所有行都是同一个窗口。

Bounded ROWS OVER Window

Bounded ROWS OVER Window 每一行元素都视为新的计算行,即,每一行都是一个新的窗口。

语义

我们以3个元素(2 PRECEDING)的窗口为例,如下图:
image

上图所示窗口 user 1 的 w5和w6, user 2的 窗口 w2 和 w3,虽然有元素都是同一时刻到达,但是他们仍然是在不同的窗口,这一点有别于RANGE OVER Window。

语法

Bounded ROWS OVER Window 语法如下:

1
2
3
4
5
6
7
8
SELECT 
agg1(col1) OVER(
[PARTITION BY (value_expression1,..., value_expressionN)]
ORDER BY timeCol
ROWS
BETWEEN (UNBOUNDED | rowCount) PRECEDING AND CURRENT ROW) AS colName,
...
FROM Tab1
  • value_expression - 进行分区的字表达式;
  • timeCol - 用于元素排序的时间字段;
  • rowCount - 是定义根据当前行开始向前追溯几行元素。
SQL 示例

利用item_tab测试数据,我们统计同类商品中当前和当前商品之前2个商品中的最高价格。

1
2
3
4
5
6
7
8
9
10
SELECT  
itemID,
itemType,
onSellTime,
price,
MAX(price) OVER (
PARTITION BY itemType
ORDER BY onSellTime
ROWS BETWEEN 2 preceding AND CURRENT ROW) AS maxPrice
FROM item_tab

Result

itemID itemType onSellTime price maxPrice
ITEM001 Electronic 2017-11-11 10:01:00 20 20
ITEM002 Electronic 2017-11-11 10:02:00 50 50
ITEM003 Electronic 2017-11-11 10:03:00 30 50
ITEM004 Electronic 2017-11-11 10:03:00 60 60
ITEM005 Electronic 2017-11-11 10:05:00 40 60
ITEM006 Electronic 2017-11-11 10:06:00 20 60
ITEM007 Electronic 2017-11-11 10:07:00 70 70
ITEM008 Clothes 2017-11-11 10:08:00 20 20

Bounded RANGE OVER Window

Bounded RANGE OVER Window 具有相同时间值的所有元素行视为同一计算行,即,具有相同时间值的所有行都是同一个窗口。

语义

我们以3秒中数据(INTERVAL ‘2’ SECOND)的窗口为例,如下图:
image

注意: 上图所示窗口 user 1 的 w6, user 2的 窗口 w3,元素都是同一时刻到达,他们是在同一个窗口,这一点有别于ROWS OVER Window。

语法

Bounded RANGE OVER Window的语法如下:

1
2
3
4
5
6
7
8
SELECT 
agg1(col1) OVER(
[PARTITION BY (value_expression1,..., value_expressionN)]
ORDER BY timeCol
RANGE
BETWEEN (UNBOUNDED | timeInterval) PRECEDING AND CURRENT ROW) AS colName,
...
FROM Tab1
  • value_expression - 进行分区的字表达式;
  • timeCol - 用于元素排序的时间字段;
  • timeInterval - 是定义根据当前行开始向前追溯指定时间的元素行;
SQL 示例

我们统计同类商品中当前和当前商品之前2分钟商品中的最高价格。

1
2
3
4
5
6
7
8
9
10
SELECT  
itemID,
itemType,
onSellTime,
price,
MAX(price) OVER (
PARTITION BY itemType
ORDER BY rowtime
RANGE BETWEEN INTERVAL '2' MINUTE preceding AND CURRENT ROW) AS maxPrice
FROM item_tab
Result(Bounded RANGE OVER Window)
itemID itemType onSellTime price maxPrice
ITEM001 Electronic 2017-11-11 10:01:00 20 20
ITEM002 Electronic 2017-11-11 10:02:00 50 50
ITEM003 Electronic 2017-11-11 10:03:00 30 60
ITEM004 Electronic 2017-11-11 10:03:00 60 60
ITEM005 Electronic 2017-11-11 10:05:00 40 60
ITEM006 Electronic 2017-11-11 10:06:00 20 40
ITEM007 Electronic 2017-11-11 10:07:00 70 70
ITEM008 Clothes 2017-11-11 10:08:00 20 20

特别说明

OverWindow最重要是要理解每一行数据都确定一个窗口,同时目前在Apache Flink中只支持按时间字段排序。并且OverWindow开窗与GroupBy方式数据分组最大的不同在于,GroupBy数据分组统计时候,在SELECT中除了GROUP BY的key,不能直接选择其他非key的字段,但是OverWindow没有这个限制,SELECT可以选择任何字段。比如一张表table(a,b,c,d)4个字段,如果按d分组求c的最大值,两种写完如下:

  • GROUP BY - SELECT d, MAX(c) FROM table GROUP BY d
  • OVER Window = SELECT a, b, c, d, MAX(c) OVER(PARTITION BY d, ORDER BY ProcTime())
    如上 OVER Window 虽然PARTITION BY d,但SELECT 中仍然可以选择 a,b,c字段。但在GROUPBY中,SELECT 只能选择 d 字段。

Group Window

根据窗口数据划分的不同,目前Apache Flink有如下3种Bounded Winodw:

  • Tumble - 滚动窗口,窗口数据有固定的大小,窗口数据无叠加;
  • Hop - 滑动窗口,窗口数据有固定大小,并且有固定的窗口重建频率,窗口数据有叠加;
  • Session - 会话窗口,窗口数据没有固定的大小,根据窗口数据活跃程度划分窗口,窗口数据无叠加。

说明: Aapche Flink 还支持UnBounded的 Group Window,也就是全局Window,流上所有数据都在一个窗口里面,语义非常简单,这里不做详细介绍了。

Tumble

语义

Tumble 滚动窗口有固定size,窗口数据不重叠,具体语义如下:
image

语法

Tumble 滚动窗口对应的语法如下:

1
2
3
4
5
6
7
8
9
SELECT 
[gk],
[TUMBLE_START(timeCol, size)],
[TUMBLE_END(timeCol, size)],
agg1(col1),
...
aggn(colN)
FROM Tab1
GROUP BY [gk], TUMBLE(timeCol, size)
  • [gk] - 决定了流是Keyed还是/Non-Keyed;
  • TUMBLE_START - 窗口开始时间;
  • TUMBLE_END - 窗口结束时间;
  • timeCol - 是流表中表示时间字段;
  • size - 表示窗口的大小,如 秒,分钟,小时,天。
SQL 示例

利用pageAccess_tab测试数据,我们需要按不同地域统计每2分钟的淘宝首页的访问量(PV)。

1
2
3
4
5
6
7
SELECT  
region,
TUMBLE_START(rowtime, INTERVAL '2' MINUTE) AS winStart,
TUMBLE_END(rowtime, INTERVAL '2' MINUTE) AS winEnd,
COUNT(region) AS pv
FROM pageAccess_tab
GROUP BY region, TUMBLE(rowtime, INTERVAL '2' MINUTE)
Result
region winStart winEnd pv
BeiJing 2017-11-11 02:00:00.0 2017-11-11 02:02:00.0 1
BeiJing 2017-11-11 02:10:00.0 2017-11-11 02:12:00.0 2
ShangHai 2017-11-11 02:00:00.0 2017-11-11 02:02:00.0 1
ShangHai 2017-11-11 04:10:00.0 2017-11-11 04:12:00.0 1

Hop

Hop 滑动窗口和滚动窗口类似,窗口有固定的size,与滚动窗口不同的是滑动窗口可以通过slide参数控制滑动窗口的新建频率。因此当slide值小于窗口size的值的时候多个滑动窗口会重叠。

语义

Hop 滑动窗口语义如下所示:
image

语法

Hop 滑动窗口对应语法如下:

1
2
3
4
5
6
7
8
9
SELECT 
[gk],
[HOP_START(timeCol, slide, size)] ,
[HOP_END(timeCol, slide, size)],
agg1(col1),
...
aggN(colN)
FROM Tab1
GROUP BY [gk], HOP(timeCol, slide, size)
  • [gk] 决定了流是Keyed还是/Non-Keyed;
  • HOP_START - 窗口开始时间;
  • HOP_END - 窗口结束时间;
  • timeCol - 是流表中表示时间字段;
  • slide - 是滑动步伐的大小;
  • size - 是窗口的大小,如 秒,分钟,小时,天;
SQL 示例

利用pageAccessCount_tab测试数据,我们需要每5分钟统计近10分钟的页面访问量(PV).

1
2
3
4
5
6
SELECT  
HOP_START(rowtime, INTERVAL '5' MINUTE, INTERVAL '10' MINUTE) AS winStart,
HOP_END(rowtime, INTERVAL '5' MINUTE, INTERVAL '10' MINUTE) AS winEnd,
SUM(accessCount) AS accessCount
FROM pageAccessCount_tab
GROUP BY HOP(rowtime, INTERVAL '5' MINUTE, INTERVAL '10' MINUTE)
Result
winStart winEnd accessCount
2017-11-11 01:55:00.0 2017-11-11 02:05:00.0 186
2017-11-11 02:00:00.0 2017-11-11 02:10:00.0 396
2017-11-11 02:05:00.0 2017-11-11 02:15:00.0 243
2017-11-11 02:10:00.0 2017-11-11 02:20:00.0 33
2017-11-11 04:05:00.0 2017-11-11 04:15:00.0 129
2017-11-11 04:10:00.0 2017-11-11 04:20:00.0 129

Session

Seeeion 会话窗口 是没有固定大小的窗口,通过session的活跃度分组元素。不同于滚动窗口和滑动窗口,会话窗口不重叠,也没有固定的起止时间。一个会话窗口在一段时间内没有接收到元素时,即当出现非活跃间隙时关闭。一个会话窗口 分配器通过配置session gap来指定非活跃周期的时长.

语义

Session 会话窗口语义如下所示:

image

语法

Seeeion 会话窗口对应语法如下:

1
2
3
4
5
6
7
8
9
SELECT 
[gk],
SESSION_START(timeCol, gap) AS winStart,
SESSION_END(timeCol, gap) AS winEnd,
agg1(col1),
...
aggn(colN)
FROM Tab1
GROUP BY [gk], SESSION(timeCol, gap)
  • [gk] 决定了流是Keyed还是/Non-Keyed;
  • SESSION_START - 窗口开始时间;
  • SESSION_END - 窗口结束时间;
  • timeCol - 是流表中表示时间字段;
  • gap - 是窗口数据非活跃周期的时长;
SQL 示例

利用pageAccessSession_tab测试数据,我们按地域统计连续的两个访问用户之间的访问时间间隔不超过3分钟的的页面访问量(PV).

1
2
3
4
5
6
7
SELECT  
region,
SESSION_START(rowtime, INTERVAL '3' MINUTE) AS winStart,
SESSION_END(rowtime, INTERVAL '3' MINUTE) AS winEnd,
COUNT(region) AS pv
FROM pageAccessSession_tab
GROUP BY region, SESSION(rowtime, INTERVAL '3' MINUTE)
Result
region winStart winEnd pv
BeiJing 2017-11-11 02:10:00.0 2017-11-11 02:13:00.0 1
ShangHai 2017-11-11 02:01:00.0 2017-11-11 02:08:00.0 4
ShangHai 2017-11-11 02:10:00.0 2017-11-11 02:14:00.0 2
ShangHai 2017-11-11 04:16:00.0 2017-11-11 04:19:00.0 1

UDX

Apache Flink 除了提供了大部分ANSI-SQL的核心算子,也为用户提供了自己编写业务代码的机会,那就是User-Defined Function,目前支持如下三种 User-Defined Function:

  • UDF - User-Defined Scalar Function
  • UDTF - User-Defined Table Function
  • UDAF - User-Defined Aggregate Funciton

UDX都是用户自定义的函数,那么Apache Flink框架为啥将自定义的函数分成三类呢?是根据什么划分的呢?Apache Flink对自定义函数进行分类的依据是根据函数语义的不同,函数的输入和输出不同来分类的,具体如下:

UDX INPUT OUTPUT INPUT:OUTPUT
UDF 单行中的N(N>=0)列 单行中的1列 1:1
UDTF 单行中的N(N>=0)列 M(M>=0)行 1:N(N>=0)
UDAF M(M>=0)行中的每行的N(N>=0)列 单行中的1列 M:1(M>=0)

UDF

  • 定义
    用户想自己编写一个字符串联接的UDF,我们只需要实现ScalarFunction#eval()方法即可,简单实现如下:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
object MyConnect extends ScalarFunction {
@varargs
def eval(args: String*): String = {
val sb = new StringBuilder
var i = 0
while (i < args.length) {
if (args(i) == null) {
return null
}
sb.append(args(i))
i += 1
}
sb.toString
}
}
  • 使用
1
2
3
val fun = MyConnect
tEnv.registerFunction("myConnect", fun)
val sql = "SELECT myConnect(a, b) as str FROM tab"

UDTF

  • 定义
    用户想自己编写一个字符串切分的UDTF,我们只需要实现TableFunction#eval()方法即可,简单实现如下:

ScalarFunction#eval()`

1
2
3
4
5
6
7
8
9
10
11
12
13
class MySplit extends TableFunction[String] {
def eval(str: String): Unit = {
if (str.contains("#")){
str.split("#").foreach(collect)
}
}

def eval(str: String, prefix: String): Unit = {
if (str.contains("#")) {
str.split("#").foreach(s => collect(prefix + s))
}
}
}
  • 使用
1
2
3
val fun = new MySplit()
tEnv.registerFunction("mySplit", fun)
val sql = "SELECT c, s FROM MyTable, LATERAL TABLE(mySplit(c)) AS T(s)"

UDAF

  • 定义
    UDAF 要实现的接口比较多,我们以一个简单的CountAGG为例,做简单实现如下:
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
class CountAccumulator extends JTuple1[Long] {
f0 = 0L
}


class MyCount
extends AggregateFunction[JLong, CountAccumulator] {



def accumulate(acc: CountAccumulator): Unit = {
acc.f0 += 1L
}



def retract(acc: CountAccumulator): Unit = {
acc.f0 -= 1L
}

def accumulate(acc: CountAccumulator, value: Any): Unit = {
if (value != null) {
acc.f0 += 1L
}
}

def retract(acc: CountAccumulator, value: Any): Unit = {
if (value != null) {
acc.f0 -= 1L
}
}

override def getValue(acc: CountAccumulator): JLong = {
acc.f0
}

def merge(acc: CountAccumulator, its: JIterable[CountAccumulator]): Unit = {
val iter = its.iterator()
while (iter.hasNext) {
acc.f0 += iter.next().f0
}
}

override def createAccumulator(): CountAccumulator = {
new CountAccumulator
}

def resetAccumulator(acc: CountAccumulator): Unit = {
acc.f0 = 0L
}

override def getAccumulatorType: TypeInformation[CountAccumulator] = {
new TupleTypeInfo(classOf[CountAccumulator], BasicTypeInfo.LONG_TYPE_INFO)
}

override def getResultType: TypeInformation[JLong] =
BasicTypeInfo.LONG_TYPE_INFO
}
  • 使用
1
2
3
val fun = new MyCount()
tEnv.registerFunction("myCount", fun)
val sql = "SELECT myCount(c) FROM MyTable GROUP BY a"

上面我们介绍了Apache Flink SQL核心算子的语法及语义,这部分将选取Bounded EventTime Tumble Window为例为大家编写一个完整的包括Source和Sink定义的Apache Flink SQL Job。假设有一张淘宝页面访问表(PageAccess_tab),有地域,用户ID和访问时间。我们需要按不同地域统计每2分钟的淘宝首页的访问量(PV). 具体数据如下:

region userId accessTime
ShangHai U0010 2017-11-11 10:01:00
BeiJing U1001 2017-11-11 10:01:00
BeiJing U2032 2017-11-11 10:10:00
BeiJing U1100 2017-11-11 10:11:00
ShangHai U0011 2017-11-11 12:10:00

大家都知道,在 Flink 中,通过 Table API 和 SQL 实现的流处理逻辑,最终会翻译为基于 DataStreamAPI 实现的 DataStream 作业,返回这个作业输出的 DataStream (writeToSink 本质上也是先得到 DataStream 作业,再为其输出 DataStream 加上一个DataStreamSink) 。

从一段 SQL 到 DataStream 作业,其过程简单描述如下:

  1. 在 TableEnvironment,即“表环境”,将数据源注册为动态表。例如,通过表环境的接口`registerDataStream`, 作为源的DataStream,即数据流, 在表环境注册为动态表

  2. 通过表环境的接口 `sqlQuery`,将 SQL 构造为 Table 对象

  3. 通过toAppendStream/toRetractedStream接口,即翻译接口,将 Table 对象表达的作业逻辑,翻译为 DataStream 作业。

图片

在调用翻译接口,将 Table 对象翻译为 DataStream 作业时,通过翻译接口传入的 TTL 配置,递归传递到各个计算节点的翻译、构造逻辑里,使得翻译出来的 DataStream 算子的内部状态按照该 TTL 配置及时清理。

【参考文献】

  1. 在数据流中使用SQL查询:Apache Flink中的动态表的持续查询

  2. Flink Table API & SQL编程指南

Flink-SQL语法

Apache Flink SQL training

group window

groupByWindow会直接生成回撤流

1
2
3
4
5
6
7
8
9
10
11
12
13
insert       
into
dim_result_lct_activy_config
select
Fact_id,
LAST_VALUE(Fact_name),
regexp_Replace( LAST_VALUE(Fact_start_time), '-|:|\s','') as startTime,
regexp_Replace(LAST_VALUE(Fact_end_time),'-|:|\s','') as endTime,
LAST_VALUE(Fstate)
from
db_act_config_t_act_logic_config
group by
Fact_id

这是一个同步数据的demo,db_act_config_t_act_logic_config 是kafka数据源,来自源MySQL的变更数据;dim_result_lct_activy_config是目的表,Fact_id为主键。

group window生成retract stream

1
insert into mysql_sink select fkey,count(1) as cnt from kafka_source group by fkey

上述语句是一个group window, 每从kafka中过来一条数据,都会产生两条记录(Tuple2<Row,Boolean>), 删除旧记录,添加新记录。

group window会产生 retract stream, 下游系统必须支持retract stream,(目前共有三种sink: AppendStreamSink, UpsertStreamSink, RetractStreamSink )

Flink-connector-JDBC 使用的是JDBCUpsertTableSink.java写入MySQL, 支持Retract

https://github.com/apache/flink/tree/master/flink-connectors/flink-jdbc/src/main/java/org/apache/flink/api/java/io/jdbc

Flink-connector-kafka 实现的是 AppendStreamSink,只支持insert,不支持retract. 会报错

AppendStreamTableSink requires that Table has only insert changes

1
insert into mysql_sink select fkey,count(1) as cnt from kafka_source

如果不带group by, 无法推导出unique key, 无法按照unique key更新

http://apache-flink.147419.n8.nabble.com/FlinkSQL-Upsert-Retraction-MySQL-td2785.html

1
2
3
4
5
6
7
8
9
10
11
/** 
* Get dialect upsert statement, the database has its own upsert syntax, such as Mysql
* using DUPLICATE KEY UPDATE, and PostgresSQL using ON CONFLICT... DO UPDATE SET..
*
* @return None if dialect does not support upsert statement, the writer will degrade to
* the use of select + update/insert, this performance is poor.
*/
default Optional<String> getUpsertStatement(
String tableName, String[] fieldNames, String[] uniqueKeyFields) {
return Optional.empty();
}

不同的数据库产品有不同的语句,所以默认实现是delete +insert

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
@Override 
public void executeBatch() throws SQLException {
if (keyToRows.size() > 0) {
for (Map.Entry<Row, Tuple2<Boolean, Row>> entry : keyToRows.entrySet()) {
Row pk = entry.getKey();
Tuple2<Boolean, Row> tuple = entry.getValue();
if (tuple.f0) {
processOneRowInBatch(pk, tuple.f1);
} else {
setRecordToStatement(deleteStatement, pkTypes, pk);
deleteStatement.addBatch();
}
}
internalExecuteBatch();
deleteStatement.executeBatch();
keyToRows.clear();
}
}

image-20210928204455343

image-20210928204646114

image-20210928204703778

Over window

SQL窗口函数 传统SQL窗口函数的介绍

1
2
3
4
5
6
7
8
9
10
11
12
13
14
select 
to_char(SYSTIMESTAMP(),'yyyymmddhh24miss') fetl_time,
*
from
(
select
*,
row_number() over(partition by fid order by fmodify_time desc,exp_time_stample_order desc) rn
from
db.table1
where
fdate=20210101
) t
where rn=1
1
2
3
4
5
6
7
8
9
10
11
SELECT COUNT(amount) OVER (
PARTITION BY user
ORDER BY proctime
ROWS BETWEEN 2 PRECEDING AND CURRENT ROW)
FROM Orders;
SELECT COUNT(amount) OVER w, SUM(amount) OVER w
FROM Orders
WINDOW w AS (
PARTITION BY user
ORDER BY proctime
ROWS BETWEEN 2 PRECEDING AND CURRENT ROW) ;

OVER 窗口应用示例

首先通过 DDL 定义源数据表和结果表,如下输入是用户行为消息,输出到计算结果消息。

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
CREATE TABLE `user_action` (
`user_id` VARCHAR,
`page_id` VARCHAR,
`action_type` VARCHAR,
`event_time` TIMESTAMP,
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
'connector.type' = 'kafka',
'connector.topic' = 'user_action',
'connector.version' = '0.11',
'connector.properties.0.key' = 'bootstrap.servers',
'connector.properties.0.value' = 'xxx:9092',
'connector.startup-mode' = 'latest-offset',
'update-mode' = 'append',
'...' = '...'
);

CREATE TABLE `agg_result` (
`user_id` VARCHAR,
`page_id` VARCHAR,
`result_type` VARCHAR,
`result_value` BIGINT
) WITH (
'connector.type' = 'kafka',
'connector.topic' = 'agg_result',
'...' = '...'
);

场景一,实时触发的最近2小时用户+页面维度的点击量,注意窗口是向前2小时,类似于实时触发的滑动窗口。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
insert into
agg_result
select
user_id,
page_id,
'click-type1' as result_type
count(1) OVER (
PARTITION BY user_id, page_id
ORDER BY event_time
RANGE BETWEEN INTERVAL '2' HOUR PRECEDING AND CURRENT ROW
) as result_value
from
user_action
where
action_type = 'click'

场景二,实时触发的当天用户+页面维度的浏览量,这就是开篇问题解法,其中多了一个日期维度分组条件,这样就做到输出结果从滑动时间转为固定时间(根据时间区间分组),因为 WATERMARK 机制,今天并不会有昨天数据到来(如果有都被自动抛弃),因此只会输出今天的分组结果。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
insert into
agg_result
select
user_id,
page_id,
'view-type1' as result_type
count(1) OVER (
PARTITION BY user_id, page_id, DATE_FORMAT(event_time, 'yyyyMMdd')
ORDER BY event_time
RANGE BETWEEN INTERVAL '1' DAY PRECEDING AND CURRENT ROW
) as result_value
from
user_action
where
action_type = 'view'

场景三,实时触发的当天用户+页面点击率 CTR(Click-Through-Rate),这相比前面增加了多个 OVER 聚合计算,可以将窗口定义写在最后。注意示例中缺少了类型转换,因为除法结果是 decimal,也缺少精度处理函数 ROUND。

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
insert into
agg_result
select
user_id,
page_id,
'ctr-type1' as result_type,
sum(
case
when action_type = 'click' then 1 else 0
end
) OVER w
/
if(
sum(
case
when action_type = 'view' then 1 else 0
end
) OVER w = 0,
1,
sum(
case
when action_type = 'view' then 1 else 0
end
) OVER w
)
as result_value
from
user_action
where
1 = 1
WINDOW w AS (
PARTITION BY user_id, page_id, DATE_FORMAT(event_time,'yyyyMMdd')
ORDER BY event_time
RANGE BETWEEN INTERVAL '1' DAY PRECEDING AND CURRENT ROW
)

此外,OVER 窗口聚合还可以支持查询子句、关联查询、UNION ALL 等组合,并可以实现对关联出来的列进行聚合等复杂情况。

实时TopN

SQL实时TopN

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
create table source_kafka 
(
userID String,
eventType String,
eventTime String,
productID String
) with (
'connector.type' = 'kafka',
'connector.version' = '0.10',
'connector.properties.bootstrap.servers' = 'kafka01:9092',
'connector.properties.zookeeper.connect' = 'kafka01:2181',
'connector.topic' = 'test_1',
'connector.properties.group.id' = 'c1_test_1',
'connector.startup-mode' = 'latest-offset',
'format.type' = 'json'
);
create table sink_mysql
(
datetime STRING,
productID STRING,
userID STRING,
clickPV BIGINT
) with (
'connector.type' = 'jdbc',
'connector.url' = 'jdbc:mysql://localhost:3306/bigdata',
'connector.table' = 't_product_click_topn',
'connector.username' = 'root',
'connector.password' = 'bigdata',
'connector.write.flush.max-rows' = '50',
'connector.write.flush.interval' = '2s',
'connector.write.max-retries' = '3'
);
INSERT INTO sink_mysql
SELECT datetime, productID, userID, clickPV
FROM (
SELECT *,
ROW_NUMBER() OVER (PARTITION BY datetime, productID ORDER BY clickPV desc) AS rownum
FROM (
SELECT SUBSTRING(eventTime,1,13) AS datetime,
productID,
userID,
count(1) AS clickPV
FROM source_kafka
GROUP BY SUBSTRING(eventTime,1,13), productID, userID
) a
) t
WHERE rownum <= 3;

OVER 窗口问题和优化

在底层实现中,所有细分 OVER 窗口的数据都是共享的,只存一份,这点不像滑动窗口会保存多份窗口数据。但是 OVER 窗口会把所有数据明细存在状态后端中(内存、RocksDB 或 HDFS),每一次窗口计算后会清除过期数据。因此如果向前窗口时间较大,或数据明细过多,可能会占用大量内存,即使通过 RocksDB 存在磁盘上,也有因为磁盘访问慢导致性能下降进而产生反压问题。在实现源码 RowTimeRangeBoundedPrecedingFunction 可以看到虽然每次窗口计算时新增聚合值和减少过期聚合值是增量式的,不用遍历全部窗口明细,但是为了计算过期数据,即超过 PRECEDING 的数据,仍然需要把存储的那些时间戳全部拿出来遍历,判断是否过期,以及是否要减少聚合值。我们尝试了通过数据有序性减少查询操作,但是效果并不明显,目前主要是配置调优和加大任务分片数进行优化。

简析Flink状态生存时间(State TTL)机制的底层实现

简析Flink状态生存时间(State TTL)机制的底层实现

前言

从Flink 1.6版本开始,社区为状态引入了TTL(time-to-live,生存时间)机制,支持Keyed State的自动过期,有效解决了状态数据在无干预情况下无限增长导致OOM的问题。State TTL的用法很简单,官方文档中给出的示例代码如下。

1
2
3
4
5
6
7
8
StateTtlConfig ttlConfig =
StateTtlConfig
.newBuilder(Time.seconds(1))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.build();
ValueStateDescriptor<String> stateDescriptor = new ValueStateDescriptor<>("text state", String.class);
stateDescriptor.enableTimeToLive(ttlConfig);

那么State TTL的背后又隐藏着什么样的思路呢?下面就从设置类StateTtlConfig入手开始研究(Flink代码版本为1.9.3)。

StateTtlConfig

该类中有5个成员属性,它们就是用户需要指定的全部参数了。

1
2
3
4
5
6
private final UpdateType updateType
private final StateVisibility stateVisibility
private final TtlTimeCharacteristic ttlTimeCharacteristic
private final Time ttl
private final CleanupStrategies cleanupStrategies

其中,ttl参数表示用户设定的状态生存时间。而UpdateType、StateVisibility和TtlTimeCharacteristic都是枚举,分别代表状态时间戳的更新方式、过期状态数据的可见性,以及对应的时间特征。它们的含义在注释中已经解释得很清楚了。

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
/** 
* This option value configures when to update last access timestamp which prolongs state TTL.
*/
public enum UpdateType {
/** TTL is disabled. State does not expire. */
Disabled,
/** Last access timestamp is initialised when state is created and updated on every write operation.
当每次写操作时,更新时间戳
*/
OnCreateAndWrite,
/** The same as <code>OnCreateAndWrite</code> but also updated on read.
当每次读写操作时,更新时间戳
*/
OnReadAndWrite
}
/** * This option configures whether expired user value can be returned or not. */
public enum StateVisibility {
/** Return expired user value if it is not cleaned up yet. */
ReturnExpiredIfNotCleanedUp,
/** Never return expired user value. */
NeverReturnExpired
}
/** * This option configures time scale to use for ttl. */
public enum TtlTimeCharacteristic {
/** Processing time, see also <code>org.apache.flink.streaming.api.TimeCharacteristic.ProcessingTime</code>. */
ProcessingTime
}

Flink目前仅支持基于处理时间的State TTL,事件时间会在不久的将来支持。

CleanupStrategies内部类则用来规定过期状态的特殊清理策略,用户在构造StateTtlConfig时,可以通过调用以下方法之一指定。

  • cleanupFullSnapshot()
    当对状态做全量快照时清理过期数据,对开启了增量检查点(incremental checkpoint)的RocksDB状态后端无效,对应源码中的EmptyCleanupStrategy。
    为什么叫做“空的”清理策略呢?因为该选项只能保证状态持久化时不包含过期数据,但TaskManager本地的过期状态则不作任何处理,所以无法从根本上解决OOM的问题,需要定期重启作业。
  • cleanupIncrementally(int cleanupSize, boolean runCleanupForEveryRecord)
    增量清理过期数据,默认在每次访问状态时进行清理,将runCleanupForEveryRecord设为true可以附加在每次写入/删除时清理。cleanupSize指定每次触发清理时检查的状态条数。
    仅对基于堆的状态后端有效,对应源码中的IncrementalCleanupStrategy。
  • cleanupInRocksdbCompactFilter(long queryTimeAfterNumEntries)
    当RocksDB做compaction操作时,通过Flink定制的过滤器(FlinkCompactionFilter)过滤掉过期状态数据。参数queryTimeAfterNumEntries用于指定在写入多少条状态数据后,通过状态时间戳来判断是否过期。
    该策略仅对RocksDB状态后端有效,对应源码中的RocksdbCompactFilterCleanupStrategy。CompactionFilter是RocksDB原生提供的机制,其说明可见这里

如果不调用上述方法,则采用默认的后台清理策略,下文有讲。

TtlStateFactory、TtlStateContext

在所有Keyed State状态后端的抽象基类AbstractKeyedStateBackend中,创建并记录一个状态实例的方法如下。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
@Override
@SuppressWarnings("unchecked")
public <N, S extends State, V> S getOrCreateKeyedState(
final TypeSerializer<N> namespaceSerializer, StateDescriptor<S, V> stateDescriptor) throws Exception {
checkNotNull(namespaceSerializer, "Namespace serializer");
checkNotNull(keySerializer, "State key serializer has not been configured in the config. " +
"This operation cannot use partitioned state.");
InternalKvState<K, ?, ?> kvState = keyValueStatesByName.get(stateDescriptor.getName());
if (kvState == null) {
if (!stateDescriptor.isSerializerInitialized()) {
stateDescriptor.initializeSerializerUnlessSet(executionConfig);
}
kvState = TtlStateFactory.createStateAndWrapWithTtlIfEnabled( namespaceSerializer, stateDescriptor, this, ttlTimeProvider);
keyValueStatesByName.put(stateDescriptor.getName(), kvState);
publishQueryableStateIfEnabled(stateDescriptor, kvState);
}
return (S) kvState;
}

可见是调用了TtlStateFactory.createStateAndWrapWithTtlIfEnabled()方法来真正创建。顾名思义,TtlStateFactory是产生TTL状态的工厂类。

1
2
3
4
5
6
7
8
9
10
11
12
public static <K, N, SV, TTLSV, S extends State, IS extends S> IS createStateAndWrapWithTtlIfEnabled( 
TypeSerializer<N> namespaceSerializer,
StateDescriptor<S, SV> stateDesc,
KeyedStateBackend<K> stateBackend,
TtlTimeProvider timeProvider
) throws Exception {
Preconditions.checkNotNull(namespaceSerializer);
Preconditions.checkNotNull(stateDesc);
Preconditions.checkNotNull(stateBackend);
Preconditions.checkNotNull(timeProvider);
return stateDesc.getTtlConfig().isEnabled() ? new TtlStateFactory<K, N, SV, TTLSV, S, IS>( namespaceSerializer, stateDesc, stateBackend, timeProvider) .createState() : stateBackend.createInternalState(namespaceSerializer, stateDesc);
}

由上可知,如果我们为状态描述符StateDescriptor加入了TTL,那么就会调用TtlStateFactory.createState()方法创建一个带有TTL的状态实例;否则,就调用StateBackend.createInternalState()创建一个普通的状态实例。TtlStateFactory.createState()的代码如下。

1
2
3
4
5
6
7
8
9
10
11
12
13
@SuppressWarnings("unchecked")
private IS createState() throws Exception {
SupplierWithException<IS, Exception> stateFactory = stateFactories.get(stateDesc.getClass());
if (stateFactory == null) {
String message = String.format("State %s is not supported by %s", stateDesc.getClass(), TtlStateFactory.class);
throw new FlinkRuntimeException(message);
}
IS state = stateFactory.get();
if (incrementalCleanup != null) {
incrementalCleanup.setTtlState((AbstractTtlState<K, N, ?, TTLSV, ?>) state);
}
return state;
}

其中,stateFactories是一个Map结构,维护了各种状态描述符与对应产生该种状态对象的工厂方法映射。所有的工厂方法都被包装成了Supplier(Java 8提供的函数式接口),所以在上述createState()方法中,可以通过Supplier.get()方法来实际执行createTtl.State()工厂方法,并获得新的状态实例。

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
this.stateFactories = createStateFactories();
@SuppressWarnings("deprecation")
private Map<Class<? extends StateDescriptor>, SupplierWithException<IS, Exception>> createStateFactories() {
return Stream.of(
Tuple2.of(ValueStateDescriptor.class, (SupplierWithException<IS, Exception>) this::createValueState),
Tuple2.of(ListStateDescriptor.class, (SupplierWithException<IS, Exception>) this::createListState),
Tuple2.of(MapStateDescriptor.class, (SupplierWithException<IS, Exception>) this::createMapState),
Tuple2.of(ReducingStateDescriptor.class, (SupplierWithException<IS, Exception>) this::createReducingState),
Tuple2.of(AggregatingStateDescriptor.class, (SupplierWithException<IS, Exception>) this::createAggregatingState),
Tuple2.of(FoldingStateDescriptor.class, (SupplierWithException<IS, Exception>) this::createFoldingState)
).collect(Collectors.toMap(t -> t.f0, t -> t.f1));
}

@SuppressWarnings("unchecked")
private IS createValueState() throws Exception {
ValueStateDescriptor<TtlValue<SV>> ttlDescriptor = new ValueStateDescriptor<>(
stateDesc.getName(),
new TtlSerializer<>(LongSerializer.INSTANCE, stateDesc.getSerializer())
);
return (IS) new TtlValueState<>(createTtlStateContext(ttlDescriptor));
}

@SuppressWarnings("unchecked")
private <T> IS createListState() throws Exception {
ListStateDescriptor<T> listStateDesc = (ListStateDescriptor<T>) stateDesc;
ListStateDescriptor<TtlValue<T>> ttlDescriptor = new ListStateDescriptor<>(
stateDesc.getName(),
new TtlSerializer<>(LongSerializer.INSTANCE, listStateDesc.getElementSerializer())
);
return (IS) new TtlListState<>(createTtlStateContext(ttlDescriptor));
}
// 以下略去...

可见,带有TTL的状态类名其实就是普通状态类名加上Ttl前缀,只是没有公开给用户而已。并且在生成Ttl$State时,还会通过createTtlStateContext()方法生成TTL状态的上下文。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
@SuppressWarnings("unchecked")
private <OIS extends State, TTLS extends State, V, TTLV> TtlStateContext<OIS, V>
createTtlStateContext(StateDescriptor<TTLS, TTLV> ttlDescriptor) throws Exception {
ttlDescriptor.enableTimeToLive(stateDesc.getTtlConfig());
// also used by RocksDB backend for TTL compaction filter config
OIS originalState = (OIS) stateBackend.createInternalState(
namespaceSerializer,
ttlDescriptor,
getSnapshotTransformFactory()
);
return new TtlStateContext<>(
originalState,
ttlConfig,
timeProvider,
(TypeSerializer<V>) stateDesc.getSerializer(),
registerTtlIncrementalCleanupCallback((InternalKvState<?, ?, ?>) originalState)
);
}

TtlStateContext的本质是对以下几个实例做了封装。

  • 原始State(通过StateBackend.createInternalState()方法创建)及其序列化器(通过StateDescriptor.getSerializer()方法取得);
  • StateTtlConfig,前文已经讲过;
  • TtlTimeProvider,用来提供判断状态过期标准的时间戳。当前只是简单地代理了System.currentTimeMillis(),没有任何其他代码;
  • 一个Runnable类型的回调方法,通过registerTtlIncrementalCleanupCallback()方法产生,用于状态数据的增量清理,后面会看到它的用途。

接下来就具体看看TTL状态是如何实现的。

AbstractTtlState、AbstractTtlDecorator

在解说之前,先放一幅类图。

所有Ttl.State都是AbstractTtlState的子类,而AbstractTtlState又是装饰器AbstractTtlDecorator的子类。AbstractTtlDecorator提供了最基本的TTL逻辑,代码不长,全部抄录如下。

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
abstract class AbstractTtlDecorator<T> {
/** Wrapped original state handler. */
final T original;
final StateTtlConfig config;
final TtlTimeProvider timeProvider;
/** Whether to renew expiration timestamp on state read access. */
final boolean updateTsOnRead;
/** Whether to renew expiration timestamp on state read access. */
final boolean returnExpired;
/** State value time to live in milliseconds. */
final long ttl;
AbstractTtlDecorator( T original, StateTtlConfig config, TtlTimeProvider timeProvider) {
// ......
}
// 查到已过期数据,返回null
<V> V getUnexpired(TtlValue<V> ttlValue) {
return ttlValue == null || (expired(ttlValue) && !returnExpired) ? null : ttlValue.getUserValue();
}
// 过期掉数据
<V> boolean expired(TtlValue<V> ttlValue) {
return TtlUtils.expired(ttlValue, ttl, timeProvider);
}
// 包装数据
<V> TtlValue<V> wrapWithTs(V value) {
return TtlUtils.wrapWithTs(value, timeProvider.currentTimestamp());
}
// 重新包装数据
<V> TtlValue<V> rewrapWithNewTs(TtlValue<V> ttlValue) {
return wrapWithTs(ttlValue.getUserValue());
}
<SE extends Throwable, CE extends Throwable, CLE extends Throwable, V> V getWithTtlCheckAndUpdate(
SupplierWithException<TtlValue<V>, SE> getter,
ThrowingConsumer<TtlValue<V>, CE> updater,
ThrowingRunnable<CLE> stateClear
) throws SE, CE, CLE {
TtlValue<V> ttlValue = getWrappedWithTtlCheckAndUpdate(getter, updater, stateClear);
return ttlValue == null ? null : ttlValue.getUserValue();
}
<SE extends Throwable, CE extends Throwable, CLE extends Throwable, V> TtlValue<V> getWrappedWithTtlCheckAndUpdate(
SupplierWithException<TtlValue<V>, SE> getter,
ThrowingConsumer<TtlValue<V>, CE> updater,
ThrowingRunnable<CLE> stateClear
) throws SE, CE, CLE {
TtlValue<V> ttlValue = getter.get();
if (ttlValue == null) {
return null;
} else if (expired(ttlValue)) {
stateClear.run();
if (!returnExpired) {
return null;
}
} else if (updateTsOnRead) {
updater.accept(rewrapWithNewTs(ttlValue));
}
return ttlValue;
}
}

它的成员属性比较容易理解,例如,updateTsOnRead表示在读取状态值时也更新时间戳(即UpdateType.OnReadAndWrite),returnExpired表示即使状态过期,在被真正删除之前也返回它的值(即StateVisibility.ReturnExpiredIfNotCleanedUp)。

状态值与TTL的包装(成为TtlValue)以及过期检测都由工具类TtlUtils来负责,思路很简单,代码如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
public class TtlUtils {
static <V> boolean expired(@Nullable TtlValue<V> ttlValue, long ttl, TtlTimeProvider timeProvider) {
return expired(ttlValue, ttl, timeProvider.currentTimestamp());
}
static <V> boolean expired(@Nullable TtlValue<V> ttlValue, long ttl, long currentTimestamp) {
return ttlValue != null && expired(ttlValue.getLastAccessTimestamp(), ttl, currentTimestamp);
}
static boolean expired(long ts, long ttl, TtlTimeProvider timeProvider) {
return expired(ts, ttl, timeProvider.currentTimestamp());
}
public static boolean expired(long ts, long ttl, long currentTimestamp) {
return getExpirationTimestamp(ts, ttl) <= currentTimestamp;
}
private static long getExpirationTimestamp(long ts, long ttl) {
long ttlWithoutOverflow = ts > 0 ? Math.min(Long.MAX_VALUE - ts, ttl) : ttl;
return ts + ttlWithoutOverflow;
}
static <V> TtlValue<V> wrapWithTs(V value, long ts) {
return new TtlValue<>(value, ts);
}
}

TtlValue的属性只有两个:状态值和时间戳,代码略去。

AbstractTtlDecorator核心方法是获取状态值的getWrappedWithTtlCheckAndUpdate(),它接受三个参数:

  • getter:一个可抛出异常的Supplier,用于获取状态值;
  • updater:一个可抛出异常的Consumer,用于更新状态的时间戳;
  • stateClear:一个可抛出异常的Runnable,用于异步删除过期状态。

可见,在默认情况下的后台清理策略是:只有状态值被读取时,才会做过期检测,并异步清除过期的状态。这种惰性清理的机制会导致那些实际已经过期但从未被再次访问过的状态无法被删除,需要特别注意。官方文档中也已有提示:

By default, expired values are explicitly removed on read, such as ValueState#value, and periodically garbage collected in the background if supported by the configured state backend.

当确认到状态过期时,会调用stateClear的逻辑进行删除;如果需要在读取时顺便更新状态的时间戳,会调用updater的逻辑重新包装一个TtlValue。

AbstractTtlState的代码更加简单,主要的方法列举如下。

1
2
3
4
5
6
7
final Runnable accessCallback; <SE extends Throwable, CE extends Throwable, T> T getWithTtlCheckAndUpdate(
SupplierWithException<TtlValue<T>, SE> getter,
ThrowingConsumer<TtlValue<T>, CE> updater) throws SE, CE {
return getWithTtlCheckAndUpdate(getter, updater, original::clear);} @Overridepublic void clear() {
original.clear();
accessCallback.run();
}

其中,accessCallback就是TtlStateContext中注册的增量清理回调。

下面以TtlMapState为例,看看具体的TTL状态如何利用上文所述的这些实现。

TtlMapState

以下是部分代码。

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
class TtlMapState<K, N, UK, UV>
extends AbstractTtlState<K, N, Map<UK, UV>, Map<UK, TtlValue<UV>>, InternalMapState<K, N, UK, TtlValue<UV>>>
implements InternalMapState<K, N, UK, UV> {
TtlMapState(TtlStateContext<InternalMapState<K, N, UK, TtlValue<UV>>, Map<UK, UV>> ttlStateContext) {
super(ttlStateContext);
}
@Override
public UV get(UK key) throws Exception {
TtlValue<UV> ttlValue = getWrapped(key);
return ttlValue == null ? null : ttlValue.getUserValue();
}
private TtlValue<UV> getWrapped(UK key) throws Exception {
accessCallback.run();
return getWrappedWithTtlCheckAndUpdate(() -> original.get(key), v -> original.put(key, v), () -> original.remove(key));
}
@Override
public void put(UK key, UV value) throws Exception {
accessCallback.run();
original.put(key, wrapWithTs(value));
}
@Override
public void putAll(Map<UK, UV> map) throws Exception {
accessCallback.run();
if (map == null) {
return;
}
Map<UK, TtlValue<UV>> ttlMap = new HashMap<>(map.size());
long currentTimestamp = timeProvider.currentTimestamp();
for (Map.Entry<UK, UV> entry : map.entrySet()) {
UK key = entry.getKey();
ttlMap.put(key, TtlUtils.wrapWithTs(entry.getValue(), currentTimestamp));
}
original.putAll(ttlMap);
}
@Override
public void remove(UK key) throws Exception {
accessCallback.run();
original.remove(key);
}
@Override
public boolean contains(UK key) throws Exception {
TtlValue<UV> ttlValue = getWrapped(key);
return ttlValue != null;
}
// ......
}

可见,TtlMapState的增删改查操作都是在原MapState上进行,只是加上了TTL相关的逻辑,这也是装饰器模式的特点。例如,TtlMapState.get()方法调用了上述AbstractTtlDecorator.getWrappedWithTtlCheckAndUpdate()方法,传入的获取(getter)、插入(updater)和删除(stateClear)的逻辑就是原MapState的get()、put()和remove()方法。而TtlMapState.put()只是在调用原MapState的put()方法之前,将状态包装为TtlValue而已。

增量清理策略

另外需要注意,所有增删改查操作之前都需要执行accessCallback.run()方法。如果启用了增量清理策略,该Runnable会通过在状态数据上维护一个全局迭代器向前清理过期数据。如果未启用增量清理策略,accessCallback为空。前文提到过的TtlStateFactory.registerTtlIncrementalCleanupCallback() 方法如下。

1
2
3
4
5
6
7
8
9
10
11
12
private Runnable registerTtlIncrementalCleanupCallback(InternalKvState<?, ?, ?> originalState) {
StateTtlConfig.IncrementalCleanupStrategy config =
ttlConfig.getCleanupStrategies().getIncrementalCleanupStrategy();
boolean cleanupConfigured = config != null && incrementalCleanup != null;
boolean isCleanupActive = cleanupConfigured &&
isStateIteratorSupported(originalState, incrementalCleanup.getCleanupSize());
Runnable callback = isCleanupActive ? incrementalCleanup::stateAccessed : () -> { };
if (isCleanupActive && config.runCleanupForEveryRecord()) {
stateBackend.registerKeySelectionListener(stub -> callback.run());
}
return callback;
}

实际清理的代码则位于TtlIncrementalCleanup类中,stateIterator就是状态数据的迭代器。

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
void stateAccessed() {
initIteratorIfNot();
try {
runCleanup();
} catch (Throwable t) {
throw new FlinkRuntimeException("Failed to incrementally clean up state with TTL", t);
}
}
private void initIteratorIfNot() {
if (stateIterator == null || !stateIterator.hasNext()) {
stateIterator = ttlState.original.getStateIncrementalVisitor(cleanupSize);
}
}
private void runCleanup() {
int entryNum = 0;
Collection<StateEntry<K, N, S>> nextEntries;
while ( entryNum < cleanupSize &&
stateIterator.hasNext() && !(nextEntries = stateIterator.nextEntries()).isEmpty()) {
for (StateEntry<K, N, S> state : nextEntries) {
S cleanState = ttlState.getUnexpiredOrNull(state.getState());
if (cleanState == null) {
stateIterator.remove(state);
} else if (cleanState != state.getState()) {
stateIterator.update(state, cleanState);
}
}
entryNum += nextEntries.size();
}
}

RocksDB压缩过滤清理策略

如果启用了该策略,Flink会通过维护一个RocksDbTtlCompactFiltersManager实例来管理FlinkCompactionFilter过滤器。FlinkCompactionFilter并不是在Flink工程中维护的,而是位于Data Artisans为Flink专门维护的FRocksDB库内。FLINK-10471实现了FlinkCompactionFilter及其附属逻辑,主要为C++代码,通过JNI调用。对应的commit详见GitHub,这里就不班门弄斧了。关于RocksDB的compaction相关细节,笔者之前也写过一篇长文做了些分析。