ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

flink基础知识,系统运行时架构,部署模式详解

flink基础知识,系统运行时架构,部署模式详解 一、flink基础核心知识1、什么是Flink2014年Flink作为主攻流计算的大数据引擎开始在开源大数据行业内崭露头角。区别于Storm、Spark Streaming 以及其他流式计算引擎的是它不仅是一个高吞吐、低延迟的 计算引擎同时还提供很多高级功能。比如它提供有状态的计算支持状态管理支持强一致性的数据语义以及支持Event Time,WaterMark 对消息乱序的处理等。2015年是流计算百花齐放的时代各个流计算框架层出不穷。Storm, JStorm, Heron,Flink, Spark StreamingGoogle Dataflow(后来的Beam)等等。其中Flink的一致性语义和最接近Dataflow模型的开源实现使其成为流计算框架中最耀眼的一颗。也许这也是阿里看中Flink的原因并决心投入重金去研究基于Flink的Blink框架。2、Flink的优点1、支持批处理和数据流程序处理2、优雅流畅的支持java和scala api3、同时支持高吞吐量和低延迟4、支持事件处理和无序处理通过SataStream API基于DataFlow数据流模型5、在不同的时间语义(事件时间摄取时间、处理时间)下支持灵活的窗口(时间滑动、翻滚会话自定义触发器)6、仅处理一次的容错担保7、自动反压机制8、图处理(批) 机器学习(批) 复杂事件处理(流)9、在dataSet(批处理)API中内置支持迭代程序(BSP)10、高效的自定义内存管理和健壮的切换能力在in-memory和out-of-core中11、兼容hadoop的mapreduce和storm12、集成YARN,HDFS,Hbase 和其它hadoop生态系统的组件3、flink 与 sparkStreaming的区别1、计算模型flink流计算spark微批处理2、时间语义flink事件时间、处理时间spark处理时间3、窗口flink多灵活spark少不灵活窗口必须是批次的整数倍4、状态flink有spark没有5、流式SQLflink有spark没有二、flink系统架构1、JobManager主节点集群大脑负责任务调度、资源管理、Checkpoint 协调、客户端接入1、 内嵌 ResourceManager管理整个集群的 Slot 资源2、可配置 HA多 JobManager 主备避免单点故障。2、TaskManager从节点工作节点独立 JVM 进程负责实际执行计算任务1、每个 TaskManager 管理固定数量的 Slot预配置2、启动后主动向 JobManager 注册领取任务。3、客户端提交作业的节点负责生成 JobGraph 并提交给 JobManager。4、运行流程集群启动先启动 JobManager再启动所有 TaskManagerTaskManager 主动向 JobManager 注册汇报 Slot 资源作业提交客户端执行flink run生成 JobGraph 提交给 JobManager调度执行JobManager 根据作业并行度申请 Slot将 Task 分发到对应 TaskManager任务执行TaskManager 启动线程执行 Task通过网络交互数据定期汇报状态。5、并行度1、定义同一个逻辑算子被拆分成独立的子任务SubTask同时运行的个数为并行度每个子任务处理一部分数据。一个流程序的并行度是其所有算子中最大的并行度。一个程序中不同算子可能有不同的并行度。2、并行度的 4 个设置层级优先级从高到低① 算子级最高优先级在代码中针对单个算子单独设置并行度粒度最细仅作用于当前算子。// 给 map 算子设置并行度为 4dataStream.map(x-x*2).setParallelism(4)② 执行环境级针对整个作业的所有算子设置默认并行度算子级未设置时生效StreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(6);// 整个作业所有算子默认并行度为 6③ 客户端提交级提交作业时通过命令行参数指定./bin/flink run-m JobManager地址:8081-p 8-c 全类名 jar包路径④ 集群配置级最低优先级Flink 集群配置文件flink-conf.yaml中的全局默认值所有作业未指定时生效。parallelism.default: 16、算子链Operator Chain1、定义算子链是 Flink 的执行层优化机制在满足条件的前提下把多个相邻的算子合并成一个算子链作为一个独立的执行任务Task在同一个线程中运行。算子链的并行度 链内所有算子的并行度2、优点① 减少线程切换开销多个算子在同一个线程执行避免线程间切换② 减少网络传输同一算子链内的数据直接在内存中传递不需要序列化 / 反序列化和网络通信③ 降低延迟、提升吞吐量少了跨线程 / 跨节点的交互。3、算子链的形成条件同时满足以下 4 个条件并行度相同上下游算子的并行度必须一致Forward 分区器上下游之间是一对一的数据转发forward分区没有 shuffle、没有数据重分区未禁用算子链代码中没有手动禁用算子链优化同一个 Slot 共享组上下游算子属于同一个 Slot 共享组默认都属于default组。4、自定义算子链①.disableChaining()彻底禁用当前算子的算子链前后都断开env.是全局禁用也可以单个算子禁用②.startNewChain()从当前算子开始开启一条新的算子链和上游断开。7、任务槽Task Slot1、定义任务槽Task Slot是TaskManager 中资源分配的最小单位每个 Slot 对应 TaskManager 的一部分内存资源。2、特点①内存隔离每个 Slot 拥有独立的内存空间不同 Slot 之间的内存不共享避免任务互相抢占内存②CPU 共享Slot 只隔离内存不隔离 CPUCPU 所有 Slot 共享3、Slot 的数量与配置每个 TaskManager 的 Slot 数量是在配置文件flink-conf.yaml里修改参数taskmanager.numberOfTaskSlots:1静态配置的启动后固定不变生产环境通常设置为和 TaskManager 的 CPU 核心数一致4、Slot 共享机制①什么是 Slot 共享同一个作业的不同Task不同算子/算子链可以在同一个 Slot执行前提是它们必须属于同一个 Slot 共享组默认所有算子都属于default共享组自动开启共享可以通过.slotSharingGroup(group1)手动指定分组不同分组之间不能共享 Slot。也就是说一个 Slot 里可以同时运行多个不同算子/算子链的 Task比如 source 算子的 Task 和keyBy 算子的 Task 可以放在同一个 Slot 里同时运行。②Slot 共享的优点提高资源利用率避免轻量算子独占资源造成浪费负载均衡Task 自动均匀分布到各个 Slot避免部分节点过载。三、flink部署模式1、本地模式Local Mode纯开发调试模式所有 Flink 角色都运行在同一个 JVM 进程中用多线程模拟分布式并行不具备真实分布式能力仅用于验证代码逻辑、单元测试。在 IDEA 里直接运行 Flink 代码本质就是本地模式。2、独立集群模式Standalone ModeFlink 自带的独立分布式集群模式不依赖任何外部资源调度框架Flink 自己管理集群资源是最基础的分布式部署形态。如果资源不足或出现故障没有自动扩展或重分配资源的保证必须手动处理。所以一般只用在测试或作业很少的场景1会话模式1、原理需要先启动一个集群保持一个会话在这个会话中通过客户端提交作业集群启动时所有资源就都已经确定所有提交的作业会竞争集群中的资源。比较适用于单个规模小执行时间短的大量作业。#启动集群命令./bin/start-cluster.sh#命令行提交作业./bin/flink run-m JobManager地址:8081-c 全类名 jar包路径#停止集群命令./bin/stop-cluster.sh2应用模式1、原理会话模式代码都是在客户端上执行然后由客户端提交给JobManager的。应用模式直接把应用提交到JobManager上运行。也就代表我们需要为每一个提交的应用单独启动一个JobManager集群。#应用程序的jar包必须放在lib/目录下#执行以下命令启动JobManager./bin/standalone-job.shstart--job-classname 全类名#自己启动TaskManager,需要几个启动几个./bin/taskmanager.shstart#停止集群./bin/taskmanager.sh stop./bin/standalone-job.sh stop3、YARN 部署模式Flink 把资源调度完全交给 YARN按需申请资源和 Hadoop 生态深度打通是目前大厂最主流的部署方式1会话模式1、基于yarn原理客户端把flink应用提交给yarn的ResourceManager,yarn的ResourceManager会向yarn的NodeManager申请容器在容器上flink会部署JobManager和TaskManager的实例从而启动集群。flink会根据运行在JobManager上的作业所需要的slot数量动态分配TaskManager资源。所有作业都提交到这个共享集群里运行#启动一个session集群,对yarn来讲就是启动了一个应用./bin/yarn-session.sh-d-nm sessionTest#命令行提交任务./bin/flink run-d-c 全类名 jar包路径#会自动提交到启动session集群因为/tmp/.yarn-properties文件会记录集群的ip端口直接连接#命令行集群关闭echostop|./bin/yarn-session.sh-id 应用id(application-开头的2单作业模式会话模式因为资源共享会导致很多问题所以为了更好的隔离资源考虑为每个作业启动一个集群就是单作业模式作业完成后集群会关闭所有资源也会释放。但flink本身无法直接这样运行一般需要借助一些资源管理框架来启动集群比如yarnk8s.#命令行提交任务./bin/flink run-d-t yarn-per-job-c 全类名 jar包路径#停止作业./bin/flink list-t yarn-per-job-Dyarn.application.id application开头的id#获取对应的jobid./bin/flink cancel-t yarn-per-job-Dyarn.application.id application开头的id 刚获取的jobid3应用模式前面两种模式下应用代码都是在客户端上执行然后由客户端提交给JobManager的但是这样客户端需要占用大量网络带宽去下载依赖和把二进制数据发送给JobManager很多情况下我们提交作业用的是同一个客户端就会加重客户端所在节点的资源消耗。所以我们不要客户端了直接把应用提交到JobManager上运行。也就代表我们需要为每一个提交的应用单独启动一个JobManager集群。#命令行提交任务./bin/flink run-application-t yarn-application-c 全类名 jar包路径#停止作业./bin/flink list-t yarn-application-Dyarn.application.id application开头的id#获取对应的jobid./bin/flink cancel-t yarn-application-Dyarn.application.id application开头的id 刚获取的jobid四、批处理与流处理实例1、批处理publicclassWordCountBath{publicstaticvoidmain(String[]args)throwsException{//创建执行环境ExecutionEnvironmentenvExecutionEnvironment.getExecutionEnvironment();//读取数据DataSourceStringlinedataenv.readTextFile(文件路径);//切分转换数据 (flink word ,转换成flink,1word,1FlatMapOperatorString,Tuple2String,IntegerwordAndOnelinedata.flatMap(newFlatMapFunctionString,Tuple2String,Integer(){OverridepublicvoidflatMap(Stringvalue,CollectorTuple2String,Integerout)throwsException{//按照空格切分单词String[]wordsvalue.split( );//将单词转换为flink,1)for(Stringword:words){Tuple2String,IntegerwordTuple2Tuple2.of(word,1);//使用Collector向下游发送数据out.collect(wordTuple2);}}});//按照word分组UnsortedGroupingTuple2String,IntegerwordAndOneGroupBywordAndOne.groupBy(0);//分组内聚合AggregateOperatorTuple2String,IntegersumwordAndOneGroupBy.sum(1);//1是位置表示第二个元素//输出sum.print();}}2、流处理publicclassWordCountStream{publicstaticvoidmain(String[]args)throwsException{//创建执行环境StreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();//读取数据DataStreamSourceStringlinedataenv.readTextFile(文件路径);//有界流DataStreamSourceStringsocketdataenv.socketTextStream(hadoop102,8888);//无界流//切分转换数据 (flink word ,转换成flink,1word,1SingleOutputStreamOperatorTuple2String,IntegerwordAndOnelinedata.flatMap(newFlatMapFunctionString,Tuple2String,Integer(){OverridepublicvoidflatMap(Stringvalue,CollectorTuple2String,Integerout)throwsException{//按照空格切分单词String[]wordsvalue.split( );//将单词转换为flink,1)for(Stringword:words){Tuple2String,IntegerwordTuple2Tuple2.of(word,1);//使用Collector向下游发送数据out.collect(wordTuple2);}}});//按照word分组KeyedStreamTuple2String,Integer,StringwordAndOneKeyBywordAndOne.keyBy(newKeySelectorTuple2String,Integer,String(){OverridepublicStringgetKey(Tuple2String,Integervalue)throwsException{returnvalue.f0;}});//分组内聚合SingleOutputStreamOperatorTuple2String,IntegersumwordAndOneKeyBy.sum(1);//1是位置表示第二个元素//输出sum.print();//执行:类似spark streaming最后的.start()env.execute();}}
返回列表