【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