Flink总结_flink中的初始化数据是什么意思啊-程序员宅基地

技术标签: flink  checkpoint  Operator  Java+大数据之旅  state  

Flink总结


一、初步了解什么是Flink?

Flink是一个实时的流式计算引擎,与sparkStreaming不同的是底层是流式引擎,并且有用事件窗口和时间窗口两种窗口,可以进行离线和实时计算,有着完美的容错机制,以及数据延迟机制,在支持高吞吐的同时保证低延迟,并提出了时间语义的概念,将数据分为有界流和无界流,且拥有FlinkSQL方便操作与学习成本。

1、Flink的编程模型

Flink API分层

  • 1、Stateful Stream Processing:是Flink最底层的接口,提供了对时间和状态的细粒度控制,虽然灵活度高,但学习成本高,要求编码能力高
  • 2、DataStream DataSet API:提供了一些封装好的算子,方便使用计算处理分为两种,流式-DataStream API 和 DataSet API 批处理
  • 3、SQL& Table API : 通过构建Table环境,将数据注册成表,直接通过SQL进行编写即可
  • 4、扩展库:复杂事件处理CEP,Gelly做图计算的,是一个可扩展的图形处理和分析库。

Flink 组成:数据源+数据转换+数据输出
Data Source + Transformations + Data Sink
Flink程序整体的流程可以由多个数据源或者多个输出Slink,中间会经过多个算子进行数据的过滤,形成一个有向无环图DAG
在这里插入图片描述


2、Flink的算子Operator

Spark的算子分为:控制算子,行动算子,转换算子。Flink算子划分如下;

  • ① 基本转换算子:map()/filter()/flatmap()
  • ② 键控流转换算子:keyby()/滚动聚合算子(sum/min/max/minBy)/reduce(x+x)
  • ③ 多流转换算子:union():对多条数据合并输出要求数据类型相同不去重。connect():对两条不同的数据流进行合并
  • ④ 分布式算子:Random():将上游数据随机分发给下游。Rescale():将上游数据平分到下游。Rebalance():将上游数据依次分发到下游。Global:将上游数据每一份分发到下游第一个分区。Broadcast():将上游数据所有数据复制发送到下游算子的任务中。
3、富函数

富函数:每个函数处理数据之前都需要进行初始化工作,以及数据处理的事后清理,每个DataStream API提供的所有转换算子都由其富函数版本:
常用函数:RichMapFunction、RichFlatMapFunction、RichFilterFunction
富函数主要提供了额外方法:

  • open():即初始化方法,通常用来只需要一次的初始化工作
  • close():做最后的清理工作
  • getRuntimeContext():提供了函数的一些信息,并行度,子任务等以及分区状态的方法

二、Flink集群架构
1、角色分配以及流程

在这里插入图片描述
流程:
由App发送任务给分发器Dispatcher,再由分发器对任务进行分发,提交给JobManager,JobManager负责本次任务,JM向ResourceManager资源管理者申请资源,RM会将每个集群的资源情况获取到,并分配给JM资源,再由JM将任务分发给子节点上的TaskManager进行执行,TM开始完成任务。


2、TaskSlot与Parallelism

TaskSlot:任务槽,即用于完成任务所用的资源,会根据任务的并行度进行申请资源
Parallelism:并行度,分为算子并行度,环境并行度,客户端并行度,系统并行度
Flink的执行图分层:

  • StreamGraph:根据用户的Stream API编写的代码生成拓扑结构图
  • Job Graph:将多个符合条件的节点chain在一起作为一个节点减少节点之间的IO传输消耗,以及序列化和反序列化、(形成一个操作链)
  • ExecutionGraph:即调度层,最核心的地方由Job Graph的基础上生成
  • 物理执行图:通过具体的组件算子进行计算。

3、Flink的并行度
  • 算子级别:setParallelism()方法定义并行度
  • 执行环境级别:创建环境后.setParallelism()方法
  • 客户端级别:即使用客户端提交任务时指定-p参数来设置并行度
  • 系统级别:通过修改flink的parallelism.default文件来设置并行度

4、窗口机制

首先窗口概念:通过对数据基于时间或者时间的划分,进行计算,便是窗口。
窗口分类:

  • 滑动窗口:滑动窗口在规定时间内进行滑动,会出现重复数据计算
  • 滚动窗口:滚动窗口通过规定时间划分窗口,不会出现重复数据计算
  • 会话窗口:会话窗口不会重叠,没有固定的开始和结束,当窗口一段时间没有接收到数据,则会关闭窗口
  • 全局窗口:将所有相同key的数据分配到单个窗口中计算结果

窗口功能分类:

  • 时间窗口:即设置窗口一次处理多长时间数据,后者窗口滑动、滚动的时间,
  • 事件窗口:即基于事件,一个窗口处理几条事件作为窗口的划分

窗口函数分类:

  • 增量函数:增量指在之前的上个窗口结果的基础上进行当前数据的计算
  • 全量函数:全量指不仅将当前的数据进行计算还有加上历史数据整体进行计算

