资讯动态

Java面试——Flink原理及应用(一)

发布时间:2026/8/27 8:20:39 来源:尧图企业网站定制
Flink原理及应用1、Flink的原理1.1、Flink的特点1.2、Flink的核心概念1.2.1、Flink Cluster1.2.2、Flink Master1.2.3、Flink Job Manager1.2.4、Flink Task Manager1.2.5、Job1.2.6、Flink Graph1.2.7、Flink Operator和Operator Chain1.2.8、Flink Task和SubTask1.2.9、Event1.2.10、Function1.2.11、Flink Record1.2.12、Flink State Backend1.3、Flink架构介绍1.3.1、Job Manager的职责1.3.2、Task Manager的职责1.3.3、客户端1.3.4、Flink应用程序的运行流程1.3.5、Flink Task Slot资源分配1.3.6、Flink任务和算子1.3.7、Flink状态存储1.3.8、Flink运行模式1.4、Flink的事件驱动模型1.4.1、什么是事件驱动模型1.4.2、事件驱动模型的特点1.4.3、Flink的事件驱动模型的特点1.5、Flink的数据分析应用1.5.1、批量数据分析1.5.2、流式数据分析1.5.3、Flink中的流式数据分析1.6、Flink的数据清洗和数据管道1.6.1、数据清洗定时处理1.6.2、数据管道实时处理1.6.3、数据管道优势实时性高1.6.4、Flink数据管道1.7、Flink数据流处理基本概念1.7.1、数据流1.7.2、状态1.7.3、时间1.8、API分类1.8.1、Process Function1.8.2、DataStream API1.8.3、SQL和Table API1.9、Flink基于状态的内存计算1.10、Flink的编程模型1.10.1、数据流1.10.2、并行度1.11、Flink窗口计算1.11.1、翻滚窗口1.11.2、滑动窗口1.11.3、会话窗口1.11.4、全局窗口1.12、Flink故障恢复1.12.1、重启策略1.12.2、故障恢复策略1、Flink的原理Flink是一个分布式计算引擎主要用于有界数据流和无界数据流的有状态的数据分析和处理。Flink擅长处理有界数据流和无界数据流其精确的时间控制和状态化使Flink能够安全并快速地处理海量数据。Flink将数据抽象为有界数据流和无界数据流。1无界数据流无界数据流只定义了数据流的开始没有定义数据流的结束即随着时间的推移数据会源源不断地产生具体如图11-1所示。比如滴滴打车中车辆的实时位置信息只要车辆在运行就有数据源源不断地上报。无界数据流必须被实时地处理因为数据会源源不断地产生如果处理不及时就会产生数据积压。在数据的摄取过程中无界数据流需要指定数据的时间顺序例如Event Time事件时间​、Process Time处理时间​、Ingestion Time摄取时间​以便能够保障数据处理的顺序性和完整性。2有界数据流有界数据流既定义了数据流的开始时间也定义了数据流的结束时间具体如图11-2所示。有界数据流在一批数据获取完成后再执行计算其数据在内部可以被排序因此有界数据流不需要指定数据的顺序。1.1、Flink的特点Flink是有状态的计算框架也就是说运行中的Flink应用程序即使因为某些原因执行出错也可以通过快照快速还原到故障点继续执行计算。Flink被广泛应用于事件驱动、流式分析、批量分析、数据清洗Extract-Transform-LoodETL等领域可用于处理事件流数据、用户行为数据、物联网设备数据、系统日志数据等各种类型的数据。Flink应用程序以事件流或数据库文件系统、关系型数据库、非关系型数据库、Key-Value数据库等的形式获取数据在内存中对数据执行分布式计算并最终将结果写入应用程序、事件日志或其他数据库系统。Flink可运行于Kubernetes、YARN、Mesos等环境中其支持的存储有HDFS、S3、NFS等数据库具体如图所示。Flink的主要特点为支持丰富的流式计算应用场景保障良好的数据正确性基于分层API的设计理念运维方便支持超大规模计算和延迟低、吞吐量高。1支持丰富的流式计算应用场景Flink支持事件驱动模型、流式计算、批量计算、数据管道和数据清洗等多种应用场景。Flink被广泛应用于阿里巴巴、AWS、滴滴出行等互联网公司的大规模计算中。2保障良好的数据正确性Flink支持Exactly-One精确一次数据一致性保障用户也可以自定义事件时间以保证数据处理的顺序性避免因网络延迟或者其他原因导致数据乱序。3基于分层API的设计理念根据抽象程度的不同Flink将API分为High-Level Analytics API、Stream AndBatch Data Processing API、Stateful Event-DrivenApplications三种不同的API。每一种API在简洁性和表达力上都有着不同侧重。4运维方便Flink部署灵活支持单机和集群模式部署可运行于Kubernetes、YARN、Mesos等多种资源管理框架上同时支持定时创建CheckPoint以方便程序故障恢复支持手动创建SavePoint以方便应用程序安全升级。5支持超大规模计算Flink基于Task Manager水平扩展当资源不足时添加计算资源即可添加集群的过程不会影响其他程序的运行。6延迟低、吞吐量高Flink基于内存计算其延迟低吞吐量高。1.2、Flink的核心概念1.2.1、Flink ClusterFlink Cluster集群是用于运行Flink应用程序的分布式系统一个Flink集群由ZooKeeper、Job Manager和Task Manager 3个角色组成。在高可用模式下一般ZooKeeper为至少3个节点的集群Job Manager为至少2个节点的集群Job Manager高可用模式为一主多备。在正常情况下主节点提供服务当主节点宕机时一个备节点会升级为主节点对外提供服务。Task Manager为具体的计算节点一个集群中有一个或多个TaskManager。1.2.2、Flink MasterFlink Master指集群的管理节点一个Flink Master由Flink Resource Manager资源管理​、FlinkDispatcher分发和Flink Job Manager 3个角色组成。1.2.3、Flink Job ManagerFlink Job Manager是Flink的任务管理节点用于任务的提交、分发和运行状态的监控。一个集群可以有一个或多个高可用模式下Job Manager。1.2.4、Flink Task ManagerFlink Task Manager是Flink集群的计算节点Flink的任务被Job Manager调度在多个Task Manager上执行。多个Task Manager上的Task相互交换计算结果以完成数据流的计算。1.2.5、JobJob任务指一个运行中的Flink应用程序Job可以通过Job Manager以命令行的方式提交到集群也可以通过Flink监控页面提交到集群。1.2.6、Flink GraphFlink Graph图指Flink流式计算程序所组成的数据流程图。Flink Graph分为Logical Graph和PhysicalGraph。Logical Graph描述的是应用程序通常指基于Java或者Scala实现的Flink程序定义的数据流之间的逻辑关系与之对应的是逻辑上的Operator算子​、Input、Output、DataStream和DataSet。PhysicalGraph是Logical Graph基于分布式运行环境转换后物理上的计算逻辑图与之对应的是物理计算上的Task、Input、Output、DataStream和DataSet。1.2.7、Flink Operator和Operator ChainFlink Operator即Flink算子是Flink Logical Graph中的 一个节点用于执行一个Flink Function。一个完整的数据流通常包含一个Source Operator数据摄取​、Process Function数据计算和Sink Operator数据输出​。多个相邻Operator相互连接组成一个Operator Chain在一个Operator Chain内部Operator的数据可以相互直接访问不需要经过Flink集群的序列化和网络传输。1.2.8、Flink Task和SubTaskFlink Task是Physical Graph的一个节点对应物理上的一个计算单元一个Task由多个SubTask组成每个SubTask都对应数据流上的一个处理函数。1.2.9、EventEvent事件指Flink运行时数据模型的状态变化Event在流式计算和批量计算接口中被输入、输出以完成状态的记录和传递。1.2.10、FunctionFunction函数是由应用程序实现的逻辑计算单元Function一般通过实现Flink的Function接口或继承Function类来定义。常用的Function有MapFunction、ReduceFunction、ProcessFunction、RichFunction等。1.2.11、Flink RecordFlink Record记录指数据流中的元素。1.2.12、Flink State BackendFlink State Backend定义了Task Manager上运行的Job的状态存储方式例如JVM堆内存存储、RocksDB、FileSystem​同时定义了SavePoint和CheckPoint的存储规则和存储方式。1.3、Flink架构介绍Flink由Job Manager、Task Manager和客户端组成。Job Manager是管理节点负责集群任务的提交、分配和资源管理Task Manager是具体执行任务的计算节点客户端用于作业的提交。1.3.1、Job Manager的职责Job Manager负责协调分布式计算节点也被称为Master节点。它负责调度任务、协调CheckPoint、故障恢复等。Job Manager将一个作业分为多个Task并通过Actor系统与Task Manager进行相互通信用于Task的部署、停止、取消。在高可用部署下会有多个Job Manager其中有一个Leader、多个Flower。Leader总是处于Active状态为集群提供服务Flower处于Standby状态在Leader宕机后会从Flower中选出一个作为Leader继续为集群提供服务。Job Manager选举通过ZooKeeper实现。1.3.2、Task Manager的职责Task Manager也被称为Worker节点用于执行JobManager分配的Task准确来说是SubTask​。TaskManager将系统资源CPU、网络、内存分为多个TaskSlot计算槽​Task运行在具体的Task Slot上TaskManager通过Actor系统与Job Manager进行相互通信定期将Task的运行状态和Task Manager的运行状态提交 给Job Manager。多个Task Manager上的Task通过DataStream进行状态计算和结果交互。1.3.3、客户端客户端不是运行时Runtime环境的一部分它主要用于提交作业给Job Manager。在作业提交完成后客户端可以断开连接也可以保持连接来接收作业的运行状态。Flink的架构如图所示。1.3.4、Flink应用程序的运行流程由图可知Flink应用程序的运行流程如下。1编写应用程序数据流作业可以基于Java或者Scala。2构建DAG和优化流程。3通过客户端命令提交作业到集群的Job ManagerLeader节点。4Job Manager将任务提交的结果和运行状态反馈给客户端。5Job Manager根据每个Task Manager上资源的使用情况将作业拆分为多个Task并通过Actor系统将Task部署到具体的Task Manager节点上6Task Manager在Task Slot上运行Task并定时将Task Manager的运行状态和Task的运行状态发送给JobManagerJob Manager根据Task Manager上资源的使用情况和Task的运行状态对集群进行调度。7Job Manager和ZooKeeper进行交互以完成JobManager的选举和故障恢复。1.3.5、Flink Task Slot资源分配Job Manager将一个Task拆分为多个SubTask分配到不同Task Manager上运行。Task Manager是一个JVM进程内部以多线程的方式运行不同SubTask。在Task Manager内部将操作系统计算资源分为多个Slot每个Slot都代表一份固定的计算资源TaskManager为每个SubTask都分配一定的Slot进行运算多个SubTask之间不会竞争资源很好地实现了资源隔离。需要强调的是Slot的资源并不是CPU的资源划分而是内存资源的划分。可以通过设置Task Manager中Slot的数量来配置多个SubTask的资源隔离方式。如果Task Manager中Slot的数 量为1则意味着Task Manager上的所有SubTask共享所有资源例如将Task Manager运行在一个独立的Docker中​。如果Task Manager中Slot的数量为多个则多个SubTask在同一个Task Manager的JVM中运行TaskManager将JVM内存分为多份分配给每个Slot每个Slot上都运行不同SubTask以实现SubTask之间的内存资源隔离SubTask在同一个JVM中共享TCP连接通过多路复用技术和心跳信息同时各个SubTask之间可以共享数据集和数据结构从而减少每个Task的开销。FlinkTask Slot的计算资源如图所示。官方建议将Slot的数量设置为与CPU核的数量相同这样Flink能够更好地使用操作系统资源。当使用超线程时每个Slot都将会占用2个或更多的硬件线程上下文。1.3.6、Flink任务和算子Flink将多个SubTask连接在一起组成一个OperatorChain每个Operator Chain都对应一个Task如果Operator Chain只包含一个SubTask则该Task内只有一个SubTask。在Flink中每个Task都运行在一个线程上Flink将连续相邻的多个SubTask组成一个Operator Chain能有效减少多个SubTask之间的线程数据交换、数据共享、线程上下文切换所带来的时间和性能的消耗从而提高系统整体 的吞吐量。Flink任务和算子的结构如图所示。1.3.7、Flink状态存储Flink在状态存储中以Key-Value的形式对数据的状态进行存储存储的方式有内存、RocksDB、本地磁盘、HDFS、S3、OSS等多种方式。状态存储还可按照应用程序的配置定时将计算状态存储到CheckPoint中以便基于状态进行故障恢复。Flink还支持以命令行或API的方式手动将当前运行中的集群状态存储为SavePoint。Flink应用程序在启动时可以从指定的SavePoint启动即基于历史状态启动​。SavePoint常用于在程序升级过程中保障数据状态的一致性和完整性。SavePoint和CheckPoint的不同在于CheckPoint是程序按照配置定期自动生成的并且在新的CheckPoint生成后旧的CheckPoint将不再被需要而SavePoint需要用户手动触发用户可以从任何一个SavePoint状态中恢复计算。Flink CheckPoint的存储流程如图所示。1Job Manager根据用户配置定期触发CheckPoint存储并将执行CheckPoint的任务下发给Task Manager。2Task Manager在收到保存CheckPoint任务后进行快照备份Save CheckPoint​。3Task Manager在执行完快照备份后将备份结果及状态上报给Job Manager。1.3.8、Flink运行模式Flink有多种运行模式具体包括Local、StandaloneCluster、On YARN、On Mesos如表所示。1.4、Flink的事件驱动模型1.4.1、什么是事件驱动模型事件驱动模型是基于事件流的有状态的计算模型。它接收源源不断的事件并根据事件的不同类型更新不同状态来触发不同计算。事件驱动模型和一般的计算存储分离模型的最大差别在于计算存储分离模型需要将数据存储在远程对象存储系统如S3、OSS、OBS​、事务型数据库MySQL、Postgres​、分布式内存系统Redis、Memcached、LeverDB中。即计算存储分离模型的所有数据计算都基于本地内存和磁盘而数据均被存储在远程存储系统中这样做的好处是计算和存储分离以便计算和存储可以各自扩展而相互不影响。计算存储分离模型的架构如图所示。事件驱动模型是基于状态化流处理完成的它并未将计算和存储分离而是在计算过程中通过访问本地存储操作系统内存或磁盘获取数据以便尽可能快地完成计算。事件驱动模型系统通过定期向远程持久化存储写入CheckPoint和SavePoint来实现状态回滚、故障恢复和程序升级架构如图所示。1.4.2、事件驱动模型的特点事情驱动模型的特点是高效。事件驱动模型由于不需要频繁访问远程数据大部分数据操作在本地内存中完成少部分数据操作在磁盘上完成因此具有更高的吞吐量和更低的延迟。同时事件驱动模型会定期增量地将数据处理的状态以CheckPoint的形式存储在远程持久化存储中以便程序状态回滚和故障恢复。1.4.3、Flink的事件驱动模型的特点Flink底层基于有状态的事件数据处理而设计提供了丰富的状态操作并提供了Exactly-One数据一致性保障和海量规模至少TB级的状态数据计算。此外Flink提供了多种窗口计算灵活性极高。1.5、Flink的数据分析应用数据分析的目的是从数据中获取有价值的信息一般将数据分析分为批量数据分析和流式数据分析。1.5.1、批量数据分析批量数据分析指以批量查询的方式将文件系统HDFS​、对象存储系统S3/OSS/OBS或数据库MySQL/Postgres中的数据加载到计算模型中进行计算然后将计算结果写入外部数据库存储或以报表的形式输出。Hadoop的MapReduce和Spark的Batch Process都属于批量数据分析。由于批量计算一次只能计算一部分数据一段时间内的数据或者符合某些查询条件的数据​因此为了获取最新的数据计算结果必须将数据重新加载到内存中进行计算并输出。批量数据分析如图所示。1.5.2、流式数据分析1什么是流式数据分析。流式数据分析主要用于接入实时事件流不断地处理到达系统的数据并实时更新计算状态和结果其特点为数据是“无穷无尽”的即随着时间的流逝新的数据会不断地流入并被分析。流式数据分析的中间结果会以状态的形式存储在内存中或者定期存储到本地磁盘上计算结果往往最终被发送到高效的消息系统、Key-Value数据库或实时报表。我们知道从CPU到内存、本地磁盘再到网络计算性能是呈指数级下降的因此流式数据分析设计的“大忌”是在计算过程中依赖远程的某个运行效率并不高的接口或数据库从而造成系统的瓶颈进而影响整个系统的运行效率。流式数据分析如图所示。2流式数据分析的特点。①高效与批量数据分析相比流式数据分析不用频繁地执行数据的批量查询和导入因此性能更高延迟更低。②易于维护批量数据分析通常根据业务的不同将任务分为多组件周期性的调度执行各个任务之间根据时间上的先后顺序和依赖关系组成一个数据处理流当其中一个任务执行出错时将影响后续所有其他任务的执行。而流式数据分析整体运行在一个流式数据分析框架上也就是说一个流式数据分析可以独立地完成从数据接入到连续的数据计算和最终计算结果的输出。因此流式数据分析可以独立地完成一个工作流任务的各个节点而不依赖其他任务的计算结果每个节点的运行状态都被记录在流式数据分析框架中比如Flink的CheckPoint​在其中某个节点执行出错后可以快速地从故障节点重新计算更加易于维护。1.5.3、Flink中的流式数据分析Flink为流式数据分析和批量数据分析都提供了良好的支持。Flink内置了一个符合ANSI标准的SQL接口将批量、流式查询的语义抽象在一起。无论批量查询还是实时事件流查询相同的SQL查询都会得到一致的处理结果。同时Flink支持丰富的用户自定义函数允许在SQL中执行定制化代码。如果需要进一步定制逻辑则可以通过Flink DataStream API和DataSet API进行底层的控制。此外Flink的Gelly库为基于批量数据集的大规模高性能图分析提供了算法和构建模块的支持。1.6、Flink的数据清洗和数据管道系统之间数据的转换和迁移常用的两种方式是数据清洗和数据管道。1.6.1、数据清洗定时处理数据清洗通常调用周期性的任务对数据进行清洗处理通常从事务型的数据库中获取数据通过提取、转换、加载将最终结果输出到分析型数据库或其他数据仓库中。数据清洗过程如图所示。1.6.2、数据管道实时处理数据管道和数据清洗的功能类似均用于数据的转换和迁移不同的是数据管道中的数据以实时流的模式运行因此其运行不会停止。数据管道主要用于从一个源源不断的数据流上接收事件数据流并通过实时转化将数据以低延迟的方式传递到目标节点目标节点可以是数据库、文件系统或者日志系统。数据管道常用于日志监控、实时数据格式转换等。数据管道的数据处理过程如图所示。1.6.3、数据管道优势实时性高数据管道主要用于端到端的数据转换和迁移主要特点是吞吐量高和延迟低另外与数据清洗不同的是数据管道用于持续性的消息处理因此其处理数据的实时性比较高。1.6.4、Flink数据管道应用程序可以通过Flink的DataStream API、SQL接口或者自定义函数实现数据管道应用。同时Flink提供了多种连接器如Kafka、Kinesis、ElasticSearch、JDBC数据库系统等用于数据管道应用的开发。1.7、Flink数据流处理基本概念Flink是一个针对无界数据流和有界数据流进行有状态计算的框架以数据流、状态State​、时间Time3个核心组件为基础构建整个框架。1.7.1、数据流流是流式数据处理的基本要素。1有界流和无界流Flink数据流可以分为有界流和无界流有界流定义了数据的结束时间无界流没有定义数据的结束时间。2实时数据流和历史数据流Flink中的所有数据都是以流的方式产生的但是用户处理数据的方式可以分为实时数据流处理和历史数据流处理。实时数据流处理指将接收的数据进行“立即”处理而历史数据流处理指先将数据流持久化到存储系统中例如内存、文件系统、数据库​然后对一段时间内的数据执行批量处理。1.7.2、状态每一个具有一定复杂度的流处理应用都应该是有状态的。任何运行基本业务逻辑的流处理应用都需要在一定时间内存储所接收的事件或中间结果以供在后续某个时间点例如收到下一个事件或者经过一段特定时间进行访问并进行后续处理。事件状态如图所示。状态是Flink的核心概念。基于状态Flink提供了多种特性。1多种状态基础类型Flink为多种不同数据结构提供了对应的状态基础类型如原子值、列表及映射。2插件化的状态管理Flink通过State Backend实现状态的管理并按照应用程序的配置进行CheckPoint。Flink支持内存、RocksDBRocksDB是一种高效的嵌入式、持久化Key-Value存储引擎​、本地磁盘或S3这些状态管理组件以插件化的形式存在Flink中应用程序也可以通过自定义State Backend完成状态存储。3精确一次语义Flink的CheckPoint和故障恢复算法保证了故障发生后应用状态的一致性。因此Flink能够在应用程序发生故障时快速恢复而不会造成数据不一致。4大规模数据状态管理Flink基于异步增量式的Checkpoint算法可以进行TB级应用数据存储状态的管理。5可弹性伸缩Flink能够通过在节点扩容或缩容对状态进行重新分布支持应用状态的分布式横向扩缩。1.7.3、时间时间是流处理应用的一个重要组成部分因为数据总在某个时间点产生或事件总在某个时间点发生。Flink提供了丰富的时间语义支持时间语义分为Event Time事件时间​、Ingestion Time摄取时间​、Process Time处理时间​。不同的时间语义如图所示。1事件时间使用事件时间语义的流处理应用根据事件本身自带的时间戳进行结果计算。因此无论处理的是历史记录事件还是实时事件事件时间模式的处理总能保证结果的准确性和一致性。2摄取时间数据进入Flink数据源的时间。3处理时间除了事件时间和摄取时间Flink还支持处理时间语义。处理时间模式根据处理引擎的时钟触发计算一般应用于对数据处理延时性要求较高并且能够容忍因网络延迟带来数据乱序问题的场景。由于网络延迟或者其他原因在数据计算完成后仍可能会有数据到达这样的数据被称为迟到数据。Flink提供了WaterMark机制以衡量数据处理的进展。基于该机制当应用程序收到迟到数据时可以将这些数据重新定向到旁路输出Side Output或者更新之前完成计算的结果。1.8、API分类根据抽象程度的不同Flink将API封装成ProcessFunction、DataStream API和SQL/Table API如图所示。其中Process Function是Flink底层接口接口内部通过Event、State、Time来实现有状态的事件计算DataStream API是Flink基于Process Function的进一步封装其内部通过定义不同的Stream和Window来实现流式计算和批量计算最上层的是SQL/Table APISQL/Table API将数据封装成一个Dynamic Table以便应用更加简单、方便地实现业务数据的分析。1.8.1、Process FunctionProcess Function是Flink所有应用底层的API。它将实时收到的事件数据流分配到不同窗口中进行计算。对于每条事件数据Flink都提供了时间和状态细粒度的控制应用程序可以获取每条计算中事件数据的状态并更新该状态值同时基于时间窗口应用程序可以定义触发器在未来的某个时刻执行数据计算。Process Function定义的代码如下。上述代码通过继承Process Function类实现一个ProcessFunction函数计算其中open方法用于初始化计算需要的全局资源或配置在每个实例初始化时都执行一次processElement用于执行具体数据的处理逻辑Value为算子的输入数据Collector为算子的输出数据Context为算子的上下文信息保存了数据的状态、时间等信息。在数据计算完成后通过out.collect​将计算结果数据输 出到下一个算子。Process Function函数计算在被定义好后只需要在数据流中加入对应的Process Function实例即可具体实现如下。上述代码定义了一个DataStream并将NumberClassFunction应用到数据流中。1.8.2、DataStream APIDataStream API为流处理提供了许多通用的操作。这些操作包括窗口计算、数据的转换操作、执行外部数据库查询等。DataStream API预先定义了map​​、reduce​​、aggregate​等函数。应用程序可以通过扩展实现自定义接口或使用Java、Scala的lambda表达式实现自定义函数。下面的代码展示了如何捕获会话时间范围内的所有点击数据流事件并对每一次会话的点击量进行计数。1.8.3、SQL和Table APIFlink支持SQL和Table API两种关系型API。在无边界的实时数据流和有边界的历史记录数据流上关系型API会以相同的语义执行查询并产生相同的结果。SQL和TableAPI借助Apache Calcite进行查询的解析、校验及优化。SQL和Table API可以与DataStream API和DataSet API无缝集成并支持用户自定义的标量函数、聚合函数及表值函数。Flink的关系型API旨在简化数据分析、数据流处理和ETL应用。下面的代码示例展示了如何使用SQL语句查询会话时间范围内的所有点击数据流事件并对每一次会话的点击量都进行计数。1.9、Flink基于状态的内存计算Flink是有状态的分布式计算框架在运行中任务的状态始终保留在内存中如果状态大小超过可用内存则会被保存在能高效访问的磁盘中同时Flink可以将状态信息存储在远程存储中。运行中的任务通过访问保存在内存或高效磁盘上的状态来完成下一步的状态计算因此大部分计算都在内存中完成其吞吐量高计算延迟低。Flink通过定期异步地对本地状态进行持久化存储来保证发生故障时数据状态的精确一致性。Flink的内存计算模型如图所示。1.10、Flink的编程模型1.10.1、数据流Flink的编程模型由流Stream和转换Transformation组成。流由一系列记录Record组成转换是针对流的操作通常以一个或者多个流作为数据源输入Flink集群经过Map或Reduce等计算再将计算结果输出到另外一个流。当运行Flink应用程序时Flink会根据流的依赖关系生成一个数据流。每个数据流从一个或者多个数据源开始以一个或者多个数据输出结束。数据流类似DAG具体如图所示。1.10.2、并行度一个Flink程序由多个Task组成。一个Task包括多个并行执行的实例且每一个实例都处理Task输入数据的一个子集。一个Task的并行实例数被称为该Task的并行度Parallelism​具体如图所示。在开发过程中可以通过改变特定算子或者整个程序的并行度来提高系统的运行效率和吞吐量。Flink默认算子的并行度为1。在单个算子上设置的并行度作用于单个算子在执行环境上设置的并行度作用于所有算子具体设置如下。同时可以在将作业提交到Flink集群后在客户端通过命令行设置并行度具体实现如下。还可以在conf/flink-conf.yaml中设置parallelism.default参数来设置Flink集群的并行度。Flink Job运行过程中并行度的优先级为算子并行度执行环境并行度客户端并行度Flink集群默认并行度。1.11、Flink窗口计算Flink中的窗口将源源不断的流数据切分为有限大小的分段每个分段中都包含一部分数据Flink基于该分段上的数据执行计算。Flink中的窗口可以分为翻滚窗口TumblingWindow-无重叠滚动窗口Sliding Window-有重叠会话窗口Session Window-活动间隙以及全局窗口Global Window​。1.11.1、翻滚窗口翻滚窗口的尺寸是固定的不能重叠。例如如果指定一 个大小为5s的翻滚窗口则Flink会每5s启动一个新窗口并针对窗口内的数据执行计算翻滚窗口如图所示。对翻滚窗口的定义如下。1.11.2、滑动窗口滑动窗口分配器将每个元素都分配到固定大小的窗口中。滑动窗口由窗口大小和滑动长度两个参数控制。因此如果滑动长度小于窗口大小则滑动窗口内的数据将会重叠。在这种情况下元素将被分配到多个窗口。例如将窗口大小设置为10s将滑动长度设置为5s。这样每10s都会生成一个新的窗口该窗口将包含最近10s内到达的数据。滑动窗口如图所示。对滑动窗口的定义如下。1.11.3、会话窗口会话窗口分配器通过活动会话分配元素。与翻滚窗口和滑动窗口相比会话窗口不会重叠也没有固定的开始时间和结束时间。当会话窗口在一段时间内没有接收元素时窗口会关闭。会话窗口通过设置会话间隙不活动的时间长度来确定窗口大小当该时间段到期时当前会话关闭后续元素将被分配到新的会话窗口。会话窗口如图所示。对会话窗口的定义如下。1.11.4、全局窗口全局窗口将具有相同Key的所有元素都分配到同一个全局窗口。全局窗口只有在定义触发器后才能触发计算。否则只会将数据分配给不同的窗口但不会执行任何计算。全局窗口如图所示。对全局窗口的定义如下。1.12、Flink故障恢复当Task发生故障时Flink需要重启出错的Task及其他受到影响的Task以便作业恢复到正常执行状态。Flink通过重启策略Restart Strategies和故障恢复策略Failover Strategies来控制Task重启。重启策略决定是否可以重启及重启的间隔故障恢复策略决定哪些Task需要重启。1.12.1、重启策略Flink作业如果没有定义重启策略则会遵循集群启动时加载的默认重启策略。如果提交作业时设置了重启策略则该策略将覆盖集群的默认重启策略。我们可以通过Flink的配置文件flink-conf.yaml来设置集群的默认重启策略通过配置参数restart-strategy定义重启策略。如果没有启用CheckPoint则采用“不重启”策略。如果启用了CheckPoint且没有配置重启策略则采用固定延迟重启策略。除了定义默认重启策略我们还可以为每个Flink作业都单独定义重启策略。这个重启策略通过在程序中的ExecutionEnvironment对象上调用setRestartStrategy方法来设置。当然对于StreamExecutionEnvironment也同样适用。Flink重启策略有3种没有重启策略Disable​、固定延迟重启Fixeddelay​、故障率重启Failurerate​。每个重启策略都有自己的一组配置参数来控制其行为。1固定延迟重启。固定延迟重启策略按照给定的次数尝试重启作业。如果尝试超过了给定的最大次数则作业将最终失败。在连续两次重启尝试之间重启策略等待一段固定长度的时间。在flink-conf.yaml中设置任务发生故障时延迟10s执行、重试3次的配置参数如下。在应用程序中配置固定延迟重启策略的代码如下。2故障率重启策略。故障率重启策略在故障发生之后重启作业但是当故障率每个时间间隔发生故障的次数超过设定的限制时作业会最终失败。在连续两次重启尝试之间重启策略会等待一段固定长度的时间。在flink-conf.yaml中设置如下配置参数。在应用程序中配置故障率重启策略的代码如下。1.12.2、故障恢复策略Flink目前支持两种不同的故障恢复策略全局重启故障恢复Full​、按分区重启Region​。该策略需要在Flink配置文件flink-conf.yaml的jobmanager.execution.failoverstrategy配置项进行配置。具体配置如下。1全局重启故障恢复在全局重启故障恢复策略下当Task发生故障时会重启作业中的所有Task来进行故障恢复。2按分区重启该策略会将作业中的所有Task划分为多个Region。当有Task发生故障时它会尝试找出进行故障恢复需要重启的最小Region集合。相比于全局重启故障恢复策略按分区重启策略在一些场景下需要重启的Task会更少。此处的Region指以Pipeline的形式进行数据交换的Task集合Batch形式的数据交换会构成Region的边界。Region的重启判断逻辑如下。1出错Task所在的Region需要重启。2如果要重启的Region需要消费的数据有部分无法访问丢失或损坏​则产生该部分数据的Region也需要重启。3需要重启的Region的下游Region也需要重启。这是出于保障数据一致性的考虑因为一些非确定性的计算或者分发会导致同一个Result Partition每次产生时包含的数据都不相同。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价