要懂得一个体系,一般都是大年夜架构开端。我们关怀的问题是:体系安排成功后各个节点都启动了哪些办事,各个办事之间又是怎么交互和调和的。下方是 Flink 集群启动后架构图。

当 Flink 集群启动后,起首会启动一个 JobManger 和一个或多个的 TaskManager。由 Client 提交义务给 JobManager,JobManager 再调剂义务到各个 TaskManager 去履行,然后 TaskManager 将心跳和统计信息报告请示给 JobManager。TaskManager 之间以流的情势进行数据的传输。上述三者均为自力的 JVM 过程。
- Client 为提交 Job 的客户端,可所以运行在任何机械上(与 JobManager 情况连通即可)。提交 Job 后,Client 可以停止过程(Streaming的义务),也可以不停止并等待结不雅返回。
- JobManager 重要负责调剂 Job 并调和 Task 做 checkpoint,职责上很像 Storm 的 Nimbus。大年夜 Client 处接收到 Job 和 JAR 包等资本后,会生成优化后的履行筹划,并以 Task 的单位调剂到各个 TaskManager 去履行。
- TaskManager 在启动的时刻就设置好了槽位数(Slot),每个 slot 能启动一个 Task,Task 为线程。大年夜 JobManager 处接收须要安排的 Task,安排启动后,与本身的上游建立 Netty 连接,吸法术据并处理。
可以看到 Flink 的义务调剂是多线程模型,并且不合Job/Task混淆在一个 TaskManager 过程中。固然这种方法可以有效进步 CPU 应用率,然则小我不太爱好这种设计,因为不仅缺乏资本隔离机制,同时也不便利调试。类似 Storm 的过程模型,一个JVM 中只跑该 Job 的 Tasks 实际应用中更为合理。
Job 例子
本文所示例子为 flink-1.0.x 版本
我们应用 Flink 自带的 examples 包中的 SocketTextStreamWordCount ,这是一个大年夜 socket 流中统计单词出现次数的例子。
-
起首,应用 netcat 启动本地办事器:
$ nc -l 9000
-
然后提交 Flink 法度榜样
起首我们看到,JobGraph 之上除了 StreamGraph 还有 OptimizedPlan。OptimizedPlan 是由 Batch API 转换而来的。StreamGraph 是由 Stream API 转换而来的。为什么 API 不直接转换成 JobGraph?因为,Batch 和 Stream 的图构造和优化办法有很大年夜的差别,比如 Batch 有很多履行前的预分析用来竽暌古化图的履行,而这种优化并不普适于 Stream,所以经由过程 OptimizedPlan 来做 Batch 的优化会更便利和清楚,也不会影响 Stream。JobGraph 的义务就是同一 Batch 和 Stream 的图,用来描述清跋扈一个拓扑图的构造,并且做了 chaining 的优化,chaining 是普适于 Batch 和 Stream 的,所以在这一层做掉落。ExecutionGraph 的义务是便利调剂和各个 tasks 状况的监控和跟踪,所以 ExecutionGraph 是并行化的 JobGraph。而“物理履行图”就是最终分布式在各个机械上运行着的tasks了。所以可以看到,这种解耦方法极大年夜处所便了我们在各个层所做的工作,各个层之间是互相隔离的。
$ bin/flink run examples/streaming/SocketTextStreamWordCount.jar \ --hostname 10.218.130.9 \ --port 9000
SocketTextStreamWordCount 的具体代码如下:
public static void main(String[] args) throws Exception{ // 检查输入 final ParameterTool params = ParameterTool.fromArgs(args); ... // set up the execution environment final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // get input data DataStream<String> text = env.socketTextStream(params.get("hostname"), params.getInt("port"), '\n', 0); DataStream<Tuple2<String, Integer>> counts = // split up the lines in pairs (2-tuples) containing: (word,1) text.flatMap(new Tokenizer()) // group by the tuple field "0" and sum up tuple field "1" .keyBy(0) .sum(1); counts.print(); // execute program env.execute("WordCount from SocketTextStream Example");}我们将最后一行代码 env.execute 调换成 System.out.println(env.getExecutionPlan()); 并在本地运行该代码(并发度设为2),可以获得该拓扑的逻辑履行筹划图的 JSON 串,将该 JSON 串粘贴到 http://flink.apache.org/visualizer/ 中,能可视化该履行图。

但这并不是最终在 Flink 中运行的履行图,只是一个表示拓扑节点关系的筹划图,在 Flink 中对应了 SteramGraph。别的,提交拓扑后(并发度设为2)还能在 UI 中看到另一张履行筹划图,如下所示,该图对应了 Flink 中的 JobGraph。
在netcat端输入单词并监控 taskmanager 的输出可以看到单词统计的结不雅。

Graph
看起来竽暌剐点乱,怎么竽暌剐这么多不一样的图。实际上,还有更多的图。Flink 中的履行图可以分成四层:StreamGraph -> JobGraph -> ExecutionGraph -> 物理履行图。
- StreamGraph: 是根据用户经由过程 Stream API 编写的代码生成的最初的图。用来表示法度榜样的拓扑构造。
- JobGraph: StreamGraph经由优化后生成了 JobGraph,提交给 JobManager 的数据构造。重要的优化为,将多个相符前提的节点 chain 在一路作为一个节点,如许可以削减数据在节点之间流动所须要的序列化/反序列化/传输消费。
推荐阅读
时至今日,微办事相干的话题不堪列举,上百次的会议,在线评论辩论以及相干文┞仿。你可以假设大年夜家已经熟悉到其长处以及与之俱来的风险。然而,有很多组织没有事先预备就迈入这个潮流了。天然,这也就导致了在架>>>详细阅读
本文标题:Flink 原理与实现:架构和拓扑概览
地址:http://www.17bianji.com/lsqh/36085.html
1/2 1

网友点评
精彩导读
科技快报
品牌展示