【Flink】概念、入门、部署、架构、流处理API、window、时间语义、Wartermark、ProcessFunction、状态编程、容错、Table API和SQL、CEP、面试题
一、Flink简介
1、概述
Apache Flink 是为分布式、高性能、随时可用以及准确 的流处理应用程序打造的开源流处理框架
2、重要特点
2.1 事件驱动型(Event-driven):从一个或多个事件流提取数据,并 根据到来的事件触发计算、状态更新或其他外部动作
2.2 流处理(无界、实时)与批处理
数据流分为无界数据流和有界数据流(对数据排序,也被称为批处理)
2.3 分层API
有状态流通过过程函数(Process Function) 被嵌入到 DataStream API 中
DataStream API(有界或无界流数据)以及 DataSet API(有界数据集)为数据处理提供了通用的构建模块,比如由用户定义的多种形式的 转换(transformations),连接(joins),聚合(aggregations),窗口操作(windows) 等等。
Table API 是以表为中心的声明式编程,提供可比较的操作,执行之前会经过内置优化器进行优化
3、模块
最高 层 级 的 抽 象 是 SQL,语 法 与 表 达能 力 上 与 Table API 类似
Flink Table API 和 Flink SQL 也并不完善,大多都由各大厂商自己定制
主要学习DataStream API
几大模块:Flink Table & SQL、Flink Gelly(图计算)、Flink CEP(复杂事件处理)
二、快速上手
1、搭建maven工程FlinkTutorial
pom导包:flink-scala_2.11、flink-streaming-scala_2.11
添加Scala文件夹
2、批处理wordcount--DataSet
val inputDS: DataSet[String] = env.readTextFile(inputPath)
val wordCountDS: AggregateDataSet[(String, Int)] = inputDS.flatMap(_.split(" ")).map((_, 1)).groupBy(0).sum(1)
3、流处理wordcount--DataStream
val textDstream: DataStream[String] = env.socketTextStream(host, port)
val dataStream: DataStream[(String, Int)] = textDstream.flatMap(_.split("\\s")).filter(_.nonEmpty).map((_, 1)).keyBy(0).sum(1)
三、Flink部署
1、Standalone 模式
分发配置到集群、8081端口监控管理
数据分发到taskmanager机器
执行程序:./flink run -c com.atguigu.wc.StreamWordCount –p 2 FlinkTutorial-1.0-SNAPSHOT-jar-with-dependencies.jar --host lcoalhost –port 7777
查看结果和计算过程
2、Yarn模式:Session-Cluster 和 Per-Job-Cluster 模式
Hadoop版本>2.2且安装hdfs
Session-Cluster :先启动集群,然后再提交作业,集群会常驻在 yarn 集群中
启动:./yarn-session.sh -n 2 -s 2 -jm 1024 -tm 1024 -nm test -d
执行任务:./flink run -c com.atguigu.wc.StreamWordCount FlinkTutorial-1.0-SNAPSHOT-jar-with-dependencies.jar --host lcoalhost –port 7777
Per-Job-Cluster:每提交一个作业申请一次资源,独享 Dispatcher 和 ResourceManager,按需接受资源申请;任务结束集群消失
直接执行任务:./flink run –m yarn-cluster -c com.atguigu.wc.StreamWordCount FlinkTutorial-1.0-SNAPSHOT-jar-with-dependencies.jar --host lcoalhost –port 7777
3、Kubernetes部署
分别启动Flink的docker组件:JobManager、TaskManager、JobManagerService
启动Session-Cluster:kubectl create -f jobmanager-service.yaml、jobmanager-deployment.yaml、taskmanager-deployment.yaml
通过jobmanager配置访问Flink UI界面
四、Flink运行架构
1、运行时的组件
作 业 管 理 器 ( JobManager ) 、 资 源 管 理 器 ( ResourceManager ) 、 任 务 管 理 器 (TaskManager),以及分发器(Dispatcher)
作业管理器:主进程、将作业图(JobGraph)转化为数据流图/执行图;请求资源、分发、协调
资源管理器:将有空闲插槽的 TaskManager 分配给 JobManager,发起会话、中止释放资源
任务管理器:注册插槽、与同一程序的taskM交换数据
分发器:跨作业运行、提供rest接口
2、任务提交流程(都有ResourceManager)
Yarn 模式任务提交流程
使用yarn的resourcemanager而非自身的
3、任务调度原理
TaskManger 与 Slots:JVM进程、TaskManger上一个固定的子集
程序与数据流(DataFlow):Flink程序-Source 、Transformation 和 Sink,转换 运算(transformations)跟 dataflow 中的算子(operator)是一一对应的关系
执行图(ExecutionGraph):直接映射成的数据流图是 StreamGraph,也被称为逻辑流图,需要转换为物理视图
StreamGraph -> JobGraph -> ExecutionGraph -> 物理执行图
并行度(Parallelism):特定算子的子任务(subtask)的个数,One-to-one类似于窄依赖,Redistributing类似于宽依赖
任务链(Operator Chains):相同并行度的 one to one 操作算子形成一个task,减少线 程之间的切换和基于缓存区的数据交换
五、Flink流处理API
1、Environment
getExecutionEnvironment,客户端
createLocalEnvironment
createRemoteEnvironment
2、Source
从集合读取数据
从文件读取数据
以 kafka 消息队列的数据作为来源
自定义 Source
3、Transform
4、支持的数据类型
5、实现UDF-更细粒度的控制流
6、Sink
六、Flink中的Window
1、Window
2、Window API
七、时间语义与Wartermark
1、Flink中的时间语义
2、EventTime的引入
3、WaterMark
4、EventTime在Window中的应用
八、ProcessFunction API(底层API)
1、KeyedProcessFunction
2、TimerService 和 定时器(Timers)
3、侧输出流(SideOutput)
4、CoProcessFunction
九、状态编程与容错机制
1、有状态的算子和应用程序
2、状态一致性
3、 检查点(checkpoint)
4、选择一个状态后端(state backend)
十、Table API与SQL
1、需要引入的 pom 依赖
2、了解 TableAPI
3、TableAPI 的窗口聚合操作
4、SQL 如何编写
十一、Flink CEP简介
1、什么是复杂事件处理 CEP
2、Flink CEP
十二、常见面试题汇总
架构、监控、数据高峰处理、特性解决、CEP