详解水位线原理—>点击跳转

  • 水位线注意点:单个线程(单数据源)的时候每次获取当前事务中最大的事务时间减去延迟时间来获取水位线,而并发情况下的水位线会获取到最小的水位线向下游广播同步,也是对齐机制。

5、水位线之后迟到的数据怎么办?

现实中很难有一个很完美的水位线将所有的延迟数据都进行挽回,水位线不仅要考虑效率,还要考虑将数据丢失概率降低,从整体的性价比来考量,故此Flink提供了一些机制进行弥补:

  • 直接将延迟数据丢弃
  • 将迟到的数据输出到单独的数据流中,即使用sideOutputLateData(new OutputTag<>())实现测输出
  • 根据迟到的事件更新并发处结果

三、Flink的状态

数据流被分为有状态和无状态,Flink中的算子与状态关联,所有Flink的计算是有状态的,算子会在计算时将自己的状态注册到TaskManager中。
状态分类:
算子状态、键控状态
在这里插入图片描述


1、Flink容错机制

容错机制详解—>跳转


2、State Backends & SavePoint

Flink在保存状态时,支持三种存储方式,如下:

  • MemoryStateBackend (基于内存存储)
  • FsStateBackend (基于文件系统存储)
  • RocksDBStateBackend (基于RocksDB数据库存储)

Savepoint:保存点与CheckPoint类似,一个时系统提供的,一个是用户自己定义,一般由用户进行手动的备份和恢复。


3、Flink流处理的三种语义

at most once : 至多一次,表示一条消息不管后续处理成功与否只会被消费处理一次,那么就存在数据丢失可能。
exactly once : 精确一次,表示一条消息从其消费到后续的处理成功,只会发生一次。
at least once :至少一次,表示一条消息从消费到后续的处理成功,可能会发生多次。


4、Flink之CEP概念

CEP 由一个或者多个规则组成,主要目的就是从有序简单的数据中获取到高阶特征,简单说就是通过数据的表面看数据本质,CEP可以理解为一个数据模型,数据经过CEP模型来获取一定的指标或者数据。(Pattern API )
CEP模式分类:

  • 单个模式:单个模式就是只接受一个事件
  • 循环模式:可以接受多个事件
  • 组合模式:① 严格连续 ② 松散连续 ③ 不确定的松散连续
  • 匹配后跳过策略:对于一个给定的模式,防止同一个事件可能会分配到多个成功的匹配上。

5、Flink 数据反压

Flink1.5版本之前的反压机制
在这里插入图片描述
首先由TaskA 发送数据至TaskB,在TaskA的速率远远大于TaskB时,一定会出现反压情况,首先是TaskB的InputChannel会被填满,此时会向LocalBuffer申请空间,当LocalBuffer也填满后,再向NetworkBuffer申请空间,最后NetworkBuffer没空间后,堆积到Socket,Socket堆满会给发送端发送一个状态,此时发送端停止给Socket发送,TaskA这边的Netty发现Socket满了之后会使用Buffer,最后全部全部缓存用尽,TaskA也停止发数据,实现反压。
缺点:

  • 过于依赖TCP传输,并且反压延迟过高

1.5版本之后
在这里插入图片描述
如图TaskA正常向TaskB发送数据,单每次ResultSubPartition向InputChannel发送消息的时候都会发送一个Backlog size告诉下游准备发送多少数据,下游会告诉上游是否还有足够空间Buffer,当没有足够的空间时则不进行发送。主要降低了反压生效的延迟性,同时Socket不会阻塞。


版权声明:本文为博主原创文章,遵循 CC 4.0 BY-SA 版权协议,转载请附上原文出处链接和本声明。
本文链接:https://blog.csdn.net/llAl_lAll/article/details/123361323

智能推荐

稀疏编码的数学基础与理论分析-程序员宅基地

文章浏览阅读290次,点赞8次,收藏10次。1.背景介绍稀疏编码是一种用于处理稀疏数据的编码技术,其主要应用于信息传输、存储和处理等领域。稀疏数据是指数据中大部分元素为零或近似于零的数据,例如文本、图像、音频、视频等。稀疏编码的核心思想是将稀疏数据表示为非零元素和它们对应的位置信息,从而减少存储空间和计算复杂度。稀疏编码的研究起源于1990年代,随着大数据时代的到来,稀疏编码技术的应用范围和影响力不断扩大。目前,稀疏编码已经成为计算...

EasyGBS国标流媒体服务器GB28181国标方案安装使用文档-程序员宅基地

文章浏览阅读217次。EasyGBS - GB28181 国标方案安装使用文档下载安装包下载,正式使用需商业授权, 功能一致在线演示在线API架构图EasySIPCMSSIP 中心信令服务, 单节点, 自带一个 Redis Server, 随 EasySIPCMS 自启动, 不需要手动运行EasySIPSMSSIP 流媒体服务, 根..._easygbs-windows-2.6.0-23042316使用文档

【Web】记录巅峰极客2023 BabyURL题目复现——Jackson原生链_原生jackson 反序列化链子-程序员宅基地

文章浏览阅读1.2k次,点赞27次,收藏7次。2023巅峰极客 BabyURL之前AliyunCTF Bypassit I这题考查了这样一条链子:其实就是Jackson的原生反序列化利用今天复现的这题也是大同小异,一起来整一下。_原生jackson 反序列化链子

