Flink Clients

先了解整个流程,有了全局视角之后,后续会详述细节。
先来看看 flink datastream 任务的执行过程:
DataStream:
使用时要在 flink datastream api 提供的各种 udf(比如 flatmap,keyedProcessFunction 等)中自定义处理逻辑,具体的业务执行逻辑都是敲代码、 java 文件写的,然后编译在 jvm 中执行,就和一个普通的 main 函数应用一模一样的流程。因为代码执行逻辑都是自己写的,所以这一部分相对好理解。
SQL:
Java 编译器不能识别和编译一条 SQL 进行执行,那么一条 SQL 是咋执行的呢?
我们逆向思维进行考虑,如果想让一条 Flink SQL 按照我们的预期在 jvm 中执行,需要哪些过程。
如下图所示,描绘了上述逻辑:

12
那么这个和 flink 实际实现有啥异同呢?
flink 大致是这样做的,虽在 flink 本身的中间还有一些其他的流程,后来的版本也不是基于 datastream,但是整体的处理逻辑还是和上述一致的。
所以不了解整体流程的同学可以先按照上述流程进行理解。
按照 博主的脑洞 来总结一条 sql 的使命就是:sql -> AST -> codegen(java code) -> 让我们 run 起来好吗

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

28
标准的一条 flink sql 运行起来的流程如下:
Notes:刚开始对其中的 SqlNode,RelNode 概念可能比较模糊。先理解整个流程,后续会详细介绍这些概念。
sql 解析阶段:calcite parser 解析(sql -> AST,AST 即 SqlNode Tree)
SqlNode 验证阶段:calcite validator 校验(SqlNode -> SqlNode,语法、表达式、表信息)
语义分析阶段:SqlNode 转换为 RelNode,RelNode 即 Logical Plan(SqlNode -> RelNode)
优化阶段:calcite optimizer 优化(RelNode -> RelNode,剪枝、谓词下推等)
物理计划生成阶段:Logical Plan 转换为 Physical Plan(等同于 RelNode 转换成 DataSet\DataStream API)
后续的运行逻辑与 datastream 一致
可以发现 flink 的实现 比 博主的脑洞 整体主要框架上面是一致的。多出来的部分主要是 SqlNode 验证阶段,优化阶段。
大致了解了 一条 flink sql 的运行流程 之后,我们来看看 calcite 这玩意到底在 flink 里干了些啥。
根据上文总结来说 calcite 在 flink sql 中担当了 sql 解析、验证、优化功能。

30
看着 calcite 干了这么多事,那 calcite 是个啥东东,它的定位是啥?
calcite 是一个动态数据的管理框架,它可以用来构建数据库系统的不同的解析的模块,但是它不包含数据存储数据处理等功能。
calcite 的目标是一种方案,适应所有的需求场景,希望能为不同计算平台和数据源提供统一的 sql 解析引擎,但是它只是提供查询引擎,而没有真正的去存储这些数据。

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

4
简单来说的话,可以先理解为 calcite 具有这几个功能(当然还有其他很牛逼的功能,感兴趣可以自查官网)。
SqlSelect、SqlNode),就可以根据这些对象做具体逻辑处理了。举个例子,如下图,一条简单的 select c,d from source where a = '6' sql,经过 calcite 的解析之后,就可以得到 AST model(SqlNode)。可以看到有 SqlSelect、SqlIdentifier、SqlIdentifier、SqlCharStringLiteral。上面的这些能力整体组成如下图所示:

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

2
1 | 重中之重,在了解原理之前,先跑起来是王道,也会帮助我们逐步理解。 |
官网已经有一个 csv 的案例了。感兴趣的可以直达 https://calcite.apache.org/docs/tutorial.html 。
跑完一个 csv demo,在详细了解 calcite 之前还需要了解下 sql,calcite 的支柱:关系代数。
sql 是基于关系代数的查询语言,是关系代数在工程上的一种很好的实现方案。在工程中,关系代数难表达,但是 sql 就易于理解。关系代数和 sql 的关系如下。
可以将一条 sql 解析为一个关系代数表达式的组合。在 sql 中的操作都可以转化成关系代数的表达式。
sql 的执行优化(所有的优化的前提都是优化前和优化后最终执行结果相同,即等价交换)是基于关系代数运算的。
总结下,有哪些常用的关系代数:

50
关系代数等价变换是 calcite optimizer 的基础理论。
下面是一些等价变换的例子。
1.连接(⋈),笛卡尔积(×)的交换律

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

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

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

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

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

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

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

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

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

