flink 1.5.6 源码分析——flink job 执行流程详解(二)
目录
一.wordcount代码逻辑分析
1.1 StreamExecutionEnvironment 运行环境
1.2 DataStreamSource 流数据源
1.3 流处理过程和Sink输出---flatMap,keyBy,sum,print
flatMap的逻辑
keyBy的逻辑
Sum的逻辑
Print的逻辑
小结
1.4 env.execute启动job
二. flink 四层图结构
2.1 WordCount的StreamGraph
本文将基于flink-1.5.6, WordCount类作为案例程序进行研究与分析。
wordcount的代码逻辑如下:
-
- 初始化StreamExecutionEnvironment
-
- 初始化DataStreamSource(并将其实例命名为Source)
-
- 源模块代表Flink的数据输入流。
-
- 为源模块增加一系列操作。
-
- 每一个操作都会生成一个新的DataStream对象,并对该对象执行相应的操作。
-
- 最终的DataStream引入汇出模块。
-
- 汇出模块负责整合并输出整个流程的结果。
-
- 通过调用env.execute()方法开始执行Flink作业。
一.wordcount代码逻辑分析
1.1 StreamExecutionEnvironment 运行环境
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
涉及Flink系统的上下文包括其默认配置设置以及当前运行时的状态。这些因素例如根据CPU核心数量设定默认并行度,并采用'execute'操作作为启动过程。而运行环境的作用则是记录和管理整个数据流处理流程中的各个操作节点及其相关信息。源节点与操作节点之间的关系完整连接到目标节点,并且完整记录着各阶段返回的数据类型以及系统的整体交互关联性。
1.2 DataStreamSource 流数据源
text = env.fromElements(WordCountData.WORDS);
DataStreamSource 是Flink流计算的起源,并与其它类型的流具有类似的特征。其源码主要实现了以下功能:
- 1.通过TypeExtractor获取DataStreamSource产生的数据类型结果为String。
- 2.建立源函数,并利用该源函数初始化相应的源操作。
- 3.基于源操作构建并配置相应的源转换。
- 4.借助源转换工具来设置DataStreamSource的初始状态。
部分名词解释 | 名词| 解释 |
| --- | --- |
|---|---|
| Function | 流数据的基本处理单元,数据处理逻辑封装在function中,例如:SourceFunction封装Source中用来获取数据的代码,FlatMapFunction封装了处理数据FlatMap逻辑的代码 |
| Operator | 流操作,将一个Function封装为不同类型的Operator,例如:单输出Operator OneInputStreamOperator ,双输出Operator TwoInputStreamOperator等,是function的不同输出类型的处理逻辑封装。 |
| Transformation | 转换,将Operator封装为Transformation,用于形容一组类型的输入,经过Transformation以后,出来一组不同类型的输出。是Operator的更高级封装,也可以分为单输出,双输出,Source,sink等多种Transformation |
| DataStream | DataStream代表一个数据流,里面封装transformation,表示流的处理逻辑。还有若干基于该流往下游的方法。例如map,flatmap。 |
| SourceFunction | Function的一种封装,负责Source处理逻辑的Function,例如FromElementsFunction,看了一下,SourceFunction的实现,flink源码中就有22种 |
| SourceOperator | Operator的一种封装,封装SourceFunction的Operator,可以封装不同的SourceFunction |
| SourceTransformation | Transformation的一种实现,负责封装不同类别的SourceOperator |
| DataStreamSource | 数据流的一种实现,表示数据流计算的起点,里面封装了SourceTransformation |
1.3 流处理过程和Sink输出---flatMap,keyBy,sum,print
从Source到Sink中间的处理阶段,这个阶段的方法分两种:
- 1.处理数据的方法
- 2.不处理数据的方法
值得注意的是此乃博主个人的理解,并非官方的说法。为什么要这样分类?由于flatMap这类操作通常涉及对单行数据进行展宽处理。它会将单行数据转换为多行结构的数据。而像keyBy、equalTo、partition以及rebalance等方法,则不具备这种修改功能。它们主要承担按字段匹配或分配的任务。因此,在计算过程中这些方法不会被纳入Operator
flatMap的逻辑
代码这里就不贴了,一路点进去是很清晰的。大体逻辑如下:
-
- 通过FlatMapFunction的类型推断获取ResultType变量。
-
- 将该Function转换为Operator,并将其与ResultType一并传递给transform方法。
transform方法是公共方法,主要做下面的几件事情:
- 实现将当前流及其operator instance基于其transformation生成flatmap的对象。
- 基于flatmap transformation object生产返回的结果数据集。
- 将flatmap transformation objec加入 StreamExecutionEnvironment 的 transformations 数组中。
在上述流程中,首先将当前流this.transformation以及Operator用于构建flatmap所需的Transformation对象。每个Transformation对象都会携带一个input字段(即current flow this.transformation),它表示当前转换操作所依赖的上游转换操作。这样一来,在获取所有Transformation对象之后便能够明确它们之间的上下级关系。随后将Transformation组件封装至DataStream中,并将其加入到StreamExecutionEnvironment中。这样一来每个DataStream都将包含自身所执行的处理逻辑(transformations→Operator→functions, 装饰模式)而StreamExecutionEnvironment则会存储所有的Transformations组件,并基于每个Transformations对象的input字段来揭示它们之间的上下级关联关系
keyBy的逻辑
keyBy的逻辑如下:
以下是按照要求对原文进行同义改写的文本
不同于Transform方法,在生成keyedStream后会返回结果,并未将相应的Transformation添加至StreamExecutionEnvironment;这一设计背后的原因值得探讨;
Sum的逻辑
我看了一下代码后发现,sum的行为与flatMap大致一致,遵循函数→操作→转换→DataStream的流程,并且将转换操作放入了StreamExecutionEnvironment中进行处理。值得注意的是,在此过程中,转换操作的具体输入来源于keyBy的PartitionTransformation组件。
审视StreamExecutionEnvironment的变换类型是否存在遗漏的设置
Print的逻辑
处理逻辑与flatMap以及sum本质上是相同的。其中使用的是一项关键组件是Flink内置的DataStreamSink,并且遵循功能→操作→转换→DataStreamSink这一流程。随后将这些转换整合进StreamExecutionEnvironment体系中。
整理一下:
在WordCount系统中所采用的流处理流程包括Source节点、FlatMap操作、KeyBy步骤以及Sum节点,并最终达到Sink节点。每个操作都会生成DataStream对象,并伴随相应的Transformation。
StreamExecutionEnvironment中保存了3个Transformation,分别是
- 此Flatmap变换器id为2,其输入为DataStreamSource变换器id1。(注意:源变换器无输入端)
- 此处Sum变换器id4接收来自Keyby变换器id3的输入。(观察发现:该Keyby变换器的上游是Flatmap转换器id2)
- 此处Sink转换er_id5连接至Sum转换er_id4。(观察发现:该Sum转换er的上游与前一描述一致)
从Source到Sink中间的处理阶段,这个阶段的方法分两种:
- 1.处理数据的方法
- 2.不处理数据的方法
基于前面所述的内容,在分析为何不打算把keyby对应的transformation整合进StreamExecutionEnvironment系统中时,默认情况下认为存在某些限制因素:作者在此基础上提出了一些初步假设。
未对数据进行分发的方法类似于keyBy和Source的方式;由于下一级必定属于Operator类型,则只需在Operator处实施Transformation,并将 upstream的数据传递给该Transformation的输入端即可。
2.添加最后一个Sink,可以得到Source->Sink的一条链路,
在某个DAG网络中,在源节点与操作节点之间以及操作节点与汇点之间形成了单向连接关系。当该网络中的某些节点之间存在多条路径时,则无法仅通过最终节点来推断所有可能的路径。这表明为了完整地了解整个系统的运行机制和数据流路径关系,在这种情况下必须记录所有参与操作的节点
并非所有Transformation都会转换为runtime层的实际操作。有些仅作为逻辑概念存在,例如union、split-select以及partition等。如图所示的转换树,在运行时会优化为下方的操作图

