0%

Flink SQL知其然知其所以然: Flink SQL + calcite

先了解整个流程,有了全局视角之后,后续会详述细节。

先来看看 flink datastream 任务的执行过程:

  • DataStream
    使用时要在 flink datastream api 提供的各种 udf(比如 flatmapkeyedProcessFunction 等)中自定义处理逻辑,具体的业务执行逻辑都是敲代码、 java 文件写的,然后编译在 jvm 中执行,就和一个普通的 main 函数应用一模一样的流程。因为代码执行逻辑都是自己写的,所以这一部分相对好理解。

  • SQL
    Java 编译器不能识别和编译一条 SQL 进行执行,那么一条 SQL 是咋执行的呢?

2.1.先发挥自己的想象力

我们逆向思维进行考虑,如果想让一条 Flink SQL 按照我们的预期在 jvm 中执行,需要哪些过程。

  1. 整体来说:参考 datastream,如果 jvm 能执行 datastream java code 编译后的 class 文件,那么加一个 sql 解析层,能将 sql 逻辑解析为 datastream 的各种算子,然后编译执行不就 vans 了。
  2. sql parser:首先得有一个 sql parser 吧,得先能识别 sql 语法,将 sql 语法转化为 AST、具体的关系代数。
  3. 关系代数到 datastream 算子的映射:sql 逻辑解析为 datastream,需要有一个解析的映射逻辑吧。sql 是基于关系代数的,可以维护一个 sql 中的每个关系代数到具体 datastream 接口的映射关系,有了这些映射关系我们就可以将 sql 映射成一段可执行的 datastream 代码。举个例子:其可以将:
  • sql select xxx 解析为类似 datastream 中的 map
  • where xxx 解析为 filter
  • group by 解析成 keyby
  • sum(xx),count(xxx)可以解析为 datastream 中的 aggregate function
  • etc…
  • 代码生成:有了 sql AST,sql 到 datasretam 算子的映射关系之后,就要进行具体的代码生成了。比如去解析 sql AST 中具体哪些字段用作 where 逻辑,哪些字段用作 group by,都需要生成对应具体的 datastream 代码。
  • 运行:经过上述流程之后,就可以将一个 sql 翻译成一个 datastream 作业了,happy 的执行。

如下图所示,描绘了上述逻辑:

img

12

那么这个和 flink 实际实现有啥异同呢?

flink 大致是这样做的,虽在 flink 本身的中间还有一些其他的流程,后来的版本也不是基于 datastream,但是整体的处理逻辑还是和上述一致的。

所以不了解整体流程的同学可以先按照上述流程进行理解。

按照 博主的脑洞 来总结一条 sql 的使命就是:sql -> AST -> codegen(java code) -> 让我们 run 起来好吗

img

26

上面手绘可能看不清,下面这张图更清楚。

img

28

标准的一条 flink sql 运行起来的流程如下:

Notes:刚开始对其中的 SqlNode,RelNode 概念可能比较模糊。先理解整个流程,后续会详细介绍这些概念。

  1. sql 解析阶段:calcite parser 解析(sql -> AST,AST 即 SqlNode Tree)

  2. SqlNode 验证阶段:calcite validator 校验(SqlNode -> SqlNode,语法、表达式、表信息)

  3. 语义分析阶段:SqlNode 转换为 RelNode,RelNode 即 Logical Plan(SqlNode -> RelNode)

  4. 优化阶段:calcite optimizer 优化(RelNode -> RelNode,剪枝、谓词下推等)

  5. 物理计划生成阶段:Logical Plan 转换为 Physical Plan(等同于 RelNode 转换成 DataSet\DataStream API)

  6. 后续的运行逻辑与 datastream 一致

可以发现 flink 的实现博主的脑洞 整体主要框架上面是一致的。多出来的部分主要是 SqlNode 验证阶段优化阶段

大致了解了 一条 flink sql 的运行流程 之后,我们来看看 calcite 这玩意到底在 flink 里干了些啥。

根据上文总结来说 calcite 在 flink sql 中担当了 sql 解析、验证、优化功能。

img

30

看着 calcite 干了这么多事,那 calcite 是个啥东东,它的定位是啥?

3.1.calcite 是啥?

calcite 是一个动态数据的管理框架,它可以用来构建数据库系统的不同的解析的模块,但是它不包含数据存储数据处理等功能。

calcite 的目标是一种方案,适应所有的需求场景,希望能为不同计算平台和数据源提供统一的 sql 解析引擎,但是它只是提供查询引擎,而没有真正的去存储这些数据。

img

61

下图是目前使用了 calcite 能力的其他组件,也可见官网 https://calcite.apache.org/docs/powered_by.html

img

4

简单来说的话,可以先理解为 calcite 具有这几个功能(当然还有其他很牛逼的功能,感兴趣可以自查官网)。

  1. 自定义 sql 解析器:比如说我们新发明了一个引擎,然后我们要在这个引擎上来创造一套基于 sql 的接口,那么我们就可以使用直接 calcite,不用自己去写一套专门的 sql 的解析器,以及执行以及优化引擎,calcite 人都有。
  2. sql parser(extends SqlAbstractParserImpl):将 sql 的各种关系代数解析为具体的 AST,这些 AST 都能对应到具体的 java model,在 java 的世界里面,对象很重要,有了这些对象(SqlSelectSqlNode),就可以根据这些对象做具体逻辑处理了。举个例子,如下图,一条简单的 select c,d from source where a = '6' sql,经过 calcite 的解析之后,就可以得到 AST model(SqlNode)。可以看到有 SqlSelectSqlIdentifierSqlIdentifierSqlCharStringLiteral
  3. sql validator(extends SqlValidatorImpl):根据语法、表达式、表信息进行 SqlNode 正确性校验。
  4. sql optimizer:剪枝、谓词下推等优化

上面的这些能力整体组成如下图所示:

img

29

实际使用 calcite 解析一条 sql,跑起来看看。

img

2

  1. 不用重复造轮子。有限的精力应该放在有价值的事情上。
  2. calcite 有针对 stream 表的解决方案。具体可见 https://calcite.apache.org/docs/stream.html

4.案例篇-calcite 的能力、案例

4.1.先用用 calcite

1
重中之重,在了解原理之前,先跑起来是王道,也会帮助我们逐步理解。

官网已经有一个 csv 的案例了。感兴趣的可以直达 https://calcite.apache.org/docs/tutorial.html

跑完一个 csv demo,在详细了解 calcite 之前还需要了解下 sql,calcite 的支柱:关系代数。

4.2.关系代数

sql 是基于关系代数的查询语言,是关系代数在工程上的一种很好的实现方案。在工程中,关系代数难表达,但是 sql 就易于理解。关系代数和 sql 的关系如下。

  1. 可以将一条 sql 解析为一个关系代数表达式的组合。在 sql 中的操作都可以转化成关系代数的表达式。

  2. sql 的执行优化(所有的优化的前提都是优化前和优化后最终执行结果相同,即等价交换)是基于关系代数运算的。

4.2.1.常用关系代数

总结下,有哪些常用的关系代数:

img

50

4.2.2.sql 优化支柱之关系代数等价变换

关系代数等价变换是 calcite optimizer 的基础理论。

下面是一些等价变换的例子。

1.连接(),笛卡尔积(×)的交换律

img

51

2.连接(),笛卡尔积(×)的结合律

img

3.投影(Π)的串接定律

img

4.选择(σ)的串接定律

img

5.选择(σ)与投影(Π)的交换

img

6.选择(σ)与笛卡尔积(×)的交换

img

7.选择(σ)与并()的交换

img

8.选择(σ)与差(-)的交换

img

9.投影(Π)与笛卡尔积(×)的交换

img

10.投影(Π)与并()的交换

img

然后看一个基于关系代数优化的实际 sql 案例:

有三个关系 A(a1,a2,a3,…)B(b1,b2,b3, … )C(a1,b1,c1,c2, … )

有一个查询请求如下:

1
SELECT A.a1 FROM A,B,C WHERE A.a1 = C.a1 AND B.b1 = C.b1 AND f(c1)

1.首先将 sql 转为关系代数的语法树。

img

36

2.优化:选择(σ)的串接定律。

img

47

img

37

3.优化:选择(σ)与笛卡尔积(×)的交换。

img

48

img

38

4.优化:投影(π)与笛卡尔积(×)的交换。

img

49

img

img

img

img

img

img

img

关于关系代数我们就有了大致的了解。

除此之外,对于更深入了解 flink sql,calcite 而言,我们还需要了解一下在 calcite 代码体系中有哪些重要 model。

4.3.calcite 必知的基础 model

calcite 中有两个最最基础、重要的 model 在我们理解 flink sql 解析流程时需要知道的。

  1. SqlNode:sql 转化而成,可以理解为直观表达 sql 层次结构的的 model
  2. RelNode:SqlNode 转化而成,可以理解为将 SqlNode 转化为关系代数,表达关系代数层次结构的 model

举个例子来说明下,下面这条 flink sql,经过解析之后的 SqlNodeRelNode 如下图:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
SELECT
sum(part_pv) as pv,
window_start
FROM (
SELECT
count(1) as part_pv,
cast(tumble_start(rowtime, INTERVAL '60' SECOND) as bigint) * 1000 as window_start
FROM
source_db.source_table
GROUP BY
tumble(rowtime, INTERVAL '60' SECOND)
, mod(id, 1024)
)
GROUP BY
window_start

img

62

可以看到 SqlNode 包含的内容是 sql 的层次结构,包括 selectListfromwheregroup by 等。

RelNode 包含的是关系代数的层次结构,每一层都有一个 input 来承接。结合上面优化案例的树状结构一样。

img

63

img

29

如上图所示,此处我们结合上节介绍的 calcite 的 model,以及 flink sql 的实现来走一遍其处理流程:

  1. sql 解析阶段(sql –> SqlNode)
  2. SqlNode 验证(SqlNode –> SqlNode)
  3. 语义分析(SqlNode –> RelNode)
  4. 优化阶段(RelNode –> RelNode)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
SELECT
sum(part_pv) as pv,
window_start
FROM (
SELECT
count(1) as part_pv,
cast(tumble_start(rowtime, INTERVAL '60' SECOND) as bigint) * 1000 as window_start
FROM
source_db.source_table
GROUP BY
tumble(rowtime, INTERVAL '60' SECOND)
, mod(id, 1024)
)
GROUP BY
window_start

其中前三步解析和转化,都在 在执行 TableEnvironment#sqlQuery 进行。

最后一步优化,在执行 sink 操作时进行,即在这个例子中是 tEnv.toRetractStream(result, Row.class)

源码公众号后台回复flink sql 知其所以然(六)| flink sql 约会 calcite获取。

4.4.2.sql 解析阶段(sql –> SqlNode)

sql 解析阶段使用 Sql Parser 将 sql 解析为 SqlNode。这一步在执行 TableEnvironment#sqlQuery 进行。

img

img

img

img

可以从上图看到 flink sql 具体实现类是 FlinkSqlParserImpl

img

68

具体 parse 得到 SqlNode 如上图。

4.4.3.SqlNode 验证(SqlNode –> SqlNode)

上面的第一步生产的 SqlNode 对象是一个未经验证的,这一步就是语法检查阶段,语法检查前需要知道元数据信息,这个检查会包括表名、字段名、函数名、数据类型的检查。进行语法检查的实现如下:

img

img

img

可以从上图看到 flink sql 校验器的具体实现类是 FlinkCalciteSqlValidator,其中包含了元数据信息,从而可以进行元数据信息检查。

4.4.4.语义分析(SqlNode –> RelNode)

这一步就是将 SqlNode 转换成 RelNode,也就是生成相应的关系代数层面的逻辑(这里一般都叫做逻辑计划:Logical Plan)。

img

img

img

4.4.5.优化阶段(RelNode –> RelNode)

这一步就是优化阶段。详细内容可以自己 debug 代码查看,此处不赘述。

img

img

4.5.calcite 怎么做到这么通用?

此处以 calcite parser 举例说明,其模块为什么这通用?其他的模块都是类似的方式。

先说结论:因为 calcite parser 模块提供了接口,具体的 parse 逻辑、规则是可以根据用户自定义进行配置的。大家可以看下图,博主画出了一张图进行详述。

img

5

如上图,引擎 sql 解析器的生成是有一个输入的,就是 用户自定义语法分析规则变量,具体引擎的 sql 解析器其实也是根据用户自定义的 解析规则 去生成的 解析器。其 解析器 的动态生成依赖 javacc 这样的组件。calcite 提供的是统一的 sql AST 模型、优化模型接口等,而具体的解析实现交给了用户自己去决定。

javacc 会根据 calcite 中定义的 Parser.jj 文件,生成具体的 sql parser 代码(如上图),这个 sql parser 的能力就是将 sql 转换成 AST (SqlNode)。关于 calcite 能力的更详细内容见 https://matt33.com/2019/03/07/apache-calcite-process-flow/

上图涉及到的文件大家可以下载 calcite 源码 https://github.com/apache/calcite.git 之后,切换到 coremodule 之后查看。

img

31

4.5.1.javacc 是啥?

javacc 是一个用 java 开发的最受欢迎的语法分析生成器。这个分析生成器工具可以读取上下文无关且有着特殊意义的语法并把它转换成可以识别且匹配该语法的 java 程序。它是 100% 的纯 java 代码,可以在多种平台上运行。

简单解释 javacc 就是它是一个通用的语法分析生产器,用户可以使用 javacc 任意定义一套 DSL 及解析器。

举个例子,如果哪天你觉得 sql 也不够简洁通用,你可以使用 javacc 自己定义一套更简洁的 user-define-ql。然后使用 javacc 作为你的 user-define-ql 的解析器。是不是很流批,可以自己去搞编译器了。

4.5.2.跑跑 javacc