一文搞懂SpringCloud,详解干货,做好笔记_spring cloud-程序员宅基地

文章浏览阅读734次,点赞9次,收藏7次。微服务架构简单的说就是将单体应用进一步拆分,拆分成更小的服务,每个服务都是一个可以独立运行的项目。这么多小服务,如何管理他们?(服务治理 注册中心[服务注册 发现 剔除])这么多小服务,他们之间如何通讯?这么多小服务,客户端怎么访问他们?(网关)这么多小服务,一旦出现问题了,应该如何自处理?(容错)这么多小服务,一旦出现问题了,应该如何排错?(链路追踪)对于上面的问题,是任何一个微服务设计者都不能绕过去的,因此大部分的微服务产品都针对每一个问题提供了相应的组件来解决它们。_spring cloud

Js实现图片点击切换与轮播-程序员宅基地

文章浏览阅读5.9k次,点赞6次,收藏20次。Js实现图片点击切换与轮播图片点击切换<!DOCTYPE html><html> <head> <meta charset="UTF-8"> <title></title> <script type="text/ja..._点击图片进行轮播图切换

tensorflow-gpu版本安装教程(过程详细)_tensorflow gpu版本安装-程序员宅基地

文章浏览阅读10w+次,点赞245次,收藏1.5k次。在开始安装前,如果你的电脑装过tensorflow,请先把他们卸载干净,包括依赖的包(tensorflow-estimator、tensorboard、tensorflow、keras-applications、keras-preprocessing),不然后续安装了tensorflow-gpu可能会出现找不到cuda的问题。cuda、cudnn。..._tensorflow gpu版本安装

随便推点

物联网时代 权限滥用漏洞的攻击及防御-程序员宅基地

文章浏览阅读243次。0x00 简介权限滥用漏洞一般归类于逻辑问题,是指服务端功能开放过多或权限限制不严格,导致攻击者可以通过直接或间接调用的方式达到攻击效果。随着物联网时代的到来,这种漏洞已经屡见不鲜,各种漏洞组合利用也是千奇百怪、五花八门,这里总结漏洞是为了更好地应对和预防,如有不妥之处还请业内人士多多指教。0x01 背景2014年4月,在比特币飞涨的时代某网站曾经..._使用物联网漏洞的使用者

Visual Odometry and Depth Calculation--Epipolar Geometry--Direct Method--PnP_normalized plane coordinates-程序员宅基地

文章浏览阅读786次。A. Epipolar geometry and triangulationThe epipolar geometry mainly adopts the feature point method, such as SIFT, SURF and ORB, etc. to obtain the feature points corresponding to two frames of images. As shown in Figure 1, let the first image be ​ and th_normalized plane coordinates

开放信息抽取(OIE)系统(三)-- 第二代开放信息抽取系统(人工规则, rule-based, 先抽取关系)_语义角色增强的关系抽取-程序员宅基地

文章浏览阅读708次,点赞2次,收藏3次。开放信息抽取(OIE)系统(三)-- 第二代开放信息抽取系统(人工规则, rule-based, 先关系再实体)一.第二代开放信息抽取系统背景​ 第一代开放信息抽取系统(Open Information Extraction, OIE, learning-based, 自学习, 先抽取实体)通常抽取大量冗余信息,为了消除这些冗余信息,诞生了第二代开放信息抽取系统。二.第二代开放信息抽取系统历史第二代开放信息抽取系统着眼于解决第一代系统的三大问题: 大量非信息性提取(即省略关键信息的提取)、_语义角色增强的关系抽取

10个顶尖响应式HTML5网页_html欢迎页面-程序员宅基地

文章浏览阅读1.1w次,点赞6次,收藏51次。快速完成网页设计,10个顶尖响应式HTML5网页模板助你一臂之力为了寻找一个优质的网页模板,网页设计师和开发者往往可能会花上大半天的时间。不过幸运的是,现在的网页设计师和开发人员已经开始共享HTML5,Bootstrap和CSS3中的免费网页模板资源。鉴于网站模板的灵活性和强大的功能,现在广大设计师和开发者对html5网站的实际需求日益增长。为了造福大众,Mockplus的小伙伴整理了2018年最..._html欢迎页面

计算机二级 考试科目,2018全国计算机等级考试调整,一、二级都增加了考试科目...-程序员宅基地

文章浏览阅读282次。原标题:2018全国计算机等级考试调整,一、二级都增加了考试科目全国计算机等级考试将于9月15-17日举行。在备考的最后冲刺阶段,小编为大家整理了今年新公布的全国计算机等级考试调整方案,希望对备考的小伙伴有所帮助,快随小编往下看吧!从2018年3月开始,全国计算机等级考试实施2018版考试大纲,并按新体系开考各个考试级别。具体调整内容如下:一、考试级别及科目1.一级新增“网络安全素质教育”科目(代..._计算机二级增报科目什么意思

conan简单使用_apt install conan-程序员宅基地

文章浏览阅读240次。conan简单使用。_apt install conan