小结
我们开发了一个Flink程序,在代码封装结构中遵循函数→操作符→转换器→DataStream(Source, Sink, Stream三种)的层级架构。
每个transformation都会记录上一个流的transformation作为input
Operator和Sink的相关变换会被记录到StreamExecutionEnvironment中,并非必须添加所有的Transformations以保持dag的完整性;反而无需添加所有的Transformations,并带来了一定的优化效果。
1.4 env.execute启动job
env.execute启动job的过程详见我的上一篇博客Flink-1.5.6 源码分析-----main方法转换为dag的流程,具体内容请参见。
二. flink 四层图结构
先贴一张网上流传很多的图

转换过程说明:
- 1.StreamExecutionEnvironment存储的transformation依次演变为StreamGraph、JobGraph、ExecutionGraph,并最终形成物理执行图。
- 在客户端完成后被提交至JobManager。
- 当JobManager的主节点(即JobMaster)处理后将每个Node依次转换为对应的ExecutionGraph并发送至相应的taskManager以生成最终的实际物理执行图。
2.1 WordCount的StreamGraph
源代码位于StreamGraphGenerator.generateInternal函数内部随后通过transform->transformOneInputTransform路径继续进行操作这一过程涉及将Transformation转换为StreamGraph的技术为此我们决定首先分析其组成部分以期更好地理解其工作原理
public class StreamGraph extends StreamingPlan {
//记录 jobName,StreamExecutionEnvironment,ExecutionConfig,CheckpointConfig等配置
// 记录StreamNode节点,使用id作为key
private Map<Integer, StreamNode> streamNodes;
//记录Source
private Set<Integer> sources;
//记录Sink
private Set<Integer> sinks;
//记录Select 节点,该节点是虚拟的,需要处理,但是不放入Streamgraph中。后面会被优化掉
private Map<Integer, Tuple2<Integer, List<String>>> virtualSelectNodes;
//记录 side output 节点,该节点是虚拟的,需要处理,但是不放入Streamgraph中。后面会被优化掉
private Map<Integer, Tuple2<Integer, OutputTag>> virtualSideOutputNodes;
//记录 partition 节点,该节点是虚拟的,需要处理,但是不放入Streamgraph中。后面会被优化掉
//wordcount 这里为 6 = (2, HASH)
private Map<Integer, Tuple2<Integer, StreamPartitioner<?>>> virtualPartitionNodes;
}
在StreamGraph架构中,其中最关键的是streamNodes(对应变换的节点)、data sources(数据源端)以及output sinks(输出端)。此外还有一些虚拟节点如side outputs、selectors和partitions用于分发等目的;这些都会被生成为虚拟节点,并最终通过优化将其移除。
StreamGraph未对边进行记录而仅记录了节点信息因此可以推测边的存在应与节点相关联建议进一步考察StreamNodes的构成情况
public class StreamNode implements Serializable {
//id
private final int id;
//并行度
private Integer parallelism = null;
//处理逻辑Operator,保存在Node中
private transient StreamOperator<?> operator;
//入度边
private List<StreamEdge> inEdges = new ArrayList<StreamEdge>();
//出度边
private List<StreamEdge> outEdges = new ArrayList<StreamEdge>();
}
StreamNodes的关键属性包括statePartitioner、outputSelectors以及入边(inEdges)与出边(outEdges),这些元素共同构成了其核心功能结构
查看StreamEdge的属性
public class StreamEdge implements Serializable {
private static final long serialVersionUID = 1L;
// 边 id
private final String edgeId;
//来源端
private final StreamNode sourceVertex;
//目标端
private final StreamNode targetVertex;
/** * The type number of the input for co-tasks.
*/
private final int typeNumber;
/** * SelectTransformation当中用户注入的名字集合
*/
private final List<String> selectedNames;
/** * SideOutputTransformation当中用户指定的OutputTag
*/
private final OutputTag outputTag;
/** * The {@link StreamPartitioner} on this {@link StreamEdge}.
* 默认ForwardPartitioner,可由PartitionTransformation指定
*/
private StreamPartitioner<?> outputPartitioner;
}
StreamEdge的关键属性包括源顶点标识符(sourceVertex)、目标顶点标识符(targetVertex)、选定名称集合(selectedNames)、输出标签(outputTag)以及输出分片器(outputPartitioner),其中输出分片器通常默认设置为ForwardPartitioner,并可通过PartitionTransformation进行配置。
关于分流器 StreamPartitioner
| StreamPartitioner class | toString | 功能 | 场景 |
|---|---|---|---|
| 分区器,将流的数据从一个Channel输出到其他的一个或多个Channel, | 流数据从上游到下游, |
|---|
| int类型,标识下游的序号,例如,下游并行度为4 则每个子任务对应的Channel为 0,1,2,3 |
|---|
||
以round-robin的形式将元素分区到下游subtask的子集中 |
|---|
||
重平衡分区器,用于实现类似于round-robin这样的轮转模式的分区器。通过累加、取模的形式来实现对输出channel的切换。 |
上下游并行度不相同时 |
|---|
||
| 全局分区器,无论有多少Channel,直接输出到0 |
|---|
||
| 按key分组,大体等于key的hash对并行度取余 channel = keyGroupId * parallelism / maxParallelism 其中 keyGroupId 是根据key计算,按下面的公式 MathUtils.murmurHash(key.hashCode()) % maxParallelism | 调用keyBy() |
|---|
||
混洗分区器,该分区器会在所有output channel中选择一个随机的进行输出 |
调用shuffle |
|---|
||
| 该分区器将记录转发给在本地运行的下游的,一对一或多对一 | 默认Forward,没有指定partiiton时,例如map->reduce |
|---|
||
| 用于支持自定义实现, | 自定义 |
|---|
||
|广播分区器,用于将该记录广播给下游的所有的subtask|broadcast时|
了解流分区器的相关信息可以通过查阅源码以及Apache Flink流分区器剖析这篇文章来获取详细解析
flink 包含了一个名为 StreamGraph 的可视化展示工具,
可访问此处:在这里。
我们可以通过以下步骤获取并分析程序的执行计划:
首先,在控制台中运行 System.out.println(env.getExecutionPlan());
然后将输出结果复制到上述网站中进行查看与处理,
如图所示即为此操作流程。

可以看到,在源程序中已经实现了将程序转换为4个operator的功能。此外,在这些operator之间的连接线上还揭示了flink添加的一些逻辑流程。