这里不介绍具体的 javacc 语法,直接以官网的 Simple1.jj 为案例。详细语法和功能可以参考官网(https://javacc.github.io/javacc/) 或者一下博客。

  1. https://www.cnblogs.com/Gavin_Liu/archive/2009/03/07/1405029.html
  2. https://www.yangguo.info/2014/12/13/%E7%BC%96%E8%AF%91%E5%8E%9F%E7%90%86-Javacc%E4%BD%BF%E7%94%A8/
  3. https://www.engr.mun.ca/~theo/JavaCC-Tutorial/javacc-tutorial.pdf

Simple1.jj 是用于识别一系列的 {相同数量的花括号},之后跟着 0 个或多个行终结符。

img

7

下面是合法的字符串例子:

{}{{{{{}}}}},etc.

下面是不合法的字符串例子:

{{{{{}{}{}}{{}{}},etc.

接下来让我们实际将 Simple1.jj 编译生成具体的规则代码。

在 pom 中加入 javacc build 插件:

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
<plugin>
<!-- This must be run AFTER the fmpp-maven-plugin -->
<groupId>org.codehaus.mojo</groupId>
<artifactId>javacc-maven-plugin</artifactId>
<version>2.4</version>
<executions>
<execution>
<phase>generate-sources</phase>
<id>javacc</id>
<goals>
<goal>javacc</goal>
</goals>
<configuration>
<sourceDirectory>${project.build.directory}/generated-sources/</sourceDirectory>
<includes>
<include>**/Simple1.jj</include>
</includes>
<!-- This must be kept synced with Apache Calcite. -->
<lookAhead>1</lookAhead>
<isStatic>false</isStatic>
<outputDirectory>${project.build.directory}/generated-sources/</outputDirectory>
</configuration>
</execution>
</executions>
</plugin>

在 compile 之后,就会在 generated-sources 下生成代码:

img

8

然后把代码 copy 到 Sources 路径下:

img

33

执行下代码,可以看到 {}{{}} 都可以校验通过,一旦出现不符合规则的 {{ 输入,就会抛出异常。

img

img

img

关于 javacc 基本上就了解个大概了。

感兴趣的可以尝试自定义一个编译器。

4.5.3.fmpp 是啥?

img

5

fmpp 就是一个基于 freemarker 的模板生产器。用户可以统一管理自己的变量,然后用 ftl 模板 + 变量 生成对应的最终文件。在 calcite 中使用 fmpp 作为变量 + 模板的统一管理器。然后基于 fmpp 来生成对应的 Parser.jj 文件。

博主画了一张图,包含了其中重要组件之间的依赖关系。

img

3

你没猜错,还是上面那些流程,fmpp(Parser.jj 模板生成) -> javacc(Parser 生成) -> calcite

在介绍 Parser 生成流程之前,先看看 flink 最终生成的 Parser:FlinkSqlParserImpl (此处使用 Blink Planner)。

5.1.FlinkSqlParserImpl

以下面这个案例出发(代码基于 flink 1.13.1 版本):

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
public class ParserTest {

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

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

env.setParallelism(10);

EnvironmentSettings settings = EnvironmentSettings
.newInstance()
.useBlinkPlanner()
.inStreamingMode()
.build();

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

DataStream<Tuple3<String, Long, Long>> tuple3DataStream =
env.fromCollection(Arrays.asList(
Tuple3.of("2", 1L, 1627254000000L),
Tuple3.of("2", 1L, 1627218000000L + 5000L),
Tuple3.of("2", 101L, 1627218000000L + 6000L),
Tuple3.of("2", 201L, 1627218000000L + 7000L),
Tuple3.of("2", 301L, 1627218000000L + 7000L),
Tuple3.of("2", 301L, 1627218000000L + 7000L),
Tuple3.of("2", 301L, 1627218000000L + 7000L),
Tuple3.of("2", 301L, 1627218000000L + 7000L),
Tuple3.of("2", 301L, 1627218000000L + 7000L),
Tuple3.of("2", 301L, 1627218000000L + 86400000 + 7000L)))
.assignTimestampsAndWatermarks(
new BoundedOutOfOrdernessTimestampExtractor<Tuple3<String, Long, Long>>(Time.seconds(0L)) {
@Override
public long extractTimestamp(Tuple3<String, Long, Long> element) {
return element.f2;
}
});

tEnv.registerFunction("mod", new Mod_UDF());

tEnv.registerFunction("status_mapper", new StatusMapper_UDF());

tEnv.createTemporaryView("source_db.source_table", tuple3DataStream,
"status, id, timestamp, rowtime.rowtime");

String sql = "SELECT\n"
+ " count(1),\n"
+ " cast(tumble_start(rowtime, INTERVAL '1' DAY) as string)\n"
+ "FROM\n"
+ " source_db.source_table\n"
+ "GROUP BY\n"
+ " tumble(rowtime, INTERVAL '1' DAY)";

Table result = tEnv.sqlQuery(sql);

tEnv.toAppendStream(result, Row.class).print();

env.execute();

}

}

debug 过程如之前分析 sql -> SqlNode 过程所示,如下图直接定位到 SqlParser:

img

21

如上图可以看到具体的 Parser 就是 FlinkSqlParserImpl

定位到具体的代码如下图所示(flink-table-palnner-blink-2.11-1.13.1.jar)。

img

34

最终 parse 的结果 SqlNode 如下图。

img

22

img

img

再来看看 FlinkSqlParserImpl 是怎么使用 calcite 生成的。

具体到 flink 中的实现,位于源码中的 flink-table.flink-sql-parser 模块(源码基于 flink 1.13.1)。

flink 是依赖 maven 插件实现的上面的整体流程。

5.2.FlinkSqlParserImpl 的生成

img

14

接下来看看整个 Parser 生成流程。

使用 maven-dependency-plugin 将 calcite 解压到 flink 项目 build 目录下。

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
<plugin>
<!-- Extract parser grammar template from calcite-core.jar and put
it under ${project.build.directory} where all freemarker templates are. -->
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-dependency-plugin</artifactId>
<executions>
<execution>
<id>unpack-parser-template</id>
<phase>initialize</phase>
<goals>
<goal>unpack</goal>
</goals>
<configuration>
<artifactItems>
<artifactItem>
<groupId>org.apache.calcite</groupId>
<artifactId>calcite-core</artifactId>
<type>jar</type>
<overWrite>true</overWrite>
<outputDirectory>${project.build.directory}/</outputDirectory>
<includes>**/Parser.jj</includes>
</artifactItem>
</artifactItems>
</configuration>
</execution>
</executions>
</plugin>

img

15

5.2.2.fmpp 生成 Parser.jj

使用 maven-resources-pluginParser.jj 代码生成。

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
<plugin>
<artifactId>maven-resources-plugin</artifactId>
<executions>
<execution>
<id>copy-fmpp-resources</id>
<phase>initialize</phase>
<goals>
<goal>copy-resources</goal>
</goals>
<configuration>
<outputDirectory>${project.build.directory}/codegen</outputDirectory>
<resources>
<resource>
<directory>src/main/codegen</directory>
<filtering>false</filtering>
</resource>
</resources>
</configuration>
</execution>
</executions>
</plugin>
<plugin>
<groupId>com.googlecode.fmpp-maven-plugin</groupId>
<artifactId>fmpp-maven-plugin</artifactId>
<version>1.0</version>
<dependencies>
<dependency>
<groupId>org.freemarker</groupId>
<artifactId>freemarker</artifactId>
<version>2.3.28</version>
</dependency>
</dependencies>
<executions>
<execution>
<id>generate-fmpp-sources</id>
<phase>generate-sources</phase>
<goals>
<goal>generate</goal>
</goals>
<configuration>
<cfgFile>${project.build.directory}/codegen/config.fmpp</cfgFile>
<outputDirectory>target/generated-sources</outputDirectory>
<templateDirectory>${project.build.directory}/codegen/templates</templateDirectory>
</configuration>
</execution>
</executions>
</plugin>

img

16

5.2.3.javacc 生成 parser

使用 javacc 将根据 Parser.jj 文件生成 Parser。

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
<plugin>
<!-- This must be run AFTER the fmpp-maven-plugin -->
<groupId>org.codehaus.mojo</groupId>
<artifactId>javacc-maven-plugin</artifactId>
<version>2.4</version>
<executions>
<execution>
<phase>generate-sources</phase>
<id>javacc</id>
<goals>
<goal>javacc</goal>
</goals>
<configuration>
<sourceDirectory>${project.build.directory}/generated-sources/</sourceDirectory>
<includes>
<include>**/Parser.jj</include>
</includes>
<!-- This must be kept synced with Apache Calcite. -->
<lookAhead>1</lookAhead>
<isStatic>false</isStatic>
<outputDirectory>${project.build.directory}/generated-sources/</outputDirectory>
</configuration>
</execution>
</executions>
</plugin>

img

17

5.2.4.看看 Parser

最终生成的 Parser 就是 FlinkSqlParserImpl

img

18

blink planner(flink-table-planner-blink) 在打包时将 flink-sql-parserflink-sql-parser-hive 打包进去。

img

35

从TimeoutException看Flink的心跳机制

0x00 摘要

本文从一个调试时候常见的异常 “TimeoutException: Heartbeat of TaskManager timed out”切入,为大家剖析Flink的心跳机制。文中代码基于Flink 1.10。

0x01 缘由

大家如果经常调试Flink,当进入断点看到了堆栈和变量内容之后,你容易陷入了沉思。当你发现了问题可能所在,高兴的让程序Resume的时候,你发现程序无法运行,有如下提示:

1
Caused by: java.util.concurrent.TimeoutException: Heartbeat of TaskManager with id 93aa1740-cd2c-4032-b74a-5f256edb3217 timed out.

这实在是很郁闷的事情。作为程序猿不能忍啊,既然异常提示中有 Heartbeat 字样,于是我们就来一起看看Flink的心跳机制,看看有没有可以修改的途径。

0x02 背景概念

2.1 四大模块

Flink有核心四大组件:Dispatcher,JobMaster,ResourceManager,TaskExecutor。

  • Dispatcher(Application Master)用于接收client提交的任务和启动相应的JobManager。其提供REST接口来接收client的application提交,负责启动JM和提交application,同时运行Web UI。

  • ResourceManager:主要用于资源的申请和分配。当TM有空闲的slot就会告诉JM,没有足够的slot也会启动新的TM。kill掉长时间空闲的TM。

  • JobMaster

    :功能主要包括(旧版本中JobManager的功能在新版本中以JobMaster形式出现,可能本文中会混淆这两个词,请大家谅解):

    • 将JobGraph转化为ExecutionGraph(physical dataflow graph,并行化)。
    • 向RM申请资源、schedule tasks、保存作业的元数据。
  • TaskManager:类似Spark的executor,会跑多个线程的task、数据缓存与交换。Flink 架构遵循 Master - Slave 架构设计原则,JobMaster 为 Master 节点,TaskManager 为Slave节点。

这四大组件彼此之间的通信需要依赖RPC实现。

2.2 Akka

Flink底层RPC基于Akka实现。Akka是一个开发并发、容错和可伸缩应用的框架。它是Actor Model的一个实现,和Erlang的并发模型很像。在Actor模型中,所有的实体被认为是独立的actors。actors和其他actors通过发送异步消息通信。

Actor模型的强大来自于异步。它也可以显式等待响应,这使得可以执行同步操作。但是强烈不建议同步消息,因为它们限制了系统的伸缩性。

2.3 RPC机制

RPC作用是:让异步调用看起来像同步调用。

Flink基于Akka构建了其底层通信系统,引入了RPC调用,各节点通过GateWay方式回调,隐藏通信组件的细节,实现解耦。Flink整个通信框架的组件主要由RpcEndpoint、RpcService、RpcServer、AkkaInvocationHandler、AkkaRpcActor等构成。

RPC相关的主要接口如下:

  • RpcEndpoint
  • RpcService
  • RpcGateway

2.3.1 RpcEndpoint:RPC的基类

RpcEndpoint是Flink RPC终端的基类,所有提供远程过程调用的分布式组件必须扩展RpcEndpoint,其功能由RpcService支持。

RpcEndpoint的子类只有四类组件:Dispatcher,JobMaster,ResourceManager,TaskExecutor,即Flink中只有这四个组件有RPC的能力,换句话说只有这四个组件有RPC的这个需求。

每个RpcEndpoint对应了一个路径(endpointId和actorSystem共同确定),每个路径对应一个Actor,其实现了RpcGateway接口,

RpcService:RPC服务提供者

RpcServer是RpcEndpoint的成员变量,为RpcService提供RPC服务/连接远程Server,其只有一个子类实现:AkkaRpcService(可见目前Flink的通信方式依然是Akka)。

RpcServer用于启动和连接到RpcEndpoint, 连接到rpc服务器将返回一个RpcGateway,可用于调用远程过程。

Flink四大组件Dispatcher,JobMaster,ResourceManager,TaskExecutor,都是RpcEndpoint的实现,所以构建四大组件时,同步需要初始化RpcServer。如JobManager的构造方式,第一个参数就是需要知道RpcService。

RpcGateway:RPC调用的网关

Flink的RPC协议通过RpcGateway来定义;由前面可知,若想与远端Actor通信,则必须提供地址(ip和port),如在Flink-on-Yarn模式下,JobMaster会先启动ActorSystem,此时TaskExecutor的Container还未分配,后面与TaskExecutor通信时,必须让其提供对应地址。

Dispatcher,JobMaster,ResourceManager,TaskExecutor 这四大组件通过各种方式实现了Gateway。以JobMaster为例,JobMaster实现JobMasterGateway接口。各组件类的成员变量都有需要通信的其他组件的GateWay实现类,这样可通过各自的Gateway实现RPC调用。

2.4 常见心跳机制

常见的心跳检测有两种:

  • socket 套接字SO_KEEPALIVE本身带有的心跳机制,定期向对方发送心跳包,对方收到心跳包后会自动回复;
  • 应用自身实现心跳机制,同样也是使用定期发送请求的方式;

Flink实现的是第二种方案。

0x03 Flink心跳机制

3.1 代码和机制

Flink的心跳机制代码在:

1
Flink-master/flink-runtime/src/main/java/org/apache/flink/runtime/heartbeat

四个接口:

1
HeartbeatListener.java          HeartbeatManager.java      HeartbeatTarget.java  HeartbeatMonitor.java

以及如下几个类:

1
2
HeartbeatManagerImpl.java   HeartbeatManagerSenderImpl.java   HeartbeatMonitorImpl.java
HeartbeatServices.java NoOpHeartbeatManager.java

Flink集群有多种业务流程,比如Resource Manager, Task Manager, Job Manager。每种业务流程都有自己的心跳机制。Flink的心跳机制只是提供接口和基本功能,具体业务功能由各业务流程自己实现。

我们首先设定 心跳系统中有两种节点:sender和receiver。心跳机制是sender和receivers彼此相互检测。但是检测动作是Sender主动发起,即Sender主动发送请求探测receiver是否存活,因为Sender已经发送过来了探测心跳请求,所以这样receiver同时也知道Sender是存活的,然后Reciver给Sender回应一个心跳表示自己也是活着的。

因为Flink的几个名词和我们常见概念有所差别,所以流程上需要大家仔细甄别,即:

  • Flink Sender 主动发送Request请求给Receiver,要求Receiver回应一个心跳;
  • Flink Receiver 收到Request之后,通过Receive函数回应一个心跳请求给Sender;

3.2 静态架构

3.2.1 HeartbeatTarget :监控目标抽象

HeartbeatTarget是对监控目标的抽象。心跳机制在行为上而言有两种动作:

  • 向某个节点发送请求。
  • 处理某个节点发来的请求。

HeartbeatTarget的函数就是这两个动作:

  • receiveHeartbeat :向某个节点(Sender)发送心跳回应,其参数heartbeatOrigin 就是 Receiver。
  • requestHeartbeat :向某个节点(Receiver)要求其回应一个心跳,其参数requestOrigin 就是 Sender。requestHeartbeat这个函数是Sender的函数,其中Sender通过RPC直接调用到Receiver。

这两个函数的参数也很简单:分别是请求的发送放和接收方,还有Payload载荷。对于一个确定节点而言,接收的和发送的载荷是同一类型的。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
public interface HeartbeatTarget<I> {
/**
* Sends a heartbeat response to the target.
* @param heartbeatOrigin Resource ID identifying the machine for which a heartbeat shall be reported.
*/
// heartbeatOrigin 就是 Receiver
void receiveHeartbeat(ResourceID heartbeatOrigin, I heartbeatPayload);

/**
* Requests a heartbeat from the target.
* @param requestOrigin Resource ID identifying the machine issuing the heartbeat request.
*/
// requestOrigin 就是 Sender
void requestHeartbeat(ResourceID requestOrigin, I heartbeatPayload);
}

3.2.2 HeartbeatMonitor : 管理heartbeat target的心跳状态

对HeartbeatTarget的封装,这样Manager对Target的操作是通过对Monitor完成,后续会在其继承类中详细说明。

1
2
3
4
5
6
7
8
9
10
11
12
public interface HeartbeatMonitor<O> {
// Gets heartbeat target.
HeartbeatTarget<O> getHeartbeatTarget();
// Gets heartbeat target id.
ResourceID getHeartbeatTargetId();
// Report heartbeat from the monitored target.
void reportHeartbeat();
//Cancel this monitor.
void cancel();
//Gets the last heartbeat.
long getLastHeartbeat();
}

3.2.3 HeartbeatManager :心跳管理者

HeartbeatManager负责管理心跳机制,比如启动/停止/报告一个HeartbeatTarget。此接口继承HeartbeatTarget。

除了HeartbeatTarget的函数之外,这接口有4个函数:

  • monitorTarget,把和某资源对应的节点加入到心跳监控列表;
  • unmonitorTarget,从心跳监控列表删除某资源对应的节点;
  • stop,停止心跳管理服务,释放资源;
  • getLastHeartbeatFrom,获取某节点的最后一次心跳数据。
1
2
3
4
5
6
public interface HeartbeatManager<I, O> extends HeartbeatTarget<I> {
void monitorTarget(ResourceID resourceID, HeartbeatTarget<O> heartbeatTarget);
void unmonitorTarget(ResourceID resourceID);
void stop();
long getLastHeartbeatFrom(ResourceID resourceId);
}

3.2.4 HearbeatListener 处理心跳结果

用户业务逻辑需要继承这个接口以处理心跳结果。其可以看做服务的输出,实现了三个回调函数。

  • notifyHeartbeatTimeout,处理节点心跳超时
  • reportPayload,处理节点发来的Payload载荷
  • retrievePayLoad。获取对某节点发下一次心跳请求的Payload载荷
1
2
3
4
5
public interface HeartbeatListener<I, O> {
void notifyHeartbeatTimeout(ResourceID resourceID);
void reportPayload(ResourceID resourceID, I payload);
O retrievePayload(ResourceID resourceID);
}

3.3 动态运行机制

之前提到Sender和Receiver,下面两个类就对应上述概念。

  • HeartbeatManagerImpl :Receiver,存在于JobMaster与TaskExecutor中;
  • HeartbeatManagerSenderImpl :Sender,继承 HeartbeatManagerImpl类,用于周期发送心跳要求,存在于JobMaster、ResourceManager中。

几个关键问题:

  1. 如何判定心跳超时? 心跳服务启动后,Flink在Monitor中通过 ScheduledFuture 会启动一个线程来处理心跳超时事件。在设定的心跳超时时间到达后才执行线程。 如果在设定的心跳超时时间内接收到组件的心跳消息,会先将该线程取消而后重新开启,重置心跳超时事件的触发。 如果在设定的心跳超时时间内没有收到组件的心跳,则会通知组件:你超时了。
  2. 何时”调用双方”发起心跳检查? 心跳检查是双向的,一方(Sender)会主动发起心跳请求,而另一方(Receiver)则是对心跳做出响应,两者通过RPC相互调用,重置对方的 Monitor 超时线程。 以JobMaster和TaskManager为例,JM在启动时会开启周期调度,向已经注册到JM中的TM发起心跳检查,通过RPC调用TM的requestHeartbeat方法,重置TM中对JM超时线程的调用,表示当前JM状态正常。在TM的requestHeartbeat方法被调用后,通过RPC调用JM的receiveHeartbeat,重置 JM 中对TM超时线程的调用,表示TM状态正常。
  3. 如何处理心跳超时? 心跳服务依赖 HeartbeatListener,当在timeout时间范围内未接收到心跳响应,则会触发超时处理线程,该线程通过调用HeartbeatListener.notifyHeartbeatTimeout方法做后续重连操作或者直接断开。

下面是一个概要(以RM & TM为例):

  • RM : 实现了ResourceManagerGateway (可以直接被RPC调用)
  • TM : 实现了TaskExecutorGateway (可以直接被RPC调用)
  • RM :有一个Sender HM : taskManagerHeartbeatManager,Sender HM 拥有用户定义的 TaskManagerHeartbeatListener
  • TM :有一个Receiver HM :resourceManagerHeartbeatManager,Receiver HM 拥有用户定义的ResourceManagerHeartbeatListener。
  • HeartbeatManager 有一个ConcurrentHashMap<ResourceID, HeartbeatMonitor> heartbeatTargets,这个Map是它监控的所有Target。
  • 对于RM的每一个需要监控的TM, 其生成一个HeartbeatTarget,进而被构造成一个HeartbeatMonitor,放置到ResourceManager.taskManagerHeartbeatManager中。
  • 每一个Target对应的Monitor中,有自己的异步任务ScheduledFuture,这个ScheduledFuture不停的被取消/重新生成。如果在某个期间内没有被取消,则通知用户定义的listener出现了timeout。

3.3.1 HearbeatManagerImpl : Receiver

HearbeatManagerImpl是receiver的具体实现。它由 心跳 被发起方(就是Receiver,例如TM) 创建,接收 **发起方(就是Sender,例如 JM)**的心跳发送请求。心跳超时 会触发 heartbeatListener.notifyHeartbeatTimeout方法。

注意:被发起方监控线程(Monitor)的开启是在接收到请求心跳(requestHeartbeat被调用后)以后才触发的,属于被动触发。

HearbeatManagerImpl主要维护了

  • 一个心跳监控列表 map : <ResourceID, HeartbeatMonitor<O>> heartbeatTargets;。这是一个KV关联。 key代表要发送心跳组件(例如:TM)的ID,value则是为当前组件创建的触发心跳超时的线程HeartbeatMonitor,两者一一对应。 当一个从所联系的machine发过来的心跳被收到时候,对应的monitor的状态会被更新(重启一个新ScheduledFuture)。当一个monitor发现了一个 heartbeat timed out,它会通知自己的HeartbeatListener。
  • 一个 ScheduledExecutor mainThreadExecutor 负责heartbeat timeout notifications。
  • heartbeatListener :处理心跳结果。

HearbeatManagerImpl 数据结构如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
@ThreadSafe
public class HeartbeatManagerImpl<I, O> implements HeartbeatManager<I, O> {

/** Heartbeat timeout interval in milli seconds. */
private final long heartbeatTimeoutIntervalMs;

/** Resource ID which is used to mark one own's heartbeat signals. */
private final ResourceID ownResourceID;

/** Heartbeat listener with which the heartbeat manager has been associated. */
private final HeartbeatListener<I, O> heartbeatListener;

/** Executor service used to run heartbeat timeout notifications. */
private final ScheduledExecutor mainThreadExecutor;

/** Map containing the heartbeat monitors associated with the respective resource ID. */
private final ConcurrentHashMap<ResourceID, HeartbeatMonitor<O>> heartbeatTargets;

/** Running state of the heartbeat manager. */
protected volatile boolean stopped;
}

HearbeatManagerImpl实现的主要函数有:

  • monitorTarget :把一个节点加入到心跳监控列表。
    • 传入参数有:ResourceId和HearbeatTarget,monitorTarget根据这两个参数,生成一个HeartbeatMonitor对象,然后把这个对象跟ResrouceId做kv关联,存入到heartbeatTargets。 一个节点可能参与多个业务流程,因此一个节点参与多个心跳流程,一个节点上运行多个不同类型的HearbeatTarget。所以一个ResourceID可能会跟不同类型的HearbeatTarget对象关联,分别加入到多个HeartbeatManager,进行不同类型的心跳监控。也因此这个函数入参是两个参数。
  • requestHeartbeat :Sender通过RPC异步调用到Receiver的这个函数 以要求receiver向requestOrigin节点(就是Sender)发起一次心跳响应,载荷是heartbeatPayLoad。其内部流程如下:
    • 首先会调用reportHeartbeat函数,作用是 通过Monitor 记录发起请求的这个时间点,然后创建一个ScheduleFuture。如果到期后,requestOrigin没有作出响应,那么就将requestOrigin节点对应的HeartbeatMonitor的state设置成TIMEOUT状态,如果到期内requestOrigin响应了,ScheduleFuture会被取消,HeartbeatMonitor的state仍然是RUNNING。
    • 其次调用reportPayload函数,把requestOrigin节点的最新的heartbeatPayload通知给heartbeatListener。heartbeatListener是外部传入的,它根据所有节点的心跳记录做监听管理。
    • 最后调用receiveHearbeat函数,响应一个心跳给Sender。

3.3.2 HeartbeatManagerSenderImpl : Sender

继承HearbeatManagerImpl,由**心跳管理的一方(例如JM)**创建,实现了run函数(即它可以作为一个单独线程运行),创建后立即开启周期调度线程,每次遍历自己管理的heartbeatTarget,触发heartbeatTarget.requestHeartbeat,要求 Target 返回一个心跳响应。属于主动触发心跳请求。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
public class HeartbeatManagerSenderImpl<I, O> extends HeartbeatManagerImpl<I, O> implements Runnable {
public void run() {
if (!stopped) {
for (HeartbeatMonitor<O> heartbeatMonitor : getHeartbeatTargets().values()) {
requestHeartbeat(heartbeatMonitor);
}
// 周期调度
getMainThreadExecutor().schedule(this, heartbeatPeriod, TimeUnit.MILLISECONDS);
}
}

// 主动发起心跳检查
private void requestHeartbeat(HeartbeatMonitor<O> heartbeatMonitor) {
O payload = getHeartbeatListener().retrievePayload(heartbeatMonitor.getHeartbeatTargetId());
final HeartbeatTarget<O> heartbeatTarget = heartbeatMonitor.getHeartbeatTarget();
// 调用 Target 的 requestHeartbeat 函数
heartbeatTarget.requestHeartbeat(getOwnResourceID(), payload);
}
}

3.3.3 HeartbeatMonitorImpl

Heartbeat monitor管理心跳目标,它启动一个ScheduledExecutor。

  • 如果在timeout时间内没有接收到心跳信号,则判定心跳超时,通知给HeartbeatListener。
  • 如果在timeout时间内接收到心跳信号,则重置当前ScheduledExecutor。
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
public class HeartbeatMonitorImpl<O> implements HeartbeatMonitor<O>, Runnable {

/** Resource ID of the monitored heartbeat target. */
private final ResourceID resourceID; // 被监控的resource ID

/** Associated heartbeat target. */
private final HeartbeatTarget<O> heartbeatTarget; //心跳目标

private final ScheduledExecutor scheduledExecutor;

/** Listener which is notified about heartbeat timeouts. */
private final HeartbeatListener<?, ?> heartbeatListener; // 心跳监听器

/** Maximum heartbeat timeout interval. */
private final long heartbeatTimeoutIntervalMs;

private volatile ScheduledFuture<?> futureTimeout;
// AtomicReference 使用
private final AtomicReference<State> state = new AtomicReference<>(State.RUNNING);
// 最近一次接收到心跳的时间
private volatile long lastHeartbeat;

// 报告心跳
public void reportHeartbeat() {
// 保留最近一次接收心跳时间
lastHeartbeat = System.currentTimeMillis();
// 接收心跳后,重置timeout线程
resetHeartbeatTimeout(heartbeatTimeoutIntervalMs);
}

// 心跳超时,触发lister的notifyHeartbeatTimeout
public void run() {
// The heartbeat has timed out if we're in state running
if (state.compareAndSet(State.RUNNING, State.TIMEOUT)) {
heartbeatListener.notifyHeartbeatTimeout(resourceID);
}
}

// 重置TIMEOUT
void resetHeartbeatTimeout(long heartbeatTimeout) {
if (state.get() == State.RUNNING) {
//先取消线程,在重新开启
cancelTimeout();
// 启动超时线程
futureTimeout = scheduledExecutor.schedule(this, heartbeatTimeout, TimeUnit.MILLISECONDS);

// Double check for concurrent accesses (e.g. a firing of the scheduled future)
if (state.get() != State.RUNNING) {
cancelTimeout();
}
}
}

3.3.3 HeartbeatServices

建立heartbeat receivers and heartbeat senders,主要是对外提供服务。这里我们可以看到:

  • HeartbeatManagerImpl就是receivers。
  • HeartbeatManagerSenderImpl就是senders。
1
2
3
4
5
6
7
8
9
10
public class HeartbeatServices {
// Creates a heartbeat manager which does not actively send heartbeats.
public <I, O> HeartbeatManager<I, O> createHeartbeatManager(...) {
return new HeartbeatManagerImpl<>(...);
}
// Creates a heartbeat manager which actively sends heartbeats to monitoring targets.
public <I, O> HeartbeatManager<I, O> createHeartbeatManagerSender(...) {
return new HeartbeatManagerSenderImpl<>(...);
}
}

0x04 初始化

4.1 心跳服务创建

心跳管理服务在Cluster入口创建。因为我们是调试,所以在MiniCluster.start调用。

1
2
3
4
5
public void start() throws Exception {
......
heartbeatServices = HeartbeatServices.fromConfiguration(configuration);
......
}

HeartbeatServices.fromConfiguration会从Configuration中获取配置信息:

  • 心跳间隔 heartbeat.interval
  • 心跳超时时间 heartbeat.timeout

这个就是我们解决最开始问题的思路:从配置信息入手,扩大心跳间隔。

1
2
3
4
5
6
7
8
9
10
11
public HeartbeatServices(long heartbeatInterval, long heartbeatTimeout) {
this.heartbeatInterval = heartbeatInterval;
this.heartbeatTimeout = heartbeatTimeout;
}

public static HeartbeatServices fromConfiguration(Configuration configuration) {
long heartbeatInterval = configuration.getLong(HeartbeatManagerOptions.HEARTBEAT_INTERVAL);
long heartbeatTimeout = configuration.getLong(HeartbeatManagerOptions.HEARTBEAT_TIMEOUT);

return new HeartbeatServices(heartbeatInterval, heartbeatTimeout);
}

0x05 Flink中具体应用

5.1 总述

5.1.1 RM, JM, TM之间关系

系统中有几个ResourceManager?整个 Flink 集群中只有一个 ResourceManager。

系统中有几个JobManager?JobManager 负责管理作业的执行。默认情况下,每个 Flink 集群只有一个 JobManager 实例。JobManager 相当于整个集群的 Master 节点,负责整个集群的任务管理和资源管理。

系统中有几个TaskManager?这个由具体启动方式决定。比如Flink on Yarn,Session模式能够指定拉起多少个TaskManager。 Per job模式中TaskManager数量是在提交作业时根据并发度动态计算,即Number of TM = Parallelism/numberOfTaskSlots。比如:有一个作业,Parallelism为10,numberOfTaskSlots为1,则TaskManager为10。

5.1.2 三者间心跳机制

Flink中ResourceManager、JobMaster、TaskExecutor三者之间存在相互检测的心跳机制:

  • ResourceManager会主动发送请求探测JobMaster、TaskExecutor是否存活。
  • JobMaster也会主动发送请求探测TaskExecutor是否存活,以便进行任务重启或者失败处理。

我们之前讲过,HeartbeatManagerSenderImpl属于Sender,HeartbeatManagerImpl属于Receiver。

  1. HeartbeatManagerImpl所处位置可以理解为client,存在于JobMaster与TaskExecutor中;
  2. HeartbeatManagerSenderImpl类,继承 HeartbeatManagerImpl类,用于周期发送心跳请求,所处位置可以理解为server, 存在于JobMaster、ResourceManager中。

ResourceManager 级别最高,所以两个HM都是Sender,监控taskManager和jobManager

1
2
3
4
5
6
public abstract class ResourceManager<WorkerType extends ResourceIDRetrievable>
extends FencedRpcEndpoint<ResourceManagerId>
implements ResourceManagerGateway, LeaderContender {
taskManagerHeartbeatManager = heartbeatServices.createHeartbeatManagerSender
jobManagerHeartbeatManager = heartbeatServices.createHeartbeatManagerSender
}

JobMaster级别中等,一个Sender, 一个Receiver,受到ResourceManager的监控,监控taskManager。

1
2
3
4
public class JobMaster extends FencedRpcEndpoint<JobMasterId> implements JobMasterGateway, JobMasterService {
taskManagerHeartbeatManager = heartbeatServices.createHeartbeatManagerSender
resourceManagerHeartbeatManager = heartbeatServices.createHeartbeatManager
}

TaskExecutor级别最低,两个Receiver,分别被JM和RM疾控。

1
2
3
4
public class TaskExecutor extends RpcEndpoint implements TaskExecutorGateway {
this.jobManagerHeartbeatManager = return heartbeatServices.createHeartbeatManager
this.resourceManagerHeartbeatManager = return heartbeatServices.createHeartbeatManager
}

以JobManager和TaskManager为例。JM在启动时会开启周期调度,向已经注册到JM中的TM发起心跳检查,通过RPC调用TM的requestHeartbeat方法,重置对JM超时线程的调用,表示当前JM状态正常。在TM的requestHeartbeat方法被调用后,通过RPC调用JM的receiveHeartbeat,重置对TM超时线程的调用,表示TM状态正常。

5.2 初始化过程

5.2.1 TaskExecutor初始化

TM初始化生成了两个Receiver HM。

1
2
3
4
5
6
7
8
9
10
11
public class TaskExecutor extends RpcEndpoint implements TaskExecutorGateway {
/** The heartbeat manager for job manager in the task manager. */
private final HeartbeatManager<AllocatedSlotReport, AccumulatorReport> jobManagerHeartbeatManager;

/** The heartbeat manager for resource manager in the task manager. */
private final HeartbeatManager<Void, TaskExecutorHeartbeatPayload> resourceManagerHeartbeatManager;

//初始化函数
this.jobManagerHeartbeatManager = createJobManagerHeartbeatManager(heartbeatServices, resourceId);
this.resourceManagerHeartbeatManager = createResourceManagerHeartbeatManager(heartbeatServices, resourceId);
}

生成HeartbeatManager时,就注册了ResourceManagerHeartbeatListener和JobManagerHeartbeatListener。

此时,两个HeartbeatManagerImpl中已经创建好对应monitor线程,只有在JM或者RM执行requestHeartbeat后,才会触发该线程的执行。

5.2.2 JobMaster的初始化

JM生成了一个Sender HM,一个Receiver HM。这里会注册 TaskManagerHeartbeatListener 和 ResourceManagerHeartbeatListener

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
public class JobMaster extends FencedRpcEndpoint<JobMasterId> implements JobMasterGateway, JobMasterService {
private HeartbeatManager<AccumulatorReport, AllocatedSlotReport> taskManagerHeartbeatManager;
private HeartbeatManager<Void, Void> resourceManagerHeartbeatManager;

private void startHeartbeatServices() {
taskManagerHeartbeatManager = heartbeatServices.createHeartbeatManagerSender(
resourceId,
new TaskManagerHeartbeatListener(),
getMainThreadExecutor(),
log);

resourceManagerHeartbeatManager = heartbeatServices.createHeartbeatManager(
resourceId,
new ResourceManagerHeartbeatListener(),
getMainThreadExecutor(),
log);
}
}

5.2.3 ResourceManager初始化

JobMaster在启动时候,会在startHeartbeatServices函数中生成两个Sender HeartbeatManager。

taskManagerHeartbeatManager :HeartbeatManagerSenderImpl对象,会反复启动一个定时器,定时扫描需要探测的对象并且发送心跳请求。

jobManagerHeartbeatManager :HeartbeatManagerSenderImpl,会反复启动一个定时器,定时扫描需要探测的对象并且发送心跳请求。

1
2
3
4
5
6
7
8
9
10
11
taskManagerHeartbeatManager = heartbeatServices.createHeartbeatManagerSender(
resourceId,
new TaskManagerHeartbeatListener(),
getMainThreadExecutor(),
log);

jobManagerHeartbeatManager = heartbeatServices.createHeartbeatManagerSender(
resourceId,
new JobManagerHeartbeatListener(),
getMainThreadExecutor(),
log);

5.3 注册过程

我们以TM与RM交互为例。TaskExecutor启动之后,需要注册到RM和JM中。

流程图如下:

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
* 1. Run in Task Manager
*
* TaskExecutor.onStart //Life cycle
* |
* +----> startTaskExecutorServices@TaskExecutor
* | //开始TM服务
* |
* +----> resourceManagerLeaderRetriever.start(new ResourceManagerLeaderListener());
* | // 开始连接到RM
* | // start by connecting to the ResourceManager
* |
* +----> notifyLeaderAddress@ResourceManagerLeaderListener
* | // 当RM状态变化之后,将回调到这里
* | // The listener for leader changes of the resource manager.
* |
* +----> reconnectToResourceManager@TaskExecutor
* | // 以下三步调用是渐进的,就是与RM联系。
* |
* +----> tryConnectToResourceManager@TaskExecutor
* |
* +----> connectToResourceManager()@TaskExecutor
* | // 主要作用是生成了 TaskExecutorToResourceManagerConnection
* |
* +----> start@TaskExecutorToResourceManagerConnection
* | // 开始RPC调用,将会调用到其基类RegisteredRpcConnection的start
* |
* +----> start@RegisteredRpcConnection
* | // RegisteredRpcConnection实现了组件之间注册联系的基本RPC
* |


* ~~~~~~~~ 这里是 Akka RPC

* 2. Run in Resource Manager
* 现在程序执行序列到达了RM, 主要是添加一个TargetRMSender HM
*
* registerTaskExecutor@ResourceManager
* |
* +----> taskExecutorGatewayFuture.handleAsync
* | // 异步调用到这里
* |
* +----> registerTaskExecutorInternal@ResourceManager
* | // RM的内部实现,将把TM注册到RM自己这里
* |
* +----> taskManagerHeartbeatManager.monitorTarget
* | // 生成HeartbeatMonitor,
* |
* +----> heartbeatTargets.put(resourceID,heartbeatMonitor);
* | // 把Monitor放到 HM in TM之中,就是说TM开始监控了RM
* |

* ~~~~~~~~ 这里是 Akka RPC

* 3. Run in Task Manager
* 现在程序回到了TM, 主要是添加一个TargetTMReceiver HM
*
* onRegistrationSuccess@TaskExecutorToResourceManagerConnection
* |
* |
* +----> onRegistrationSuccess@ResourceManagerRegistrationListener
* | // 回调函数
* |
* +----> runAsync(establishResourceManagerConnection)
* | // 异步执行
* |
* +----> establishResourceManagerConnection@TaskExecutor
* | // 说明已经和RM建立了联系,所以可以开始监控RM了
* |
* +----> resourceManagerHeartbeatManager.monitorTarget
* | // 生成HeartbeatMonitor,
* |
* +----> heartbeatTargets.put(resourceID,heartbeatMonitor);
* | // 把 RM 也注册到 TM了
* | // monitor the resource manager as heartbeat target

下面是具体文字描述。

5.3.1 TM注册到RM中

5.3.1.1 TM的操作
  • TaskExecutor启动之后,调用onStart,开始其生命周期。
  • onStart直接调用startTaskExecutorServices。
  • 启动服务的第一步就是与ResourceManager取得联系,这里注册了一个ResourceManagerLeaderListener(),用来监听RM Leader的变化。
1
2
3
4
private final LeaderRetrievalService resourceManagerLeaderRetriever;
// resourceManagerLeaderRetriever其实是EmbeddedLeaderService的实现,A simple leader election service, which selects a leader among contenders and notifies listeners.

resourceManagerLeaderRetriever.start(new ResourceManagerLeaderListener());
  • 当得到RM Leader的地址之后,会调用到回调函数notifyLeaderAddress@ResourceManagerLeaderListener,然后调用notifyOfNewResourceManagerLeader。
  • notifyOfNewResourceManagerLeader中获取到RM地址后,就通过reconnectToResourceManager与RM联系。
  • reconnectToResourceManager中间接调用到TaskExecutorToResourceManagerConnection。其作用是建立TaskExecutor 和 ResourceManager之间的联系。因为知道 ResourceManagerGateway所以才能进行RPC操作。
  • 然后在 TaskExecutorToResourceManagerConnection中,就通过RPC与RM联系。
5.3.1.2 RM的操作
  • RPC调用后,程序就来到了RM中,RM做如下操作:
  • 会注册一个新的TaskExecutor到自己的taskManagerHeartbeatManager中。
  • registerTaskExecutor@ResourceManager会通过异步调用到registerTaskExecutorInternal。
  • registerTaskExecutorInternal中首先看看是否这个TaskExecutor的ResourceID之前注册过,如果注册过就移除再添加一个新的TaskExecutor。
  • 通过 taskManagerHeartbeatManager.monitorTarget 开始进行心跳机制的注册。
1
2
3
4
5
6
7
8
9
taskManagerHeartbeatManager.monitorTarget(taskExecutorResourceId, new HeartbeatTarget<Void>() {
public void receiveHeartbeat(ResourceID resourceID, Void payload) {
// the ResourceManager will always send heartbeat requests to the
// TaskManager
}
public void requestHeartbeat(ResourceID resourceID, Void payload) {
taskExecutorGateway.heartbeatFromResourceManager(resourceID);
}
});

当注册完成后,RM中的Sender HM内部结构如下,能看出来多了一个Target:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
taskManagerHeartbeatManager = {HeartbeatManagerSenderImpl@8866} 
heartbeatPeriod = 10000
heartbeatTimeoutIntervalMs = 50000
ownResourceID = {ResourceID@8871} "040709f36ebf38f309fed518a88946af"
heartbeatListener = {ResourceManager$TaskManagerHeartbeatListener@8872}
mainThreadExecutor = {RpcEndpoint$MainThreadExecutor@8873}
heartbeatTargets = {ConcurrentHashMap@8875} size = 1
{ResourceID@8867} "630c15c9-4861-4b41-9c95-92504f458b71" -> {HeartbeatMonitorImpl@9448}
key = {ResourceID@8867} "630c15c9-4861-4b41-9c95-92504f458b71"
value = {HeartbeatMonitorImpl@9448}
resourceID = {ResourceID@8867} "630c15c9-4861-4b41-9c95-92504f458b71"
heartbeatTarget = {ResourceManager$2@8868}
scheduledExecutor = {RpcEndpoint$MainThreadExecutor@8873}
heartbeatListener = {ResourceManager$TaskManagerHeartbeatListener@8872}
heartbeatTimeoutIntervalMs = 50000
futureTimeout = {ScheduledFutureAdapter@10140}
state = {AtomicReference@9786} "RUNNING"
lastHeartbeat = 0
5.3.1.3 返回到TM

RM会通过RPC再次回到TaskExecutor,其新执行序列如下:

  • 首先RPC调用到了 onRegistrationSuccess@TaskExecutorToResourceManagerConnection。
  • 然后onRegistrationSuccess@ResourceManagerRegistrationListener中通过异步执行调用到了establishResourceManagerConnection。这说明TM已经和RM建立了联系,所以可以开始监控RM了。
  • 然后和RM操作类似,通过resourceManagerHeartbeatManager.monitorTarget 来把RM注册到自己这里。
1
2
HeartbeatMonitor<O> heartbeatMonitor = heartbeatMonitorFactory.createHeartbeatMonitor 
heartbeatTargets.put(resourceID, heartbeatMonitor);

当注册完成后,其Receiver HM结构如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
resourceManagerHeartbeatManager = {HeartbeatManagerImpl@10163} 
heartbeatTimeoutIntervalMs = 50000
ownResourceID = {ResourceID@8882} "96a9b80c-dd97-4b63-9049-afb6662ea3e2"
heartbeatListener = {TaskExecutor$ResourceManagerHeartbeatListener@10425}
mainThreadExecutor = {RpcEndpoint$MainThreadExecutor@10426}
heartbeatTargets = {ConcurrentHashMap@10427} size = 1
{ResourceID@8886} "122fa66685133b11ea26ee1b1a6cef75" -> {HeartbeatMonitorImpl@10666}
key = {ResourceID@8886} "122fa66685133b11ea26ee1b1a6cef75"
value = {HeartbeatMonitorImpl@10666}
resourceID = {ResourceID@8886} "122fa66685133b11ea26ee1b1a6cef75"
heartbeatTarget = {TaskExecutor$1@10668}
scheduledExecutor = {RpcEndpoint$MainThreadExecutor@10426}
heartbeatListener = {TaskExecutor$ResourceManagerHeartbeatListener@10425}
heartbeatTimeoutIntervalMs = 50000
futureTimeout = {ScheduledFutureAdapter@10992}
state = {AtomicReference@10667} "RUNNING"
lastHeartbeat = 0

5.3.2 TM注册到 JM

其调用基本思路与之前相同,就是TM和JM之间互相注册一个代表对方的monitor:

1
JobLeaderListenerImpl ----> establishJobManagerConnection

消息到了JM中,做如下操作。

1
2
registerTaskManager ----> taskManagerHeartbeatManager.monitorTarget
// monitor the task manager as heartbeat target

5.4 心跳过程

在任务提交之后,我们就进入了正常的心跳监控流程。我们依然用 TM 和 RM进行演示。

我们先给出一个流程图。

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
* 1. Run in Resouce Manager
*
* HeartbeatManagerSender in RM
* |
* +----> run@HeartbeatManagerSenderImpl
* | //遍历所有监控的Monitor(Target),逐一在Target上调用requestHeartbeat
* |
* +----> requestHeartbeat@HeartbeatManagerSenderImpl
* | // 将调用具体监控对象的自定义函数
* | // heartbeatTarget.requestHeartbeat(getOwnResourceID(), payload);
* |
* +----> getHeartbeatListener().retrievePayload
* | // 调用到TaskManagerHeartbeatListener@ResourceManager
* | // 这里是return null;,因为RM不会是任何人的Receiver
* |
* +----> requestHeartbeat@HeartbeatTarget
* | // 调用到Target这里,代码在ResourceManager这里,就是生成Target时候赋值的
* |
* +----> taskExecutorGateway.heartbeatFromResourceManager
* | // 会通过gateway RPC 调用到TM,这就是主动对TM发起了心跳请求
* |

* ~~~~~~~~ 这里是 Akka RPC

* 2. Run in Task Manager
* 现在程序执行序列到达了TM, 主要是 1. 重置TMMonitor线程; 2.返回一些负载信息;
*
* heartbeatFromResourceManager@TaskExecutor
* |
* +----> resourceManagerHeartbeatManager.requestHeartbeat(resourceID, null);
* | //开始要调用到 Receiver HM in Task Manager
* |
* +----> requestHeartbeat@HeartbeatManager in TM
* | // 在Receiver HM in Task Manager 这里运行
* |
* +----> reportHeartbeat@HeartbeatMonitor
* | //reportHeartbeat : 记录发起请求的这个时间点,然后resetHeartbeatTimeout
* |
* +----> resetHeartbeatTimeout@HeartbeatMonitor
* | // 如果Monitor状态依然是RUNNING,则取消之前设置的ScheduledFuture。
* | // 重新创建一个ScheduleFuture。因为如果不取消,则之前那个ScheduleFuture运行时
* | // 会调用HeartbeatMonitorImpl.run函数,run直接compareAndSet后,通知目标函数
* | // 目前已经超时,即调用heartbeatListener.notifyHeartbeatTimeout。
* | // 这里代表 JM 状态正常。
* |
* +----> heartbeatListener.reportPayload
* | // 把Target节点的最新的heartbeatPayload通知给heartbeatListener。
* | // heartbeatListerner是外部传入的,它根据所拥有的节点的心跳记录做监听管理。
* |
* +----> heartbeatTarget.receiveHeartbeat(getOwnResourceID(), heartbeatListener.retrievePayload(requestOrigin));
* |
* |
* +----> retrievePayload@ResourceManagerHeartbeatListener in TM
* | // heartbeatTarget.receiveHeartbeat参数调用的
* |
* +----> return new TaskExecutorHeartbeatPayload
* |
* |
* +----> receiveHeartbeat in TM
* | // 回到 heartbeatTarget.receiveHeartbeat,这就是TM生成Target的时候的自定义函数
* | // 就是响应一个心跳消息回给RM
* |
* +----> resourceManagerGateway.heartbeatFromTaskManager
* | // 会通过gateway RPC 调用到 ResourcManager
* |

* ~~~~~~~~ 这里是 Akka RPC

* 3. Run in Resouce Manager
* 现在程序回到了RM, 主要是 1.重置RMMonitor线程;2. 上报收到TaskExecutor的负载信息
*
* heartbeatFromTaskManager in RM
* |
* |
* +----> taskManagerHeartbeatManager.receiveHeartbeat
* | // 这是个Sender HM
* |
* +----> HeartbeatManagerImpl.receiveHeartbeat
* |
* |
* +----> HeartbeatManagerImpl.reportHeartbeat(heartbeatOrigin);
* |
* |
* +----> heartbeatMonitor.reportHeartbeat();
* | // 这里就是重置RM 这里对应的Monitor。在reportHeartbeat重置 JM monitor线程的触发,即cancelTimeout取消注册时候的超时定时任务,并且注册下一个超时检测futureTimeout;这代表TM正常执行。
* |
* +----> heartbeatListener.reportPayload
* | //把Target节点的最新的heartbeatPayload通知给 TaskManagerHeartbeatListener。heartbeatListerner是外部传入的,它根据所拥有的节点的心跳记录做监听管理。
* |
* +----> slotManager.reportSlotStatus(instanceId, payload.getSlotReport());
* | // TaskManagerHeartbeatListener中调用,上报收到TaskExecutor的负载信息
* |

下面是具体文字描述。

5.4.1 ResourceManager主动发起

5.4.1.1 Sender遍历所有监控的Monitor(Target)

心跳机制是由Sender主动发起的。这里就是 ResourceManager 的HeartbeatManagerSenderImpl中定时schedual调用,这里会遍历所有监控的Monitor(Target),逐一在Target上调用requestHeartbeat。

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
// HeartbeatManagerSenderImpl中的代码

@Override
public void run() {
if (!stopped) {
for (HeartbeatMonitor<O> heartbeatMonitor : getHeartbeatTargets().values()) {
// 这里向被监控对象节点发起一次心跳请求,载荷是heartbeatPayLoad,要求被监控对象回应心跳
requestHeartbeat(heartbeatMonitor);
}
getMainThreadExecutor().schedule(this, heartbeatPeriod, TimeUnit.MILLISECONDS);
}
}
}

// 运行时候的变量
this = {HeartbeatManagerSenderImpl@9037}
heartbeatPeriod = 10000
heartbeatTimeoutIntervalMs = 50000
ownResourceID = {ResourceID@8788} "d349506cae32cadbe99b9f9c49a01c95"
heartbeatListener = {ResourceManager$TaskManagerHeartbeatListener@8789}
mainThreadExecutor = {RpcEndpoint$MainThreadExecutor@8790}

// 调用栈如下
requestHeartbeat:711, ResourceManager$2 (org.apache.flink.runtime.resourcemanager)
requestHeartbeat:702, ResourceManager$2 (org.apache.flink.runtime.resourcemanager)
requestHeartbeat:92, HeartbeatManagerSenderImpl (org.apache.flink.runtime.heartbeat)
run:81, HeartbeatManagerSenderImpl (org.apache.flink.runtime.heartbeat)
call:511, Executors$RunnableAdapter (java.util.concurrent)
run$$$capture:266, FutureTask (java.util.concurrent)
run:-1, FutureTask (java.util.concurrent)
5.4.1.2 Target进行具体操作

具体监控对象 Target 会调用自定义的requestHeartbeat。

1
2
3
4
5
6
7
8
9
10
11
12
13
HeartbeatManagerSenderImpl

private void requestHeartbeat(HeartbeatMonitor<O> heartbeatMonitor) {
O payload = getHeartbeatListener().retrievePayload(heartbeatMonitor.getHeartbeatTargetId());
final HeartbeatTarget<O> heartbeatTarget = heartbeatMonitor.getHeartbeatTarget();

// 这里就是具体监控对象
heartbeatTarget.requestHeartbeat(getOwnResourceID(), payload);
}

heartbeatTarget = {ResourceManager$2@10688}
taskExecutorGateway = {$Proxy42@9459} "org.apache.flink.runtime.rpc.akka.AkkaInvocationHandler@6d0c8334"
this$0 = {StandaloneResourceManager@9458}

请注意,每一个Target都是由ResourceManager生成的。ResourceManager之前注册成为Monitor时候就注册了这个HeartbeatTarget。

这个HeartbeatTarget的定义如下,两个函数是:

  • receiveHeartbeat :这个是空,因为RM没有自己的Sender。
  • requestHeartbeat :这个针对TM,就是调用TM的heartbeatFromResourceManager,当然是通过RPC调用。
5.4.1.3 RPC调用

会调用到ResourceManager定义的函数requestHeartbeat,而requestHeartbeat会通过gateway调用到TM,这就是主动对TM发起了心跳请求。

1
2
3
4
5
6
7
8
9
10
11
12
taskManagerHeartbeatManager.monitorTarget(taskExecutorResourceId, new HeartbeatTarget<Void>() {
@Override
public void receiveHeartbeat(ResourceID resourceID, Void payload) {
// the ResourceManager will always send heartbeat requests to the TaskManager
}

@Override
public void requestHeartbeat(ResourceID resourceID, Void payload) {
//就是调用到这里
taskExecutorGateway.heartbeatFromResourceManager(resourceID);
}
});

5.4.2 RM通过RPC调用TM

通过taskExecutorGateway。心跳程序执行就通过RPC从RM跳跃到了TM。

taskExecutorGateway.heartbeatFromResourceManager 的意义就是:通过RPC调用回到TaskExecutor。这个是在TaskExecutorGateway就定义好的。

1
2
// TaskExecutor RPC gateway interface.
public interface TaskExecutorGateway extends RpcGateway

TaskExecutor实现了TaskExecutorGateway,所以具体在TaskExecutor内部实现了接口函数。

1
2
3
4
5
@Override
public void heartbeatFromResourceManager(ResourceID resourceID) {
//调用到了这里 ...........
resourceManagerHeartbeatManager.requestHeartbeat(resourceID, null);
}

TM中,resourceManagerHeartbeatManager 定义如下。

1
2
/** The heartbeat manager for resource manager in the task manager. */
private final HeartbeatManager<Void, TaskExecutorHeartbeatPayload> resourceManagerHeartbeatManager;

所以下面就是执行TM中的Receiver HM。在这个过程中有两个处理步骤:

  1. 调用对应HeartbeatMonitor的reportHeartbeat方法,cancelTimeout取消注册时候的超时定时任务,并且注册下一个超时检测futureTimeout;
  2. 调用monitorTarget的receiveHeartbeat方法,也就是会通过rpc调用JobMaster的heartbeatFromTaskManager方法返回一些负载信息;

具体是调用 requestHeartbeat@HeartbeatManager。在其中会

  • 调用reportHeartbeat@HeartbeatMonitor,记录发起请求的这个时间点,然后resetHeartbeatTimeout。
  • 在resetHeartbeatTimeout@HeartbeatMonitor之中,如果Monitor状态依然是RUNNING,则取消之前设置的ScheduledFuture。重新创建一个ScheduleFuture。因为如果不取消,则之前那个ScheduleFuture运行时会调用HeartbeatMonitorImpl.run函数,run直接compareAndSet后,通知目标函数目前已经超时,即调用heartbeatListener.notifyHeartbeatTimeout。
  • 调用 heartbeatListener.reportPayload,把Target节点的最新的heartbeatPayload通知给heartbeatListener。
  • 调用 heartbeatTarget.receiveHeartbeat(getOwnResourceID(), heartbeatListener.retrievePayload(requestOrigin)); 就是响应一个心跳消息回给RM。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
@Override
public void requestHeartbeat(final ResourceID requestOrigin, I heartbeatPayload) {
if (!stopped) {
log.debug("Received heartbeat request from {}.", requestOrigin);

final HeartbeatTarget<O> heartbeatTarget = reportHeartbeat(requestOrigin);

if (heartbeatTarget != null) {
if (heartbeatPayload != null) {
heartbeatListener.reportPayload(requestOrigin, heartbeatPayload);
}

heartbeatTarget.receiveHeartbeat(getOwnResourceID(), heartbeatListener.retrievePayload(requestOrigin));
}
}
}

最后会通过resourceManagerGateway.heartbeatFromTaskManager 调用到 ResourcManager。

5.4.3 TM 通过RPC回到 RM

JobMaster在接收到rpc请求后调用其heartbeatFromTaskManager方法,会调用taskManagerHeartbeatManager的receiveHeartbeat方法,在这个过程中同样有两个处理步骤:

  1. 调用对应HeartbeatMonitor的reportHeartbeat方法,cancelTimeout取消注册时候的超时定时任务,并且注册下一个超时检测futureTimeout;
  2. 调用TaskManagerHeartbeatListener的reportPayload方法,上报收到TaskExecutor的负载信息

至此一次完成心跳过程已经完成,会根据heartbeatInterval执行下一次心跳。

5.5 超时处理

5.5.1 TaskManager

首先,在HeartbeatMonitorImpl中,如果超时,会调用Listener。

1
2
3
4
5
6
public void run() {
// The heartbeat has timed out if we're in state running
if (state.compareAndSet(State.RUNNING, State.TIMEOUT)) {
heartbeatListener.notifyHeartbeatTimeout(resourceID);
}
}

这就来到了ResourceManagerHeartbeatListener,会尝试再次连接RM。

1
2
3
4
5
6
7
8
9
10
11
12
13
private class ResourceManagerHeartbeatListener implements HeartbeatListener<Void, TaskExecutorHeartbeatPayload> {

@Override
public void notifyHeartbeatTimeout(final ResourceID resourceId) {
validateRunsInMainThread();
// first check whether the timeout is still valid
if (establishedResourceManagerConnection != null && establishedResourceManagerConnection.getResourceManagerResourceId().equals(resourceId)) {
reconnectToResourceManager(new TaskManagerException(
String.format("The heartbeat of ResourceManager with id %s timed out.", resourceId)));
} else {
.....
}
}

5.5.2 ResourceManager

RM就直接简单粗暴,关闭连接。

1
2
3
4
5
6
7
8
9
10
private class TaskManagerHeartbeatListener implements HeartbeatListener<TaskExecutorHeartbeatPayload, Void> {

@Override
public void notifyHeartbeatTimeout(final ResourceID resourceID) {
validateRunsInMainThread();
closeTaskManagerConnection(
resourceID,
new TimeoutException("The heartbeat of TaskManager with id " + resourceID + " timed out."));
}
}

0x06 解决问题

心跳机制我们讲解完了,但是我们最初提到的异常应该如何解决呢?在程序最开始生成环境变量时候,通过设置环境变量的配置即可搞定:

1
2
3
Configuration conf = new Configuration();
conf.setString("heartbeat.timeout", "18000000");
final LocalEnvironment env = ExecutionEnvironment.createLocalEnvironment(conf);

0x07 参考

[flink-001]flink的心跳机制

Flink中心跳机制

flink1.8 心跳服务

你有必要了解一下Flink底层RPC使用的框架和原理

flink RPC(akka)

弄清Flink1.8的远程过程调用(RPC)

Apache Flink源码解析 (七)Flink RPC的底层实现

flink源码阅读第一篇—入口

flink-on-yarn 基础架构和启动流程

在Flink SQL中, 元数据的管理分为三层: catalog-> database-> table,
我们知道Flink SQL是依托calcite框架来进行SQL执行树生产,校验,优化等等, 所以本文讲介绍FlinkSQL是如何来结合Calcite来进行元数据管理的.

calcite开放的接口

1
2
3
4
5
6
7
public interface Schema {
Table getTable(String name);

Schema getSubSchema(String name);

....
}

如接口所示, Schema接口,可以通过table名来获得一张表, 可以通过schema名来获得一个子schema.

1
2
3
4
public interface Table {
RelDataType getRowType(RelDataTypeFactory typeFactory);
....
}

看Table的接口, 主要就是返回table的RelDataType.

Flink的相关实现

接下来,我们来看下Flink是如何实现这些接口的:

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
public class CatalogManagerCalciteSchema extends FlinkSchema {
@Override
public Schema getSubSchema(String schemaName) {
if (catalogManager.schemaExists(name)) {
return new CatalogCalciteSchema(name, catalogManager, isStreamingMode);
} else {
return null;
}
}
}
public class CatalogCalciteSchema extends FlinkSchema {
@Override
public Schema getSubSchema(String schemaName) {
if (catalogManager.schemaExists(catalogName, schemaName)) {
return new DatabasecalciteSchema(schemaName, catalogNmae, catalogManager, isStreamingMode);
}
}
}
public class DatabaseCalciteSchema extends FlinkSchema {
private final String databaseName;
private final String catalogName;
private final CatalogManager catalogManager;

@Override
public Table getTable(String tableName) {
ObjectIdentifier identifier = ObjectIdentifier.of(catalogName, databaseName, tableName);
return catalogManager.getTable(identifier)
.map(result -> {
CatalogBaseTable table = result.getTable();
FlinkStatistic statistic = getStatistic(result.isTemporary(), table, identifier);
return new CatalogSchemaTable(identifier,
table,
statistic,
catalogManager.getCatalog(catalogName)
.flatMap(Catalog::getTableFactory)
.orElse(null),
isStreamingMode,
result.isTemporary());
})
.orElse(null);
}

@Override
public Schema getSubSchema(String name) {
return null;
}
}

很容易发现,CatalogSchema返回DatabaseSchema, DatabaseSchema返回Table,
这样就容易理解,Flink的三层结构是怎样的了. 同时, 具体的元数据实际上都是在catalogManager中。

DatabaseSchema中返回的Table类型为CatalogSchemaTable,我们来看下具体的结结构是怎样的,
上文中也提到了,Table接口主为getRowType函数, 用于返回某个table的type信息。
TableSchema是Flink内部用于保存各个字段的类型信息的类, 通过相关的转化函数,转换为calcite的type类型.

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
public class CatalogSchemaTable extends AbstractTable implements TemporalTable {

private final ObjectIdentifier tableIdentifier;
private final CatalogBaseTable catalogBaseTable;
private final FlinkStatistic statistic;
private final boolean isStreamingMode;
private final boolean isTemporary;
...
private static RelDataType getRowType(RelDataTypeFactory typeFactory,
CatalogBaseTable catalogBaseTable,
boolean isStreamingMode) {
final FlinkTypeFactory flinkTypeFactory = (FlinkTypeFactory) typeFactory;
TableSchema tableSchema = catalogBaseTable.getSchema();
final DataType[] fieldDataTypes = tableSchema.getFieldDataTypes();
if (!isStreamingMode
&& catalogBaseTable instanceof ConnectorCatalogTable
&& ((ConnectorCatalogTable) catalogBaseTable).getTableSource().isPresent()) {
// If the table source is bounded, materialize the time attributes to normal TIMESTAMP type.
// Now for ConnectorCatalogTable, there is no way to
// deduce if it is bounded in the table environment, so the data types in TableSchema
// always patched with TimeAttribute.
// See ConnectorCatalogTable#calculateSourceSchema
// for details.

// Remove the patched time attributes type to let the TableSourceTable handle it.
// We should remove this logic if the isBatch flag in ConnectorCatalogTable is fixed.
// TODO: Fix FLINK-14844.
for (int i = 0; i < fieldDataTypes.length; i++) {
LogicalType lt = fieldDataTypes[i].getLogicalType();
if (lt instanceof TimestampType
&& (((TimestampType) lt).getKind() == TimestampKind.PROCTIME
|| ((TimestampType) lt).getKind() == TimestampKind.ROWTIME)) {
int precision = ((TimestampType) lt).getPrecision();
fieldDataTypes[i] = DataTypes.TIMESTAMP(precision);
}
}
}
return TableSourceUtil.getSourceRowType(flinkTypeFactory,
tableSchema,
scala.Option.empty(),
isStreamingMode);
}
}

CatalogBaseTable接口定义如下, Flink的Table的参数(schema参数,connector参数)都可以最终表示为一个map.

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
public interface CatalogBaseTable {
/**
* Get the properties of the table.
*
* @return property map of the table/view
*/
Map<String, String> getProperties();

/**
* Get the schema of the table.
*
* @return schema of the table/view.
*/
TableSchema getSchema();

/**
* Get comment of the table or view.
*
* @return comment of the table/view.
*/
String getComment();

/**
* Get a deep copy of the CatalogBaseTable instance.
*
* @return a copy of the CatalogBaseTable instance
*/
CatalogBaseTable copy();

/**
* Get a brief description of the table or view.
*
* @return an optional short description of the table/view
*/
Optional<String> getDescription();

/**
* Get a detailed description of the table or view.
*
* @return an optional long description of the table/view
*/
Optional<String> getDetailedDescription();
}

FlinkSchema的使用

上面都是的相关接口都是Flink用于适配calcite框架元数据的相关实现。
那么这些类具体是在哪里调用的? 已经什么时候会被调用到?
calcite中的schema,主要是在validate过程中, 获得对应table的字段信息, 对应的function的返回值信息,
确保SQL的字段名, 字段类型是正确的.
类的依赖关系为:
validator —> schemaReader —> schema

FlinkPlannerImpl.scala中

1
2
3
4
5
6
7
8
9
10
private def createSqlValidator(catalogReader: CatalogReader) = {
val validator = new FlinkCalciteSqlValidator(
operatorTable,
catalogReader,
typeFactory)
validator.setIdentifierExpansion(true)
// Disable implicit type coercion for now.
validator.setEnableTypeCoercion(false)
validator
}

PlanningConfigurationBuilder.java

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
private CatalogReader createCatalogReader(
boolean lenientCaseSensitivity,
String currentCatalog,
String currentDatabase) {
SqlParser.Config sqlParserConfig = getSqlParserConfig();
final boolean caseSensitive;
if (lenientCaseSensitivity) {
caseSensitive = false;
} else {
caseSensitive = sqlParserConfig.caseSensitive();
}

SqlParser.Config parserConfig = SqlParser.configBuilder(sqlParserConfig)
.setCaseSensitive(caseSensitive)
.build();

return new CatalogReader(
rootSchema,
asList(
asList(currentCatalog, currentDatabase),
singletonList(currentCatalog)
),
typeFactory,
CalciteConfig.connectionConfig(parserConfig));
}

综上所诉, 我们就知道了Flink是如何来利用calciteschema来管理Flink的table信息的.

Flink SQL Join

Join的几种形式,与实现原理

双流Join

Left join

TTL

从 Flink 1.6 版本开始,社区引入了状态 TTL(Time-To-Live)特性。在通过Flink SQL 实现流处理时,开发者可以为作业 SQL 设置TTL,实现过期状态的自动清理,从而防止作业状态无限膨胀

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
SELECT
t_date,
COUNT ( DISTINCT user_id ) AS cnt_login, -- 今日登录用户数
COUNT ( DISTINCT CASE WHEN t_date = t_debut THEN user_id END ) AS cnt_new -- 今日新用户数
FROM
(
-- 计算每个用户有史以来的最小登录时间
SELECT
t_date,
user_id,
MIN (t_date) OVER (
PARTITION BY user_id
ORDER BY proctime
ROWS BETWEEN 1 PRECEDING AND CURRENT ROW
) AS t_debut
FROM Login
) AS t
GROUP BY t_date

Query Configuration 查询配置

https://ci.apache.org/projects/flink/flink-docs-stable/dev/table/streaming/query_configuration.html

Flink Table API 和SQL接口提供参数来调整连续查询的准确性和资源消耗。参数通过 QueryConfig 对象指定。QueryConfig 可以从 TableEnvironment 获得。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
val env = StreamExecutionEnvironment.getExecutionEnvironment
val tableEnv = TableEnvironment.getTableEnvironment(env)

// 获取 query configuration
val qConfig: StreamQueryConfig = tableEnv.queryConfig

// 设置查询参数
qConfig.withIdleStateRetentionTime(Time.hours(12), Time.hours(24))

// 定义查询和 TableSink
val result: Table = ???
val sink: TableSink[Row] = ???

// TableSink 发送结果表时传递查询参数
result.writeToSink(sink, qConfig)

// 转换为 DataStream 时传递查询参数
val stream: DataStream[Row] = result.toAppendStream[Row](qConfig)

空闲状态保持时间(Idle State Retention Time)参数定义一个键的状态在一次更新之后保存多久后删除。

通过删除键的状态,连续查询会完全忘记它之前已经看过这个键。如果删除的键再次出现,则被视为具有相应键的第一个记录。对于前面的查询示例,这意味着 sessionId 的计数从0开始。

配置空闲状态保存时间有两个参数:

minimum idle state retention time,定义非活动键的状态在删除前至少保持多少时间。
maximum idle state retention time,定义非活动键的状态在删除前最多保持多少时间。

对于前面的查询示例:

1
2
3
4
val qConfig: StreamQueryConfig = ???

// 设置 idle state retention time: min = 12 hours, max = 24 hours
qConfig.withIdleStateRetentionTime(Time.hours(12), Time.hours(24))

清理状态需要额外的记录,对于 minTime 和 maxTime 较大差异的情况成本更低,因此 minTime 和 maxTime 直接必须至少相差5分钟。

维表Join

启用AsyncIO

时间表Join: Temporal Table Join

1
2
3
4
5
6
SELECT
o.amout, o.currency, r.rate, o.amount * r.rate
FROM
Orders AS o
JOIN LatestRates FOR SYSTEM_TIME AS OF o.proctime AS r
ON r.currency = o.currency

Flink SQL 代码生成

从”UDF不应有状态” 切入来剖析Flink SQL代码生成

核心类: org.apache.flink.table.runtime.CRowProcessRunner

打印出codegen语句:

1
log4j.logger.org.apache.flink.runtime.CRowProcessRunner=DEBUG, file
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
alter table
b_cdg_cft_t_trade_user_fund_merge_ods_t_test_fund_Facc_time
set
(
'consumerGroup' = 't_cdg_cft_b_cdg_cft_t_trade_user_fund_merge_ods_cg_rtdw_prod_003',
'inDebugMode' = 'true',
'consume-from-max' = 'true'
);

-- 自定义控制台输出下游数据表,方便查看运行结果
alter table
console
set
(
-- 同上,console的schema定义的是下游表的元数据
'schema' = 'Facc_time:varchar,Flistid:varchar,Fsg_sh_type:bigint,Ftotal_fee:bigint,Funion_id:bigint,Fsg_sh_name:varchar,Ftrade_id:varchar'
);
-- 结果数据写入console,可以在应用运维->TM日志中查看结果输出
insert into
console
select
Facc_time,
Flistid,
Fsg_sh_type,
Ftotal_fee,
Funion_id,
Fsg_sh_name,
Ftrade_id
from
b_cdg_cft_t_trade_user_fund_merge_ods_t_test_fund_Facc_time
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
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
2021-10-27 15:06:43,520 DEBUG [Task-46f8730428df9ecd6d7318a02bdc405e(0/1)] org.apache.flink.table.runtime.conversion.CRowToJavaTupleMapRunner  - Compiling MapFunction: DataStreamSinkConversion$14 

Code:

public class DataStreamSinkConversion$14 extends org.apache.flink.api.common.functions.RichMapFunction {


final org.apache.flink.types.Row out =
new org.apache.flink.types.Row(7);

private org.apache.flink.types.Row in1;


public DataStreamSinkConversion$14() throws Exception {


}




@Override
public void open(org.apache.flink.configuration.Configuration parameters) throws Exception {


}

@Override
public Object map(Object _in1) throws Exception {
in1 = (org.apache.flink.types.Row) _in1;

boolean isNull$13 = (java.lang.String) in1.getField(6) == null;
java.lang.String result$12;
if (isNull$13) {
result$12 = "";
}
else {
result$12 = (java.lang.String) (java.lang.String) in1.getField(6);
}


boolean isNull$11 = (java.lang.String) in1.getField(5) == null;
java.lang.String result$10;
if (isNull$11) {
result$10 = "";
}
else {
result$10 = (java.lang.String) (java.lang.String) in1.getField(5);
}


boolean isNull$5 = (java.lang.Long) in1.getField(2) == null;
long result$4;
if (isNull$5) {
result$4 = -1L;
}
else {
result$4 = (java.lang.Long) in1.getField(2);
}


boolean isNull$1 = (java.lang.String) in1.getField(0) == null;
java.lang.String result$0;
if (isNull$1) {
result$0 = "";
}
else {
result$0 = (java.lang.String) (java.lang.String) in1.getField(0);
}


boolean isNull$9 = (java.lang.Long) in1.getField(4) == null;
long result$8;
if (isNull$9) {
result$8 = -1L;
}
else {
result$8 = (java.lang.Long) in1.getField(4);
}


boolean isNull$3 = (java.lang.String) in1.getField(1) == null;
java.lang.String result$2;
if (isNull$3) {
result$2 = "";
}
else {
result$2 = (java.lang.String) (java.lang.String) in1.getField(1);
}


boolean isNull$7 = (java.lang.Long) in1.getField(3) == null;
long result$6;
if (isNull$7) {
result$6 = -1L;
}
else {
result$6 = (java.lang.Long) in1.getField(3);
}







if (isNull$1) {
out.setField(0, null);
}
else {
out.setField(0, result$0);
}



if (isNull$3) {
out.setField(1, null);
}
else {
out.setField(1, result$2);
}



if (isNull$5) {
out.setField(2, null);
}
else {
out.setField(2, result$4);
}



if (isNull$7) {
out.setField(3, null);
}
else {
out.setField(3, result$6);
}



if (isNull$9) {
out.setField(4, null);
}
else {
out.setField(4, result$8);
}



if (isNull$11) {
out.setField(5, null);
}
else {
out.setField(5, result$10);
}



if (isNull$13) {
out.setField(6, null);
}
else {
out.setField(6, result$12);
}

return out;

}

@Override
public void close() throws Exception {


}
}

2021-10-27 15:06:43,520 DEBUG [Task-f24f04d8dcfb922ebc646161f8d38c37(0/1)] org.apache.flink.table.runtime.CRowMapRunner - Compiling MapFunction: DataStreamSinkConversion$29

Code:

public class DataStreamSinkConversion$29 extends org.apache.flink.api.common.functions.RichMapFunction {


final org.apache.flink.types.Row out =
new org.apache.flink.types.Row(7);

private org.apache.flink.types.Row in1;


public DataStreamSinkConversion$29() throws Exception {


}




@Override
public void open(org.apache.flink.configuration.Configuration parameters) throws Exception {


}

@Override
public Object map(Object _in1) throws Exception {
in1 = (org.apache.flink.types.Row) _in1;

boolean isNull$28 = (java.lang.String) in1.getField(6) == null;
java.lang.String result$27;
if (isNull$28) {
result$27 = "";
}
else {
result$27 = (java.lang.String) (java.lang.String) in1.getField(6);
}


boolean isNull$26 = (java.lang.String) in1.getField(5) == null;
java.lang.String result$25;
if (isNull$26) {
result$25 = "";
}
else {
result$25 = (java.lang.String) (java.lang.String) in1.getField(5);
}


boolean isNull$20 = (java.lang.Long) in1.getField(2) == null;
long result$19;
if (isNull$20) {
result$19 = -1L;
}
else {
result$19 = (java.lang.Long) in1.getField(2);
}


boolean isNull$16 = (java.lang.String) in1.getField(0) == null;
java.lang.String result$15;
if (isNull$16) {
result$15 = "";
}
else {
result$15 = (java.lang.String) (java.lang.String) in1.getField(0);
}


boolean isNull$24 = (java.lang.Long) in1.getField(4) == null;
long result$23;
if (isNull$24) {
result$23 = -1L;
}
else {
result$23 = (java.lang.Long) in1.getField(4);
}


boolean isNull$18 = (java.lang.String) in1.getField(1) == null;
java.lang.String result$17;
if (isNull$18) {
result$17 = "";
}
else {
result$17 = (java.lang.String) (java.lang.String) in1.getField(1);
}


boolean isNull$22 = (java.lang.Long) in1.getField(3) == null;
long result$21;
if (isNull$22) {
result$21 = -1L;
}
else {
result$21 = (java.lang.Long) in1.getField(3);
}







if (isNull$16) {
out.setField(0, null);
}
else {
out.setField(0, result$15);
}



if (isNull$18) {
out.setField(1, null);
}
else {
out.setField(1, result$17);
}



if (isNull$20) {
out.setField(2, null);
}
else {
out.setField(2, result$19);
}



if (isNull$22) {
out.setField(3, null);
}
else {
out.setField(3, result$21);
}



if (isNull$24) {
out.setField(4, null);
}
else {
out.setField(4, result$23);
}



if (isNull$26) {
out.setField(5, null);
}
else {
out.setField(5, result$25);
}



if (isNull$28) {
out.setField(6, null);
}
else {
out.setField(6, result$27);
}

return out;

}

@Override
public void close() throws Exception {


}
}

一个样例

1
2
3
4
5
6
7
8
9
10
11
12
13
CREATE TABLE mysql_orders (
ei STRING,
sei STRING,
ui STRING,
etime STRING
) WITH (
'connector.type' = 'jdbc',
'connector.driver' = 'com.mysql.cj.jdbc.Driver',
'connector.url' = 'jdbc:mysql://localhost:3306/flink_demo',
'connector.table' = 'orders',
'connector.username' = 'root',
'connector.password' = 'abc.ABC.123'
)
1
insert into mysql_orders select concat(ei,'#',sei) as pk,sei,ui,etime from orders

calcite解析后:

1
2
3
org.apache.calcite.sql2rel,DEBUG,Plan after converting SqlNode to RelNode
LogicalProject(pk=[CONCAT($0, _UTF-16LE'#', $1)], sei=[$1], ui=[$2], etime=[$3])
LogicalTableScan(table=[[default_catalog, default_database, Unregistered_TableSource_566698125, source: [CsvTableSource(read fields: ei, sei, ui, etime)]]])

优化:

1
2
3
4
5
org.apache.flink.table.planner.plan.optimize.program.FlinkGroupProgram,DEBUG,optimize convert table references before rewriting sub-queries to semi-join cost 31 ms.
optimize result:
LogicalSink(name=[`default_catalog`.`default_database`.`mysql_orders`], fields=[ei, sei, ui, etime])
+- LogicalProject(ei=[CONCAT($0, _UTF-16LE'#', $1)], sei=[CAST($1):VARCHAR(2147483647) CHARACTER SET "UTF-16LE"], ui=[CAST($2):VARCHAR(2147483647) CHARACTER SET "UTF-16LE"], etime=[CAST($3):VARCHAR(2147483647) CHARACTER SET "UTF-16LE"])
+- LogicalTableScan(table=[[default_catalog, default_database, Unregistered_TableSource_566698125, source: [CsvTableSource(read fields: ei, sei, ui, etime)]]])
1
2
3
4
org.apache.calcite.plan.RelOptPlanner,DEBUG,Cheapest plan:
FlinkLogicalSink(name=[`default_catalog`.`default_database`.`mysql_orders`], fields=[ei, sei, ui, etime]): rowcount = 1.0E8, cumulative cost = {3.0E8 rows, 3.0E8 cpu, 4.8E9 io, 0.0 network, 0.0 memory}, id = 73
FlinkLogicalCalc(select=[CONCAT(ei, _UTF-16LE'#', sei) AS ei, CAST(sei) AS sei, CAST(ui) AS ui, CAST(etime) AS etime]): rowcount = 1.0E8, cumulative cost = {2.0E8 rows, 2.0E8 cpu, 4.8E9 io, 0.0 network, 0.0 memory}, id = 72
FlinkLogicalTableSourceScan(table=[[default_catalog, default_database, Unregistered_TableSource_566698125, source: [CsvTableSource(read fields: ei, sei, ui, etime)]]], fields=[ei, sei, ui, etime]): rowcount = 1.0E8, cumulative cost = {1.0E8 rows, 1.0E8 cpu, 4.8E9 io, 0.0 network, 0.0 memory}, id = 71

实际的执行计划分为三层:

Flink SQL原理之SQL执行流程

Flink使用Calcite实现了SQL的解析、转换、执行计划优化和转换,那么,Flink SQL是如何执行的呢?

Calcite处理SQL的流程

  1. SQL解析(SQL -> SqlNode): 将SQL解析为AST(抽象语法树, Calcite中用SqlNode表示)
  2. SqlNode验证(SqlNode -> SqlNode) : 根据元数据信息(表名、字段名、函数名和数据类型等)进行语法验证
  3. 语义分析(SqlNode -> RelNode/RexNode: relational expression): 根据SqlNode与元数据信息构建RelNode树,也就是逻辑计划(Logical Plan)
  4. 逻辑计划优化(RelNode->RelNode): 优化器的核心,Calcite提供了两种Planner(HepPlanner和VolcanoPlanner),按照相应的规则(Rule)进行优化
    • HepPlanner: 启发式优化器,RBO, 按照规则匹配,直到最大次数或遍历后不再match rule
    • VolcanoPlanner: CBO,一直迭代,直到找到cost最小的plan
  5. 生成物理执行计划: 将Plan映射为Flink Graph

测试样例

MySQL实体表

1
2
3
4
5
6
7
8
9
10
CREATE TABLE IF NOT EXISTS t_rt_agg_result(
ks int(11) not null primary key auto_increment ,
biz int NOT NULL COMMENT 'biz code',
ei varchar(10) NOT NULL COMMENT 'first level event name',
sei varchar(10) NULL COMMENT 'second level event name',
uv long NULL COMMENT 'user count',
pv long NULL,
time_id varchar(10) NOT NULL,
update_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
)

数据源

1
2
3
4
5
6
7
8
9
10
11
12
CREATE TABLE orders (
biz int,
ei STRING,
sei STRING,
ui STRing,
etime timestamp(3),
watermark for etime as etime - interval '5' second
) WITH (
'connector.type' = 'filesystem',
'connector.path'='/Users/zhangzuofeng1/orders.csv',
'format.type'='csv'
)

假设选用Blink planner运行两个SQL:

  • 创建MySQL表的DDL语句
  • 从Source消费数据写入MySQL的DML语句
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
CREATE TABLE agg_result (
biz int,
ei STRING,
sei STRING,
uv bigint,
pv bigint,
time_id STRING
) WITH (
'connector.type' = 'jdbc',
'connector.driver' = 'com.mysql.cj.jdbc.Driver',
'connector.url' = 'jdbc:mysql://localhost:3306/fdata',
'connector.table' = 't_agg_result',
'connector.username' = 'root',
'connector.password' = 'abc.ABC.123'
)
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
INSERT INTO MyUserTable
SELECT
biz,
ei,
sei,
count(DISTINCT ui) AS uv,
count(*) AS pv,
DATE_FORMAT(TUMBLE_START(etime, INTERVAL '1' MINUTE), 'yyyyMMddHHmm') AS timeId
FROM orders
WHERE ui IS NOT NULL
GROUP BY biz,
ei,
sei,
TUMBLE(etime, INTERVAL '1' MINUTE)

insert into agg_result
SELECT
biz,
ei,
sei,
count(DISTINCT ui) AS uv,
count(*) AS pv,
DATE_FORMAT(TUMBLE_START(etime, INTERVAL '1' MINUTE), 'yyyyMMddHHmm') AS timeId
FROM orders
WHERE ui IS NOT NULL
GROUP BY biz,
ei,
sei,
TUMBLE(etime, INTERVAL '1' MINUTE)

执行下面MySQL

1
2
tEnv.sqlUpdate(ddlSql);
tEnv.sqlUpdate(dmlSql);
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
== Abstract Syntax Tree ==
LogicalProject(ei=[$0], sei=[$1], ui=[$2], etime=[$3])
+- LogicalFilter(condition=[<>($0, _UTF-16LE'a')])
+- LogicalTableScan(table=[[default_catalog, default_database, Unregistered_TableSource_1298380324, source: [CsvTableSource(read fields: ei, sei, ui, etime)]]])

== Optimized Logical Plan ==
Calc(select=[ei, sei, ui, etime], where=[<>(ei, _UTF-16LE'a':VARCHAR(10) CHARACTER SET "UTF-16LE")])
+- TableSourceScan(table=[[default_catalog, default_database, Unregistered_TableSource_1298380324, source: [CsvTableSource(read fields: ei, sei, ui, etime)]]], fields=[ei, sei, ui, etime])

== Physical Execution Plan ==
Stage 1 : Data Source
content : Source: Custom File source

Stage 2 : Operator
content : CsvTableSource(read fields: ei, sei, ui, etime)
ship_strategy : REBALANCE

Stage 3 : Operator
content : SourceConversion(table=[default_catalog.default_database.Unregistered_TableSource_1298380324, source: [CsvTableSource(read fields: ei, sei, ui, etime)]], fields=[ei, sei, ui, etime])
ship_strategy : FORWARD

Stage 4 : Operator
content : Calc(select=[ei, sei, ui, etime], where=[(ei <> _UTF-16LE'a':VARCHAR(10) CHARACTER SET "UTF-16LE")])
ship_strategy : FORWARD
1
{"nodes":[{"id":10,"type":"Source: Custom File source","pact":"Data Source","contents":"Source: Custom File source","parallelism":1},{"id":11,"type":"CsvTableSource(read fields: biz, ei, sei, ui, etime)","pact":"Operator","contents":"CsvTableSource(read fields: biz, ei, sei, ui, etime)","parallelism":12,"predecessors":[{"id":10,"ship_strategy":"REBALANCE","side":"second"}]},{"id":12,"type":"SourceConversion(table=[default_catalog.default_database.orders, source: [CsvTableSource(read fields: biz, ei, sei, ui, etime)]], fields=[biz, ei, sei, ui, etime])","pact":"Operator","contents":"SourceConversion(table=[default_catalog.default_database.orders, source: [CsvTableSource(read fields: biz, ei, sei, ui, etime)]], fields=[biz, ei, sei, ui, etime])","parallelism":12,"predecessors":[{"id":11,"ship_strategy":"FORWARD","side":"second"}]},{"id":13,"type":"WatermarkAssigner(rowtime=[etime], watermark=[(etime - 5000:INTERVAL SECOND)])","pact":"Operator","contents":"WatermarkAssigner(rowtime=[etime], watermark=[(etime - 5000:INTERVAL SECOND)])","parallelism":12,"predecessors":[{"id":12,"ship_strategy":"FORWARD","side":"second"}]},{"id":14,"type":"Calc(select=[biz, ei, sei, ui, etime], where=[ui IS NOT NULL])","pact":"Operator","contents":"Calc(select=[biz, ei, sei, ui, etime], where=[ui IS NOT NULL])","parallelism":12,"predecessors":[{"id":13,"ship_strategy":"FORWARD","side":"second"}]},{"id":16,"type":"GroupWindowAggregate(groupBy=[biz, ei, sei], window=[TumblingGroupWindow('w$, etime, 60000)], properties=[w$start, w$end, w$rowtime, w$proctime], select=[biz, ei, sei, COUNT(DISTINCT ui) AS uv, COUNT(*) AS pv, start('w$) AS w$start, end('w$) AS w$end, rowtime('w$) AS w$rowtime, proctime('w$) AS w$proctime])","pact":"Operator","contents":"GroupWindowAggregate(groupBy=[biz, ei, sei], window=[TumblingGroupWindow('w$, etime, 60000)], properties=[w$start, w$end, w$rowtime, w$proctime], select=[biz, ei, sei, COUNT(DISTINCT ui) AS uv, COUNT(*) AS pv, start('w$) AS w$start, end('w$) AS w$end, rowtime('w$) AS w$rowtime, proctime('w$) AS w$proctime])","parallelism":12,"predecessors":[{"id":14,"ship_strategy":"HASH","side":"second"}]},{"id":17,"type":"Calc(select=[biz, ei, sei, uv, pv, (w$start DATE_FORMAT _UTF-16LE'yyyyMMddHHmm') AS timeId])","pact":"Operator","contents":"Calc(select=[biz, ei, sei, uv, pv, (w$start DATE_FORMAT _UTF-16LE'yyyyMMddHHmm') AS timeId])","parallelism":12,"predecessors":[{"id":16,"ship_strategy":"FORWARD","side":"second"}]},{"id":18,"type":"SinkConversionToTuple2","pact":"Operator","contents":"SinkConversionToTuple2","parallelism":12,"predecessors":[{"id":17,"ship_strategy":"FORWARD","side":"second"}]},{"id":19,"type":"Sink: JDBCUpsertTableSink(biz, ei, sei, uv, pv, time_id)","pact":"Data Sink","contents":"Sink: JDBCUpsertTableSink(biz, ei, sei, uv, pv, time_id)","parallelism":12,"predecessors":[{"id":18,"ship_strategy":"FORWARD","side":"second"}]}]}

Parser: Provides methods for parsing SQL objects from a SQL string.

org.apache.flink.table.planner.delegation.PlannerBase

org.apache.flink.table.planner.operations.SqlToOperationConverter#convert

SqlNode 转换为 Operation

  1. 结合元数据验证

https://cloud.tencent.com/developer/article/1803116

https://cloud.tencent.com/developer/article/1697401

https://matt33.com/2019/10/20/paper-flink-snapshot/

Asynchronous Barrier Snapshot(ABS)

默认情况下,Checkpoint机制是关闭的,需要调用env.enableCheckpointing(n)来开启,每隔n毫秒进行一次Checkpoint。Checkpoint是一种负载较重的任务,如果状态比较大,同时n值又比较小,那可能一次Checkpoint还没完成,下次Checkpoint已经被触发,占用太多本该用于正常数据处理的资源。增大n值意味着一个作业的Checkpoint次数更少,整个作业用于进行Checkpoint的资源更小,可以将更多的资源用于正常的流数据处理。同时,更大的n值意味着重启后,整个作业需要从更长的Offset开始重新处理数据。

实时流处理系统反压机制(BackPressure)综述[转]

发表于 2018-11-15 | 更新于 2018-12-03 | 分类于 BigData | 阅读次数 333

本文主要关于实时流处理系统反压机制,最近看到反压问题看到此文章很好,在此分享并mark一下。

本文主要关于实时流处理系统反压机制,最近看到反压问题看到此文章很好,在此分享并mark一下。
(¬_¬)ノ最近菜叶子没自己写见谅。
本文转自 实时流处理系统反压机制(BackPressure)综述
https://blog.csdn.net/qq_21125183/article/details/80708142
开启Back Pressure使生产环境的Spark Streaming应用更稳定、有效

反压机制(BackPressure)被广泛应用到实时流处理系统中,流处理系统需要能优雅地处理反压(backpressure)问题。
反压通常产生于这样的场景:短时负载高峰导致系统接收数据的速率远高于它处理数据的速率。
许多日常问题都会导致反压,例如,垃圾回收停顿可能会导致流入的数据快速堆积,或者遇到大促或秒杀活动导致流量陡增。
反压如果不能得到正确的处理,可能会导致资源耗尽甚至系统崩溃。反压机制就是指系统能够自己检测到被阻塞的Operator,然后系统自适应地降低源头或者上游的发送速率。

目前主流的流处理系统 Apache Storm、JStorm、Spark Streaming、S4、Apache Flink、Twitter Heron都采用反压机制解决这个问题,不过他们的实现各自不同。

实时流处理系统反压机制01

不同的组件可以不同的速度执行(并且每个组件中的处理速度随时间改变)。 例如,考虑一个工作流程,或由于数据倾斜或任务调度而导致数据被处理十分缓慢。
在这种情况下,如果上游阶段不减速,将导致缓冲区建立长队列(队列占用内存、硬盘空间,节点负载加重),或导致系统丢弃元组。
如果元组在中途丢弃,那么效率可能会有损失,因为已经为这些元组产生的计算被浪费了。
并且在一些流处理系统中比如Strom,会将这些丢失的元组重新发送,这样会导致数据的一致性问题(at least once语义),并且还会导致某些Operator状态叠加。
进而整个程序输出结果不准确。第二由于系统接收数据的速率是随着时间改变的,短时负载高峰导致系统接收数据的速率远高于它处理数据的速率的情况,也会导致Tuple在中途丢失。
所以实时流处理系统必须能够解决发送速率远大于系统能处理速率这个问题,大多数实时流处理系统采用反压(BackPressure)机制解决这个问题。

下面我们就来介绍一下不同的实时流处理系统采用的反压机制:

Strom 反压机制

Storm 1.0 以前的反压机制

对于开启了acker机制的storm程序,可以通过设置conf.setMaxSpoutPending参数来实现反压效果,如果下游组件(bolt)处理速度跟不上导致spout发送的tuple没有及时确认的数超过了参数设定的值,spout会停止发送数据,这种方式的缺点是很难调优conf.setMaxSpoutPending参数的设置以达到最好的反压效果,设小了会导致吞吐上不去,设大了会导致worker OOM;有震荡,数据流会处于一个颠簸状态,效果不如逐级反压;另外对于关闭acker机制的程序无效;

Storm Automatic Backpressure

新的storm自动反压机制(Automatic Back Pressure)通过监控bolt中的接收队列的情况,当超过高水位值时专门的线程会将反压信息写到 Zookeeper ,Zookeeper上的watch会通知该拓扑的所有Worker都进入反压状态,最后Spout降低tuple发送的速度。

实时流处理系统反压机制02

每个Executor都有一个接受队列和发送队列用来接收Tuple和发送Spout或者Bolt生成的Tuple元组。每个Worker进程都有一个单的的接收线程监听接收端口。
它从每个网络上进来的消息发送到Executor的接收队列中。Executor接收队列存放Worker或者Worker内部其他Executor发过来的消息。
Executor工作线程从接收队列中拿出数据,然后调用execute方法,发送Tuple到Executor的发送队列。
Executor的发送线程从发送队列中获取消息,按照消息目的地址选择发送到Worker的传输队列中或者其他Executor的接收队列中。
最后Worker的发送线程从传输队列中读取消息,然后将Tuple元组发送到网络中。

  1. 当Worker进程中的Executor线程发现自己的接收队列满了时,也就是接收队列达到high watermark的阈值后,因此它会发送通知消息到背压线程。
  2. 背压线程将当前worker进程的信息注册到Zookeeper的Znode节点中。具体路径就是 /Backpressure/topo1/wk1
  3. Zookeepre的Znode Watcher监视/Backpreesure/topo1下的节点目录变化情况,如果发现目录增加了znode节点说明或者其他变化。这就说明该Topo1需要反压控制,然后它会通知Topo1所有的Worker进入反压状态。
  4. 最终Spout降低tuple发送的速度。

JStorm 反压机制

JStorm做了两级的反压,第一级和Jstorm类似,通过执行队列来监测,但是不会通过ZK来协调,而是通过Topology Master来协调。
在队列中会标记high water mark和low water mark,当执行队列超过high water mark时,就认为bolt来不及处理,则向TM发一条控制消息,上游开始减慢发送速率,直到下游低于low water mark时解除反压。

此外,在Netty层也做了一级反压,由于每个Worker Task都有自己的发送和接收的缓冲区,可以对缓冲区设定限额、控制大小,如果spout数据量特别大,缓冲区填满会导致下游bolt的接收缓冲区填满,造成了反压。

实时流处理系统反压机制03

限流机制:jstorm的限流机制, 当下游bolt发生阻塞时, 并且阻塞task的比例超过某个比例时(现在默认设置为0.1),触发反压

限流方式:计算阻塞Task的地方执行线程执行时间,Spout每发送一个tuple等待相应时间,然后讲这个时间发送给Spout, 于是, spout每发送一个tuple,就会等待这个执行时间。

Task阻塞判断方式:在jstorm 连续4次采样周期中采样,队列情况,当队列超过80%(可以设置)时,即可认为该task处在阻塞状态。

SparkStreaming 反压机制

为什么引入反压机制Backpressure

默认情况下,Spark Streaming通过Receiver以生产者生产数据的速率接收数据,计算过程中会出现batch processing time > batch interval的情况,其中batch processing time 为实际计算一个批次花费时间, batch interval为Streaming应用设置的批处理间隔。
这意味着Spark Streaming的数据接收速率高于Spark从队列中移除数据的速率,也就是数据处理能力低,在设置间隔内不能完全处理当前接收速率接收的数据。如果这种情况持续过长的时间,会造成数据在内存中堆积,导致Receiver所在Executor内存溢出等问题(如果设置StorageLevel包含disk, 则内存存放不下的数据会溢写至disk, 加大延迟)。
Spark 1.5以前版本,用户如果要限制Receiver的数据接收速率,可以通过设置静态配制参数“spark.streaming.receiver.maxRate”的值来实现,此举虽然可以通过限制接收速率,来适配当前的处理能力,防止内存溢出,但也会引入其它问题。比如:producer数据生产高于maxRate,当前集群处理能力也高于maxRate,这就会造成资源利用率下降等问题。为了更好的协调数据接收速率与资源处理能力,Spark Streaming 从v1.5开始引入反压机制(back-pressure),通过动态控制数据接收速率来适配集群数据处理能力。

反压机制Backpressure

Spark Streaming Backpressure: 根据JobScheduler反馈作业的执行信息来动态调整Receiver数据接收率。通过属性“spark.streaming.backpressure.enabled”来控制是否启用backpressure机制,默认值false,即不启用。

1
sparkConf.set("spark.streaming.backpressure.enabled",”true”)

SparkStreaming 架构图如下所示:

实时流处理系统反压机制04

SparkStreaming 反压过程执行如下图所示:

在原架构的基础上加上一个新的组件RateController,这个组件负责监听“OnBatchCompleted”事件,然后从中抽取processingDelayschedulingDelay信息. Estimator依据这些信息估算出最大处理速度(rate),最后由基于Receiver的Input Stream将rate通过ReceiverTracker与ReceiverSupervisorImpl转发给BlockGenerator(继承自RateLimiter).

实时流处理系统反压机制05

direct模式-BackPressure(此部分详细说明了direct模式接收:转自-开启Back Pressure使生产环境的Spark Streaming应用更稳定、有效)

当Spark Streaming与Kafka使用Direct API集群时,我们可以很方便的去控制最大数据摄入量–通过一个被称作spark.streaming.kafka.maxRatePerPartition的参数。根据文档描述,他的含义是:Direct API读取每一个Kafka partition数据的最大速率(每秒读取的消息量)。
配置项spark.streaming.kafka.maxRatePerPartition,对防止流式应用在下边两种情况下出现流量过载时尤其重要:
1.Kafka Topic中有大量未处理的消息,并且我们设置是Kafka auto.offset.reset参数值为smallest,他可以防止第一个批次出现数据流量过载情况。
2.当Kafka 生产者突然飙升流量的时候,他可以防止批次处理出现数据流量过载情况。

但是,配置Kafka每个partition每批次最大的摄入量是个静态值,也算是个缺点。随着时间的变化,在生产环境运行了一段时间的Spark Streaming应用,每批次每个Kafka partition摄入数据最大量的最优值也是变化的。有时候,是因为消息的大小会变,导致数据处理时间变化。有时候,是因为流计算所使用的多租户集群会变得非常繁忙,比如在白天时候,一些其他的数据应用(例如Impala/Hive/MR作业)竞争共享的系统资源时(CPU/内存/网络/磁盘IO)。
背压机制可以解决该问题。背压机制是呼声比较高的功能,他允许根据前一批次数据的处理情况,动态、自动的调整后续数据的摄入量,这样的反馈回路使得我们可以应对流式应用流量波动的问题。
Spark Streaming的背压机制是在Spark1.5版本引进的,我们可以添加如下代码启用改功能:

1
2
sparkConf.set("spark.streaming.backpressure.enabled",”true”)

那应用启动后的第一个批次流量怎么控制呢?因为他没有前面批次的数据处理时间,所以没有参考的数据去评估这一批次最优的摄入量。在Spark官方文档中有个被称作spark.streaming.backpressure.initialRate的配置,看起来是控制开启背压机制时初始化的摄入量。其实不然,该参数只对receiver模式起作用,并不适用于direct模式。推荐的方法是使用spark.streaming.kafka.maxRatePerPartition控制背压机制起作用前的第一批次数据的最大摄入量。我通常建议设置spark.streaming.kafka.maxRatePerPartition的值为最优估计值的1.5到2倍,让背压机制的算法去调整后续的值。请注意,spark.streaming.kafka.maxRatePerPartition的值会一直控制最大的摄入量,所以背压机制的算法值不会超过他。
另一个需要注意的是,在第一个批次处理完成前,紧接着的批次都将使用spark.streaming.kafka.maxRatePerPartition的值作为摄入量。通过Spark UI可以看到,批次间隔为5s,当批次调度延迟31秒时候,前7个批次的摄入量是20条记录。直到第八个批次,背压机制起作用时,摄入量变为5条记录。

Heron 反压机制

实时流处理系统反压机制06

当下游处理速度跟不上上游发送速度时,一旦StreamManager 发现一个或多个Heron Instance 速度变慢,立刻对本地spout进行降级,降低本地Spout发送速度, 停止从这些spout读取数据。并且受影响的StreamManager 会发送一个特殊的start backpressure message 给其他的StreamManager ,要求他们对spout进行本地降级。 当其他StreamManager 接收到这个特殊消息时,他们通过不读取当地Spout中的Tuple来进行降级。一旦出问题的Heron Instance 恢复速度后,本地的SM 会发送stop backpressure message 解除降级。

很多Socket Channel与应用程序级别的Buffer相关联,该缓冲区由high watermark 和low watermark组成。 当缓冲区大小达到high watermark时触发反压,并保持有效,直到缓冲区大小低于low watermark。 此设计的基本原理是防止拓扑在进入和退出背压缓解模式之间快速振荡。

Flink 反压机制

Flink 没有使用任何复杂的机制来解决反压问题,因为根本不需要那样的方案!它利用自身作为纯数据流引擎的优势来优雅地响应反压问题。下面我们会深入分析 Flink 是如何在 Task 之间传输数据的,以及数据流如何实现自然降速的。 Flink 在运行时主要由 operators 和 streams 两大组件构成。每个 operator 会消费中间态的流,并在流上进行转换,然后生成新的流。对于 Flink 的网络机制一种形象的类比是,Flink 使用了高效有界的分布式阻塞队列,就像 Java 通用的阻塞队列(BlockingQueue)一样。还记得经典的线程间通信案例:生产者消费者模型吗?使用 BlockingQueue 的话,一个较慢的接受者会降低发送者的发送速率,因为一旦队列满了(有界队列)发送者会被阻塞。Flink 解决反压的方案就是这种感觉。 在 Flink 中,这些分布式阻塞队列就是这些逻辑流,而队列容量是通过缓冲池来(LocalBufferPool)实现的。每个被生产和被消费的流都会被分配一个缓冲池。缓冲池管理着一组缓冲(Buffer),缓冲在被消费后可以被回收循环利用。这很好理解:你从池子中拿走一个缓冲,填上数据,在数据消费完之后,又把缓冲还给池子,之后你可以再次使用它。

如下图所示展示了 Flink 在网络传输场景下的内存管理。网络上传输的数据会写到 Task 的 InputGate(IG) 中,经过 Task 的处理后,再由 Task 写到 ResultPartition(RS) 中。每个 Task 都包括了输入和输入,输入和输出的数据存在 Buffer 中(都是字节数据)。Buffer 是 MemorySegment 的包装类。

实时流处理系统反压机制07

  1. TaskManager(TM)在启动时,会先初始化NetworkEnvironment对象,TM 中所有与网络相关的东西都由该类来管理(如 Netty 连接),其中就包括NetworkBufferPool。根据配置,Flink 会在 NetworkBufferPool 中生成一定数量(默认2048个)的内存块 MemorySegment(关于 Flink 的内存管理,后续文章会详细谈到),内存块的总数量就代表了网络传输中所有可用的内存。NetworkEnvironment 和 NetworkBufferPool 是 Task 之间共享的,每个 TM 只会实例化一个。
  2. Task 线程启动时,会向 NetworkEnvironment 注册,NetworkEnvironment 会为 Task 的 InputGate(IG)和 ResultPartition(RP) 分别创建一个 LocalBufferPool(缓冲池)并设置可申请的 MemorySegment(内存块)数量。IG 对应的缓冲池初始的内存块数量与 IG 中 InputChannel 数量一致,RP 对应的缓冲池初始的内存块数量与 RP 中的 ResultSubpartition 数量一致。不过,每当创建或销毁缓冲池时,NetworkBufferPool 会计算剩余空闲的内存块数量,并平均分配给已创建的缓冲池。注意,这个过程只是指定了缓冲池所能使用的内存块数量,并没有真正分配内存块,只有当需要时才分配。为什么要动态地为缓冲池扩容呢?因为内存越多,意味着系统可以更轻松地应对瞬时压力(如GC),不会频繁地进入反压状态,所以我们要利用起那部分闲置的内存块。
  3. 在 Task 线程执行过程中,当 Netty 接收端收到数据时,为了将 Netty 中的数据拷贝到 Task 中,InputChannel(实际是 RemoteInputChannel)会向其对应的缓冲池申请内存块(上图中的①)。如果缓冲池中也没有可用的内存块且已申请的数量还没到池子上限,则会向 NetworkBufferPool 申请内存块(上图中的②)并交给 InputChannel 填上数据(上图中的③和④)。如果缓冲池已申请的数量达到上限了呢?或者 NetworkBufferPool 也没有可用内存块了呢?这时候,Task 的 Netty Channel 会暂停读取,上游的发送端会立即响应停止发送,拓扑会进入反压状态。当 Task 线程写数据到 ResultPartition 时,也会向缓冲池请求内存块,如果没有可用内存块时,会阻塞在请求内存块的地方,达到暂停写入的目的。
  4. 当一个内存块被消费完成之后(在输入端是指内存块中的字节被反序列化成对象了,在输出端是指内存块中的字节写入到 Netty Channel 了),会调用 Buffer.recycle() 方法,会将内存块还给 LocalBufferPool (上图中的⑤)。如果LocalBufferPool中当前申请的数量超过了池子容量(由于上文提到的动态容量,由于新注册的 Task 导致该池子容量变小),则LocalBufferPool会将该内存块回收给 NetworkBufferPool(上图中的⑥)。如果没超过池子容量,则会继续留在池子中,减少反复申请的开销。

下面这张图简单展示了两个 Task 之间的数据传输以及 Flink 如何感知到反压的:

实时流处理系统反压机制08

  1. 记录“A”进入了 Flink 并且被 Task 1 处理。(这里省略了 Netty 接收、反序列化等过程)
  2. 记录被序列化到 buffer 中。
  3. 该 buffer 被发送到 Task 2,然后 Task 2 从这个 buffer 中读出记录。

不要忘了:记录能被 Flink 处理的前提是,必须有空闲可用的 Buffer。

结合上面两张图看:Task 1 在输出端有一个相关联的 LocalBufferPool(称缓冲池1),Task 2 在输入端也有一个相关联的 LocalBufferPool(称缓冲池2)。如果缓冲池1中有空闲可用的 buffer 来序列化记录 “A”,我们就序列化并发送该 buffer。

这里我们需要注意两个场景:

  • 本地传输:如果 Task 1 和 Task 2 运行在同一个 worker 节点(TaskManager),该 buffer 可以直接交给下一个 Task。一旦 Task 2 消费了该 buffer,则该 buffer 会被缓冲池1回收。如果 Task 2 的速度比 1 慢,那么 buffer 回收的速度就会赶不上 Task 1 取 buffer 的速度,导致缓冲池1无可用的 buffer,Task 1 等待在可用的 buffer 上。最终形成 Task 1 的降速。
  • 远程传输:如果 Task 1 和 Task 2 运行在不同的 worker 节点上,那么 buffer 会在发送到网络(TCP Channel)后被回收。在接收端,会从 LocalBufferPool 中申请 buffer,然后拷贝网络中的数据到 buffer 中。如果没有可用的 buffer,会停止从 TCP 连接中读取数据。在输出端,通过 Netty 的水位值机制来保证不往网络中写入太多数据(后面会说)。如果网络中的数据(Netty输出缓冲中的字节数)超过了高水位值,我们会等到其降到低水位值以下才继续写入数据。这保证了网络中不会有太多的数据。如果接收端停止消费网络中的数据(由于接收端缓冲池没有可用 buffer),网络中的缓冲数据就会堆积,那么发送端也会暂停发送。另外,这会使得发送端的缓冲池得不到回收,writer 阻塞在向 LocalBufferPool 请求 buffer,阻塞了 writer 往 ResultSubPartition 写数据。

这种固定大小缓冲池就像阻塞队列一样,保证了 Flink 有一套健壮的反压机制,使得 Task 生产数据的速度不会快于消费的速度。我们上面描述的这个方案可以从两个 Task 之间的数据传输自然地扩展到更复杂的 pipeline 中,保证反压机制可以扩散到整个 pipeline。

反压实验

另外,官方博客中为了展示反压的效果,给出了一个简单的实验。下面这张图显示了:随着时间的改变,生产者(黄色线)和消费者(绿色线)每5秒的平均吞吐与最大吞吐(在单一JVM中每秒达到8百万条记录)的百分比。我们通过衡量task每5秒钟处理的记录数来衡量平均吞吐。该实验运行在单 JVM 中,不过使用了完整的 Flink 功能栈。

实时流处理系统反压机制09

首先,我们运行生产task到它最大生产速度的60%(我们通过Thread.sleep()来模拟降速)。消费者以同样的速度处理数据。然后,我们将消费task的速度降至其最高速度的30%。你就会看到背压问题产生了,正如我们所见,生产者的速度也自然降至其最高速度的30%。接着,停止消费task的人为降速,之后生产者和消费者task都达到了其最大的吞吐。接下来,我们再次将消费者的速度降至30%,pipeline给出了立即响应:生产者的速度也被自动降至30%。最后,我们再次停止限速,两个task也再次恢复100%的速度。总而言之,我们可以看到:生产者和消费者在 pipeline 中的处理都在跟随彼此的吞吐而进行适当的调整,这就是我们希望看到的反压的效果。

在 Storm/JStorm 中,只要监控到队列满了,就可以记录下拓扑进入反压了。但是 Flink 的反压太过于天然了,导致我们无法简单地通过监控队列来监控反压状态。Flink 在这里使用了一个 trick 来实现对反压的监控。如果一个 Task 因为反压而降速了,那么它会卡在向 LocalBufferPool 申请内存块上。那么这时候,该 Task 的 stack trace 就会长下面这样:

1
2
3
4
java.lang.Object.wait(Native Method)
o.a.f.[...].LocalBufferPool.requestBuffer(LocalBufferPool.java:163)
o.a.f.[...].LocalBufferPool.requestBufferBlocking(LocalBufferPool.java:133) <--- BLOCKING request
[...]

那么事情就简单了。通过不断地采样每个 task 的 stack trace 就可以实现反压监控。

实时流处理系统反压机制10

Flink 的实现中,只有当 Web 页面切换到某个 Job 的 Backpressure 页面,才会对这个 Job 触发反压检测,因为反压检测还是挺昂贵的。JobManager 会通过 Akka 给每个 TaskManager 发送TriggerStackTraceSample消息。默认情况下,TaskManager 会触发100次 stack trace 采样,每次间隔 50ms(也就是说一次反压检测至少要等待5秒钟)。并将这 100 次采样的结果返回给 JobManager,由 JobManager 来计算反压比率(反压出现的次数/采样的次数),最终展现在 UI 上。UI 刷新的默认周期是一分钟,目的是不对 TaskManager 造成太大的负担。

总结

Flink不需要一种特殊的机制来处理反压,因为Flink 中的数据传输相当于已经提供了应对反压的机制。因此,Flink 所能获得的最大吞吐量由其 pipeline 中最慢的组件决定。相对于 Storm/JStorm 的实现,Flink 的实现更为简洁优雅,源码中也看不见与反压相关的代码,无需 Zookeeper/TopologyMaster 的参与也降低了系统的负载,也利于对反压更迅速的响应。

本文转自 实时流处理系统反压机制(BackPressure)综述
https://blog.csdn.net/qq_21125183/article/details/80708142
开启Back Pressure使生产环境的Spark Streaming应用更稳定、有效