前言
每次我们编写完flink作业,跑任务的时候都会在flink-ui上展示一个作业的DAG图,那么这个图是如何形成的呢?本文就和你一起来揭开flink执行图生成的神秘面纱~
##总览
在flink中的执行图可以分为4层StreamGraph -> JobGraph -> ExecutionGraph -> 物理执行图。
- StreamGraph:是根据用户通过 Stream API 编写的代码生成的最初的图。用来表示程序的拓扑结构。
- JobGraph:StreamGraph经过优化后生成了 JobGraph,提交给 JobManager 的数据结构。主要的优化为,将多个符合条件的节点 chain 在一起作为一个节点,这样可以减少数据在节点之间流动所需要的序列化/反序列化/传输消耗。
- ExecutionGraph:JobManager 根据 JobGraph 生成ExecutionGraph。ExecutionGraph是JobGraph的并行化版本,是调度层最核心的数据结构。
- 物理执行图:JobManager 根据 ExecutionGraph 对 Job 进行调度后,在各个TaskManager 上部署 Task 后形成的“图”,并不是一个具体的数据结构。
今天我们就来看下streamGraph的生成
StreamGraph的生成
组件
1 | StreamGraph:根据用户通过 Stream API 编写的代码生成的最初的图。 |
flink任务从定义一个运行环境开始streamExecutionEnvironment,流计算任务起始于addSource,我们来看这个函数,addSource之后生成了一个DataStream. DataStream的构造函数参数接收一个StreamTransformation类型的对象,这个对象反映了流之间的转换操作。但是这个transformation和operation不是一一对应的。一些分区操作:union,split/select,partition只是逻辑概念,并不会在最后的dag图上显示出来。
在生成datastream之后,经历DataStream.java中定义的一些api算子,完成业务逻辑的定义,在这之中可能包含以下的转化:
假设一个场景:
1 | addsource -> map -> filter -> connect -> flatmap |
- addSource创建生成一个
SingleOutputStreamOperator本质上是一个带有transformation=”SourceTransformation”的datastream - map创建生成一个
OneInputTransformation并调用getExecutionEnvironment().addOperator(resultTransform)将其添加入env的List<StreamTransformation<?>>中 - filter 通过 map相同操作
- connect 直接返回一个
ConnectedStreams不是Datastream的子类 - flatmap 生成一个
TwoInputTransformation将其添加入env的List<StreamTransformation<?>>中,并返回一个SingleOutputStreamOperator,并且返回的Datastream中包含的是当前这个transformation - keyby 生成一个keyedStream,这里直接生成一个
PartitionTransformation替代了父类DataStream中的transformation - window 生成windowStream
- apply调用将生成一个
OneInputTransformation,增加至List<StreamTransformation<?>> - addSink 获取
SinkTransformation
其中每一次创建OneInputTransformation都是基于Datastream的当前的transformation来创建的,也就是说keyby之后的PartitionTransformation信息也加入了.
1 | new OneInputTransformation<>( |
好了到这里已经获取了各个流程的streamtransformation,最后调用execute方法,截取了流式环境下的实现:
1 | public JobExecutionResult execute(String jobName) throws Exception { |
其实主要调用的就是
1 | StreamGraphGenerator.generate(this, transformations); |
每一个OneInputTransformation都会记录他的上游的input的transformation,在StreamGraphGenerator.generate主要针对不同的transformation进行不同的转化
1 | private Collection<Integer> transform(StreamTransformation<?> transform) { |
可以看到他里面的方法都是递归调用transform(input)方法,然后通过alreadyTransformed数据结构,避免重复计算,所以我们最终看的时候最先是从source处进行的,也就是从上游到下游进行转化
- 如果已经在
alreadyTransformed数据结构中那么就直接返回transformation的id - 分别有addSource,addOperator,addSink,addCoOperator,addEdge的不同操作来生成streamGraph中的不同节点
- addEdge建立每一个transformation和他所有上游输入节点的连线
在streamgraph中还建立了几个虚拟的节点,这几个节点主要针对的是partition,split/select,sideoutput的操作。
1 | private Map<Integer, Tuple2<Integer, List<String>>> virtualSelectNodes; |
在进行这些操作时,会添加一个唯一的虚拟节点
1 | //记录了上游某个transformId到下游的partition方式 |
经过一些列的addNode以及addEdge之后,streamGraph已经生成。关于其他几个graph的生成请听下回的分解