然后看一个基于关系代数优化的实际 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 转为关系代数的语法树。

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

47

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

48

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

49







关于关系代数我们就有了大致的了解。
除此之外,对于更深入了解 flink sql,calcite 而言,我们还需要了解一下在 calcite 代码体系中有哪些重要 model。
calcite 中有两个最最基础、重要的 model 在我们理解 flink sql 解析流程时需要知道的。
举个例子来说明下,下面这条 flink sql,经过解析之后的 SqlNode,RelNode 如下图:
1 | SELECT |

62
可以看到 SqlNode 包含的内容是 sql 的层次结构,包括 selectList,from,where,group by 等。
RelNode 包含的是关系代数的层次结构,每一层都有一个 input 来承接。结合上面优化案例的树状结构一样。

63

29
如上图所示,此处我们结合上节介绍的 calcite 的 model,以及 flink sql 的实现来走一遍其处理流程:
1 | SELECT |
其中前三步解析和转化,都在 在执行 TableEnvironment#sqlQuery 进行。
最后一步优化,在执行 sink 操作时进行,即在这个例子中是 tEnv.toRetractStream(result, Row.class)。
源码公众号后台回复flink sql 知其所以然(六)| flink sql 约会 calcite获取。
sql 解析阶段使用 Sql Parser 将 sql 解析为 SqlNode。这一步在执行 TableEnvironment#sqlQuery 进行。




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

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



可以从上图看到 flink sql 校验器的具体实现类是 FlinkCalciteSqlValidator,其中包含了元数据信息,从而可以进行元数据信息检查。
这一步就是将 SqlNode 转换成 RelNode,也就是生成相应的关系代数层面的逻辑(这里一般都叫做逻辑计划:Logical Plan)。



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


此处以 calcite parser 举例说明,其模块为什么这通用?其他的模块都是类似的方式。
先说结论:因为 calcite parser 模块提供了接口,具体的 parse 逻辑、规则是可以根据用户自定义进行配置的。大家可以看下图,博主画出了一张图进行详述。

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 之后查看。

31
javacc 是一个用 java 开发的最受欢迎的语法分析生成器。这个分析生成器工具可以读取上下文无关且有着特殊意义的语法并把它转换成可以识别且匹配该语法的 java 程序。它是 100% 的纯 java 代码,可以在多种平台上运行。
简单解释 javacc 就是它是一个通用的语法分析生产器,用户可以使用 javacc 任意定义一套 DSL 及解析器。
举个例子,如果哪天你觉得 sql 也不够简洁通用,你可以使用 javacc 自己定义一套更简洁的 user-define-ql。然后使用 javacc 作为你的 user-define-ql 的解析器。是不是很流批,可以自己去搞编译器了。
这里不介绍具体的 javacc 语法,直接以官网的 Simple1.jj 为案例。详细语法和功能可以参考官网(https://javacc.github.io/javacc/) 或者一下博客。
Simple1.jj 是用于识别一系列的 {相同数量的花括号},之后跟着 0 个或多个行终结符。

7
下面是合法的字符串例子:
{},{{{{{}}}}},etc.
下面是不合法的字符串例子:
{{{{,{}{},{}},{{}{}},etc.
接下来让我们实际将 Simple1.jj 编译生成具体的规则代码。
在 pom 中加入 javacc build 插件:
1 | <plugin> |
在 compile 之后,就会在 generated-sources 下生成代码:

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

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



关于 javacc 基本上就了解个大概了。
感兴趣的可以尝试自定义一个编译器。

5
fmpp 就是一个基于 freemarker 的模板生产器。用户可以统一管理自己的变量,然后用 ftl 模板 + 变量 生成对应的最终文件。在 calcite 中使用 fmpp 作为变量 + 模板的统一管理器。然后基于 fmpp 来生成对应的 Parser.jj 文件。
博主画了一张图,包含了其中重要组件之间的依赖关系。

3
你没猜错,还是上面那些流程,fmpp(Parser.jj 模板生成) -> javacc(Parser 生成) -> calcite。
在介绍 Parser 生成流程之前,先看看 flink 最终生成的 Parser:FlinkSqlParserImpl (此处使用 Blink Planner)。
以下面这个案例出发(代码基于 flink 1.13.1 版本):
1 | public class ParserTest { |
debug 过程如之前分析 sql -> SqlNode 过程所示,如下图直接定位到 SqlParser:

21
如上图可以看到具体的 Parser 就是 FlinkSqlParserImpl。
定位到具体的代码如下图所示(flink-table-palnner-blink-2.11-1.13.1.jar)。

34
最终 parse 的结果 SqlNode 如下图。

22


再来看看 FlinkSqlParserImpl 是怎么使用 calcite 生成的。
具体到 flink 中的实现,位于源码中的 flink-table.flink-sql-parser 模块(源码基于 flink 1.13.1)。
flink 是依赖 maven 插件实现的上面的整体流程。

14
接下来看看整个 Parser 生成流程。
使用 maven-dependency-plugin 将 calcite 解压到 flink 项目 build 目录下。
1 | <plugin> |

15
使用 maven-resources-plugin 将 Parser.jj 代码生成。
1 | <plugin> |

16
使用 javacc 将根据 Parser.jj 文件生成 Parser。
1 | <plugin> |

17
最终生成的 Parser 就是 FlinkSqlParserImpl。

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

35
本文从一个调试时候常见的异常 “TimeoutException: Heartbeat of TaskManager timed out”切入,为大家剖析Flink的心跳机制。文中代码基于Flink 1.10。
大家如果经常调试Flink,当进入断点看到了堆栈和变量内容之后,你容易陷入了沉思。当你发现了问题可能所在,高兴的让程序Resume的时候,你发现程序无法运行,有如下提示:
1 | Caused by: java.util.concurrent.TimeoutException: Heartbeat of TaskManager with id 93aa1740-cd2c-4032-b74a-5f256edb3217 timed out. |
这实在是很郁闷的事情。作为程序猿不能忍啊,既然异常提示中有 Heartbeat 字样,于是我们就来一起看看Flink的心跳机制,看看有没有可以修改的途径。
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形式出现,可能本文中会混淆这两个词,请大家谅解):
TaskManager:类似Spark的executor,会跑多个线程的task、数据缓存与交换。Flink 架构遵循 Master - Slave 架构设计原则,JobMaster 为 Master 节点,TaskManager 为Slave节点。
这四大组件彼此之间的通信需要依赖RPC实现。
Flink底层RPC基于Akka实现。Akka是一个开发并发、容错和可伸缩应用的框架。它是Actor Model的一个实现,和Erlang的并发模型很像。在Actor模型中,所有的实体被认为是独立的actors。actors和其他actors通过发送异步消息通信。
Actor模型的强大来自于异步。它也可以显式等待响应,这使得可以执行同步操作。但是强烈不建议同步消息,因为它们限制了系统的伸缩性。
RPC作用是:让异步调用看起来像同步调用。
Flink基于Akka构建了其底层通信系统,引入了RPC调用,各节点通过GateWay方式回调,隐藏通信组件的细节,实现解耦。Flink整个通信框架的组件主要由RpcEndpoint、RpcService、RpcServer、AkkaInvocationHandler、AkkaRpcActor等构成。
RPC相关的主要接口如下:
RpcEndpoint是Flink RPC终端的基类,所有提供远程过程调用的分布式组件必须扩展RpcEndpoint,其功能由RpcService支持。
RpcEndpoint的子类只有四类组件:Dispatcher,JobMaster,ResourceManager,TaskExecutor,即Flink中只有这四个组件有RPC的能力,换句话说只有这四个组件有RPC的这个需求。
每个RpcEndpoint对应了一个路径(endpointId和actorSystem共同确定),每个路径对应一个Actor,其实现了RpcGateway接口,
RpcServer是RpcEndpoint的成员变量,为RpcService提供RPC服务/连接远程Server,其只有一个子类实现:AkkaRpcService(可见目前Flink的通信方式依然是Akka)。
RpcServer用于启动和连接到RpcEndpoint, 连接到rpc服务器将返回一个RpcGateway,可用于调用远程过程。
Flink四大组件Dispatcher,JobMaster,ResourceManager,TaskExecutor,都是RpcEndpoint的实现,所以构建四大组件时,同步需要初始化RpcServer。如JobManager的构造方式,第一个参数就是需要知道RpcService。
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调用。
常见的心跳检测有两种:
Flink实现的是第二种方案。
Flink的心跳机制代码在:
1 | Flink-master/flink-runtime/src/main/java/org/apache/flink/runtime/heartbeat |
四个接口:
1 | HeartbeatListener.java HeartbeatManager.java HeartbeatTarget.java HeartbeatMonitor.java |
以及如下几个类:
1 | HeartbeatManagerImpl.java HeartbeatManagerSenderImpl.java HeartbeatMonitorImpl.java |
Flink集群有多种业务流程,比如Resource Manager, Task Manager, Job Manager。每种业务流程都有自己的心跳机制。Flink的心跳机制只是提供接口和基本功能,具体业务功能由各业务流程自己实现。
我们首先设定 心跳系统中有两种节点:sender和receiver。心跳机制是sender和receivers彼此相互检测。但是检测动作是Sender主动发起,即Sender主动发送请求探测receiver是否存活,因为Sender已经发送过来了探测心跳请求,所以这样receiver同时也知道Sender是存活的,然后Reciver给Sender回应一个心跳表示自己也是活着的。
因为Flink的几个名词和我们常见概念有所差别,所以流程上需要大家仔细甄别,即:
HeartbeatTarget是对监控目标的抽象。心跳机制在行为上而言有两种动作:
HeartbeatTarget的函数就是这两个动作:
这两个函数的参数也很简单:分别是请求的发送放和接收方,还有Payload载荷。对于一个确定节点而言,接收的和发送的载荷是同一类型的。
1 | public interface HeartbeatTarget<I> { |
对HeartbeatTarget的封装,这样Manager对Target的操作是通过对Monitor完成,后续会在其继承类中详细说明。
1 | public interface HeartbeatMonitor<O> { |
HeartbeatManager负责管理心跳机制,比如启动/停止/报告一个HeartbeatTarget。此接口继承HeartbeatTarget。
除了HeartbeatTarget的函数之外,这接口有4个函数:
1 | public interface HeartbeatManager<I, O> extends HeartbeatTarget<I> { |
用户业务逻辑需要继承这个接口以处理心跳结果。其可以看做服务的输出,实现了三个回调函数。
1 | public interface HeartbeatListener<I, O> { |
之前提到Sender和Receiver,下面两个类就对应上述概念。
几个关键问题:
HeartbeatListener.notifyHeartbeatTimeout方法做后续重连操作或者直接断开。下面是一个概要(以RM & TM为例):
HearbeatManagerImpl是receiver的具体实现。它由 心跳 被发起方(就是Receiver,例如TM) 创建,接收 **发起方(就是Sender,例如 JM)**的心跳发送请求。心跳超时 会触发 heartbeatListener.notifyHeartbeatTimeout方法。
注意:被发起方监控线程(Monitor)的开启是在接收到请求心跳(requestHeartbeat被调用后)以后才触发的,属于被动触发。
HearbeatManagerImpl主要维护了
<ResourceID, HeartbeatMonitor<O>> heartbeatTargets;。这是一个KV关联。 key代表要发送心跳组件(例如:TM)的ID,value则是为当前组件创建的触发心跳超时的线程HeartbeatMonitor,两者一一对应。 当一个从所联系的machine发过来的心跳被收到时候,对应的monitor的状态会被更新(重启一个新ScheduledFuture)。当一个monitor发现了一个 heartbeat timed out,它会通知自己的HeartbeatListener。 HearbeatManagerImpl 数据结构如下:
1 | @ThreadSafe |
HearbeatManagerImpl实现的主要函数有:
继承HearbeatManagerImpl,由**心跳管理的一方(例如JM)**创建,实现了run函数(即它可以作为一个单独线程运行),创建后立即开启周期调度线程,每次遍历自己管理的heartbeatTarget,触发heartbeatTarget.requestHeartbeat,要求 Target 返回一个心跳响应。属于主动触发心跳请求。
1 | public class HeartbeatManagerSenderImpl<I, O> extends HeartbeatManagerImpl<I, O> implements Runnable { |
Heartbeat monitor管理心跳目标,它启动一个ScheduledExecutor。
1 | public class HeartbeatMonitorImpl<O> implements HeartbeatMonitor<O>, Runnable { |
建立heartbeat receivers and heartbeat senders,主要是对外提供服务。这里我们可以看到:
1 | public class HeartbeatServices { |
心跳管理服务在Cluster入口创建。因为我们是调试,所以在MiniCluster.start调用。
1 | public void start() throws Exception { |
HeartbeatServices.fromConfiguration会从Configuration中获取配置信息:
这个就是我们解决最开始问题的思路:从配置信息入手,扩大心跳间隔。
1 | public HeartbeatServices(long heartbeatInterval, long heartbeatTimeout) { |
系统中有几个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。
Flink中ResourceManager、JobMaster、TaskExecutor三者之间存在相互检测的心跳机制:
我们之前讲过,HeartbeatManagerSenderImpl属于Sender,HeartbeatManagerImpl属于Receiver。
ResourceManager 级别最高,所以两个HM都是Sender,监控taskManager和jobManager
1 | public abstract class ResourceManager<WorkerType extends ResourceIDRetrievable> |
JobMaster级别中等,一个Sender, 一个Receiver,受到ResourceManager的监控,监控taskManager。
1 | public class JobMaster extends FencedRpcEndpoint<JobMasterId> implements JobMasterGateway, JobMasterService { |
TaskExecutor级别最低,两个Receiver,分别被JM和RM疾控。
1 | public class TaskExecutor extends RpcEndpoint implements TaskExecutorGateway { |
以JobManager和TaskManager为例。JM在启动时会开启周期调度,向已经注册到JM中的TM发起心跳检查,通过RPC调用TM的requestHeartbeat方法,重置对JM超时线程的调用,表示当前JM状态正常。在TM的requestHeartbeat方法被调用后,通过RPC调用JM的receiveHeartbeat,重置对TM超时线程的调用,表示TM状态正常。
TM初始化生成了两个Receiver HM。
1 | public class TaskExecutor extends RpcEndpoint implements TaskExecutorGateway { |
生成HeartbeatManager时,就注册了ResourceManagerHeartbeatListener和JobManagerHeartbeatListener。
此时,两个HeartbeatManagerImpl中已经创建好对应monitor线程,只有在JM或者RM执行requestHeartbeat后,才会触发该线程的执行。
JM生成了一个Sender HM,一个Receiver HM。这里会注册 TaskManagerHeartbeatListener 和 ResourceManagerHeartbeatListener
1 | public class JobMaster extends FencedRpcEndpoint<JobMasterId> implements JobMasterGateway, JobMasterService { |
JobMaster在启动时候,会在startHeartbeatServices函数中生成两个Sender HeartbeatManager。
taskManagerHeartbeatManager :HeartbeatManagerSenderImpl对象,会反复启动一个定时器,定时扫描需要探测的对象并且发送心跳请求。
jobManagerHeartbeatManager :HeartbeatManagerSenderImpl,会反复启动一个定时器,定时扫描需要探测的对象并且发送心跳请求。
1 | taskManagerHeartbeatManager = heartbeatServices.createHeartbeatManagerSender( |
我们以TM与RM交互为例。TaskExecutor启动之后,需要注册到RM和JM中。
流程图如下:
1 | * 1. Run in Task Manager |
下面是具体文字描述。
1 | private final LeaderRetrievalService resourceManagerLeaderRetriever; |
1 | taskManagerHeartbeatManager.monitorTarget(taskExecutorResourceId, new HeartbeatTarget<Void>() { |
当注册完成后,RM中的Sender HM内部结构如下,能看出来多了一个Target:
1 | taskManagerHeartbeatManager = {HeartbeatManagerSenderImpl@8866} |
RM会通过RPC再次回到TaskExecutor,其新执行序列如下:
1 | HeartbeatMonitor<O> heartbeatMonitor = heartbeatMonitorFactory.createHeartbeatMonitor |
当注册完成后,其Receiver HM结构如下:
1 | resourceManagerHeartbeatManager = {HeartbeatManagerImpl@10163} |
其调用基本思路与之前相同,就是TM和JM之间互相注册一个代表对方的monitor:
1 | JobLeaderListenerImpl ----> establishJobManagerConnection |
消息到了JM中,做如下操作。
1 | registerTaskManager ----> taskManagerHeartbeatManager.monitorTarget |
在任务提交之后,我们就进入了正常的心跳监控流程。我们依然用 TM 和 RM进行演示。
我们先给出一个流程图。
1 | * 1. Run in Resouce Manager |
下面是具体文字描述。
心跳机制是由Sender主动发起的。这里就是 ResourceManager 的HeartbeatManagerSenderImpl中定时schedual调用,这里会遍历所有监控的Monitor(Target),逐一在Target上调用requestHeartbeat。
1 | // HeartbeatManagerSenderImpl中的代码 |
具体监控对象 Target 会调用自定义的requestHeartbeat。
1 | HeartbeatManagerSenderImpl |
请注意,每一个Target都是由ResourceManager生成的。ResourceManager之前注册成为Monitor时候就注册了这个HeartbeatTarget。
这个HeartbeatTarget的定义如下,两个函数是:
会调用到ResourceManager定义的函数requestHeartbeat,而requestHeartbeat会通过gateway调用到TM,这就是主动对TM发起了心跳请求。
1 | taskManagerHeartbeatManager.monitorTarget(taskExecutorResourceId, new HeartbeatTarget<Void>() { |
通过taskExecutorGateway。心跳程序执行就通过RPC从RM跳跃到了TM。
taskExecutorGateway.heartbeatFromResourceManager 的意义就是:通过RPC调用回到TaskExecutor。这个是在TaskExecutorGateway就定义好的。
1 | // TaskExecutor RPC gateway interface. |
TaskExecutor实现了TaskExecutorGateway,所以具体在TaskExecutor内部实现了接口函数。
1 | @Override |
TM中,resourceManagerHeartbeatManager 定义如下。
1 | /** The heartbeat manager for resource manager in the task manager. */ |
所以下面就是执行TM中的Receiver HM。在这个过程中有两个处理步骤:
具体是调用 requestHeartbeat@HeartbeatManager。在其中会
1 | @Override |
最后会通过resourceManagerGateway.heartbeatFromTaskManager 调用到 ResourcManager。
JobMaster在接收到rpc请求后调用其heartbeatFromTaskManager方法,会调用taskManagerHeartbeatManager的receiveHeartbeat方法,在这个过程中同样有两个处理步骤:
至此一次完成心跳过程已经完成,会根据heartbeatInterval执行下一次心跳。
首先,在HeartbeatMonitorImpl中,如果超时,会调用Listener。
1 | public void run() { |
这就来到了ResourceManagerHeartbeatListener,会尝试再次连接RM。
1 | private class ResourceManagerHeartbeatListener implements HeartbeatListener<Void, TaskExecutorHeartbeatPayload> { |
RM就直接简单粗暴,关闭连接。
1 | private class TaskManagerHeartbeatListener implements HeartbeatListener<TaskExecutorHeartbeatPayload, Void> { |
心跳机制我们讲解完了,但是我们最初提到的异常应该如何解决呢?在程序最开始生成环境变量时候,通过设置环境变量的配置即可搞定:
1 | Configuration conf = new Configuration(); |
在Flink SQL中, 元数据的管理分为三层: catalog-> database-> table,
我们知道Flink SQL是依托calcite框架来进行SQL执行树生产,校验,优化等等, 所以本文讲介绍FlinkSQL是如何来结合Calcite来进行元数据管理的.
1 | public interface Schema { |
如接口所示, Schema接口,可以通过table名来获得一张表, 可以通过schema名来获得一个子schema.
1 | public interface Table { |
看Table的接口, 主要就是返回table的RelDataType.
接下来,我们来看下Flink是如何实现这些接口的:
1 | public class CatalogManagerCalciteSchema extends FlinkSchema { |
很容易发现,CatalogSchema返回DatabaseSchema, DatabaseSchema返回Table,
这样就容易理解,Flink的三层结构是怎样的了. 同时, 具体的元数据实际上都是在catalogManager中。
DatabaseSchema中返回的Table类型为CatalogSchemaTable,我们来看下具体的结结构是怎样的,
上文中也提到了,Table接口主为getRowType函数, 用于返回某个table的type信息。
TableSchema是Flink内部用于保存各个字段的类型信息的类, 通过相关的转化函数,转换为calcite的type类型.
1 | public class CatalogSchemaTable extends AbstractTable implements TemporalTable { |
CatalogBaseTable接口定义如下, Flink的Table的参数(schema参数,connector参数)都可以最终表示为一个map.
1 | public interface CatalogBaseTable { |
上面都是的相关接口都是Flink用于适配calcite框架元数据的相关实现。
那么这些类具体是在哪里调用的? 已经什么时候会被调用到?
calcite中的schema,主要是在validate过程中, 获得对应table的字段信息, 对应的function的返回值信息,
确保SQL的字段名, 字段类型是正确的.
类的依赖关系为:
validator —> schemaReader —> schema
FlinkPlannerImpl.scala中
1 | private def createSqlValidator(catalogReader: CatalogReader) = { |
PlanningConfigurationBuilder.java
1 | private CatalogReader createCatalogReader( |
综上所诉, 我们就知道了Flink是如何来利用calcite的schema来管理Flink的table信息的.
Join的几种形式,与实现原理
从 Flink 1.6 版本开始,社区引入了状态 TTL(Time-To-Live)特性。在通过Flink SQL 实现流处理时,开发者可以为作业 SQL 设置TTL,实现过期状态的自动清理,从而防止作业状态无限膨胀
1 | SELECT |
https://ci.apache.org/projects/flink/flink-docs-stable/dev/table/streaming/query_configuration.html
Flink Table API 和SQL接口提供参数来调整连续查询的准确性和资源消耗。参数通过 QueryConfig 对象指定。QueryConfig 可以从 TableEnvironment 获得。
1 | val env = StreamExecutionEnvironment.getExecutionEnvironment |
空闲状态保持时间(Idle State Retention Time)参数定义一个键的状态在一次更新之后保存多久后删除。
通过删除键的状态,连续查询会完全忘记它之前已经看过这个键。如果删除的键再次出现,则被视为具有相应键的第一个记录。对于前面的查询示例,这意味着 sessionId 的计数从0开始。
配置空闲状态保存时间有两个参数:
minimum idle state retention time,定义非活动键的状态在删除前至少保持多少时间。
maximum idle state retention time,定义非活动键的状态在删除前最多保持多少时间。
对于前面的查询示例:
1 | val qConfig: StreamQueryConfig = ??? |
清理状态需要额外的记录,对于 minTime 和 maxTime 较大差异的情况成本更低,因此 minTime 和 maxTime 直接必须至少相差5分钟。
1 | SELECT |
1 | bin/sql-client.sh embedded |
从”UDF不应有状态” 切入来剖析Flink SQL代码生成
核心类: org.apache.flink.table.runtime.CRowProcessRunner
打印出codegen语句:
1 | log4j.logger.org.apache.flink.runtime.CRowProcessRunner=DEBUG, file |
1 | alter table |
1 | 2021-10-27 15:06:43,520 DEBUG [Task-46f8730428df9ecd6d7318a02bdc405e(0/1)] org.apache.flink.table.runtime.conversion.CRowToJavaTupleMapRunner - Compiling MapFunction: DataStreamSinkConversion$14 |
1 | CREATE TABLE mysql_orders ( |
1 | insert into mysql_orders select concat(ei,'#',sei) as pk,sei,ui,etime from orders |
calcite解析后:
1 | org.apache.calcite.sql2rel,DEBUG,Plan after converting SqlNode to RelNode |
优化:
1 | org.apache.flink.table.planner.plan.optimize.program.FlinkGroupProgram,DEBUG,optimize convert table references before rewriting sub-queries to semi-join cost 31 ms. |
1 | org.apache.calcite.plan.RelOptPlanner,DEBUG,Cheapest plan: |
实际的执行计划分为三层:
Flink使用Calcite实现了SQL的解析、转换、执行计划优化和转换,那么,Flink SQL是如何执行的呢?
MySQL实体表
1 | CREATE TABLE IF NOT EXISTS t_rt_agg_result( |
数据源
1 | CREATE TABLE orders ( |
假设选用Blink planner运行两个SQL:
1 | CREATE TABLE agg_result ( |
1 | INSERT INTO MyUserTable |
执行下面MySQL
1 | tEnv.sqlUpdate(ddlSql); |
1 | == Abstract Syntax Tree == |
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
https://matt33.com/2019/10/20/paper-flink-snapshot/
默认情况下,Checkpoint机制是关闭的,需要调用env.enableCheckpointing(n)来开启,每隔n毫秒进行一次Checkpoint。Checkpoint是一种负载较重的任务,如果状态比较大,同时n值又比较小,那可能一次Checkpoint还没完成,下次Checkpoint已经被触发,占用太多本该用于正常数据处理的资源。增大n值意味着一个作业的Checkpoint次数更少,整个作业用于进行Checkpoint的资源更小,可以将更多的资源用于正常的流数据处理。同时,更大的n值意味着重启后,整个作业需要从更长的Offset开始重新处理数据。
发表于 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都采用反压机制解决这个问题,不过他们的实现各自不同。

不同的组件可以不同的速度执行(并且每个组件中的处理速度随时间改变)。 例如,考虑一个工作流程,或由于数据倾斜或任务调度而导致数据被处理十分缓慢。
在这种情况下,如果上游阶段不减速,将导致缓冲区建立长队列(队列占用内存、硬盘空间,节点负载加重),或导致系统丢弃元组。
如果元组在中途丢弃,那么效率可能会有损失,因为已经为这些元组产生的计算被浪费了。
并且在一些流处理系统中比如Strom,会将这些丢失的元组重新发送,这样会导致数据的一致性问题(at least once语义),并且还会导致某些Operator状态叠加。
进而整个程序输出结果不准确。第二由于系统接收数据的速率是随着时间改变的,短时负载高峰导致系统接收数据的速率远高于它处理数据的速率的情况,也会导致Tuple在中途丢失。
所以实时流处理系统必须能够解决发送速率远大于系统能处理速率这个问题,大多数实时流处理系统采用反压(BackPressure)机制解决这个问题。
下面我们就来介绍一下不同的实时流处理系统采用的反压机制:
对于开启了acker机制的storm程序,可以通过设置conf.setMaxSpoutPending参数来实现反压效果,如果下游组件(bolt)处理速度跟不上导致spout发送的tuple没有及时确认的数超过了参数设定的值,spout会停止发送数据,这种方式的缺点是很难调优conf.setMaxSpoutPending参数的设置以达到最好的反压效果,设小了会导致吞吐上不去,设大了会导致worker OOM;有震荡,数据流会处于一个颠簸状态,效果不如逐级反压;另外对于关闭acker机制的程序无效;
新的storm自动反压机制(Automatic Back Pressure)通过监控bolt中的接收队列的情况,当超过高水位值时专门的线程会将反压信息写到 Zookeeper ,Zookeeper上的watch会通知该拓扑的所有Worker都进入反压状态,最后Spout降低tuple发送的速度。

每个Executor都有一个接受队列和发送队列用来接收Tuple和发送Spout或者Bolt生成的Tuple元组。每个Worker进程都有一个单的的接收线程监听接收端口。
它从每个网络上进来的消息发送到Executor的接收队列中。Executor接收队列存放Worker或者Worker内部其他Executor发过来的消息。
Executor工作线程从接收队列中拿出数据,然后调用execute方法,发送Tuple到Executor的发送队列。
Executor的发送线程从发送队列中获取消息,按照消息目的地址选择发送到Worker的传输队列中或者其他Executor的接收队列中。
最后Worker的发送线程从传输队列中读取消息,然后将Tuple元组发送到网络中。
high watermark的阈值后,因此它会发送通知消息到背压线程。/Backpressure/topo1/wk1下JStorm做了两级的反压,第一级和Jstorm类似,通过执行队列来监测,但是不会通过ZK来协调,而是通过Topology Master来协调。
在队列中会标记high water mark和low water mark,当执行队列超过high water mark时,就认为bolt来不及处理,则向TM发一条控制消息,上游开始减慢发送速率,直到下游低于low water mark时解除反压。
此外,在Netty层也做了一级反压,由于每个Worker Task都有自己的发送和接收的缓冲区,可以对缓冲区设定限额、控制大小,如果spout数据量特别大,缓冲区填满会导致下游bolt的接收缓冲区填满,造成了反压。

限流机制:jstorm的限流机制, 当下游bolt发生阻塞时, 并且阻塞task的比例超过某个比例时(现在默认设置为0.1),触发反压
限流方式:计算阻塞Task的地方执行线程执行时间,Spout每发送一个tuple等待相应时间,然后讲这个时间发送给Spout, 于是, spout每发送一个tuple,就会等待这个执行时间。
Task阻塞判断方式:在jstorm 连续4次采样周期中采样,队列情况,当队列超过80%(可以设置)时,即可认为该task处在阻塞状态。
默认情况下,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),通过动态控制数据接收速率来适配集群数据处理能力。
Spark Streaming Backpressure: 根据JobScheduler反馈作业的执行信息来动态调整Receiver数据接收率。通过属性“spark.streaming.backpressure.enabled”来控制是否启用backpressure机制,默认值false,即不启用。
1 | sparkConf.set("spark.streaming.backpressure.enabled",”true”) |
SparkStreaming 架构图如下所示:

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

当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 | 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条记录。

当下游处理速度跟不上上游发送速度时,一旦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 是如何在 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 的包装类。

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

不要忘了:记录能被 Flink 处理的前提是,必须有空闲可用的 Buffer。
结合上面两张图看:Task 1 在输出端有一个相关联的 LocalBufferPool(称缓冲池1),Task 2 在输入端也有一个相关联的 LocalBufferPool(称缓冲池2)。如果缓冲池1中有空闲可用的 buffer 来序列化记录 “A”,我们就序列化并发送该 buffer。
这里我们需要注意两个场景:
这种固定大小缓冲池就像阻塞队列一样,保证了 Flink 有一套健壮的反压机制,使得 Task 生产数据的速度不会快于消费的速度。我们上面描述的这个方案可以从两个 Task 之间的数据传输自然地扩展到更复杂的 pipeline 中,保证反压机制可以扩散到整个 pipeline。
另外,官方博客中为了展示反压的效果,给出了一个简单的实验。下面这张图显示了:随着时间的改变,生产者(黄色线)和消费者(绿色线)每5秒的平均吞吐与最大吞吐(在单一JVM中每秒达到8百万条记录)的百分比。我们通过衡量task每5秒钟处理的记录数来衡量平均吞吐。该实验运行在单 JVM 中,不过使用了完整的 Flink 功能栈。

首先,我们运行生产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 | java.lang.Object.wait(Native Method) |
那么事情就简单了。通过不断地采样每个 task 的 stack trace 就可以实现反压监控。

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应用更稳定、有效