工业数字化—IoT建设学习笔记(3):大数据应用——Flink
目录
之前提到IoT系统的工作流程,在数据应用部分有提到大数据应用,那么这边学习笔记就介绍一下大数据应用中一个避不开的名词——Flink。
那么Flink到底是什么呢?
1、Flink基本概念
Flink 是一个分布式的流处理和批处理计算框架,核心定位是 “流优先”—— 也就是说它把批处理也看作是 “有界的流处理”,这是它和传统大数据框架(如 Spark)最核心的区别之一。
简单来说:
- 流处理:处理无界数据(比如实时产生的日志、订单、传感器数据,数据源源不断,没有结束点);
- 批处理:处理有界数据(比如一次性读取一个文件、一张数据库表,数据量固定)。

Flink 主要用于:
- 实时数据处理(如实时风控、实时大屏、实时数仓);
- 离线数据计算(替代部分 Hadoop/Spark 批处理场景);
- 复杂事件处理(CEP,比如检测连续的异常行为);
- 状态化计算(比如实时统计用户近 1 小时的点击量)。
Flink 起源于 Stratosphere 项目,Stratosphere 是在 2010~2014 年由 3 所地处柏林的大学和欧洲的一些其他的大学共同进行的研究项目,2014 年 4 月 Stratosphere 的代码被复制并捐赠给了 Apache 软件基金会,参加这个孵化项目的初始成员是 Stratosphere 系统的核心开发人员,2014 年 12 月,Flink 一跃成为 Apache 软件基金会的顶级项目。
第 1 代:Hadoop MapReduc 批处理 Mapper、Reducer 2;
第 2 代:DAG 框架(Oozie 、Tez),Tez + MapReduce 批处理 1 个 Tez = MR(1) + MR(2) + … + MR(n) 相比 MR 效率有所提升;
第 3 代:Spark 批处理、流处理、SQL 高层 API 支持 自带 DAG 内存迭代计算、性能较之前大幅提;
第 4 代:Flink 批处理、流处理、SQL 高层 API 支持 自带 DAG 流式计算性能更高、可靠性更高。
2、Flink 的核心工作机制
Flink 的工作机制可以拆解为 “架构层” 和 “运行层” 两个维度,下面用通俗易懂的方式解释:
2.1 核心架构(组件层面)
Flink 采用 “主从架构”,核心组件包括:JobManager(主节点)、TaskManager(从节点)、Client(客户端)

- JobManager(主节点):相当于 “大脑”,负责全局管理和调度:
- 接收用户提交的作业(Job),并将作业拆解为 “任务图”(StreamGraph → JobGraph → ExecutionGraph);
- 调度 TaskManager 执行任务,监控任务状态,失败时重启任务;
- 管理作业的元数据、状态快照(Checkpoint)。
- TaskManager(从节点):相当于 “手脚”,负责实际计算:
- 启动多个 “Task Slot”(任务插槽,是资源隔离的最小单位),每个 Slot 运行一个或多个任务;
- 执行数据的处理逻辑(比如 map、filter、窗口计算);
- 存储作业的状态(State),并配合 JobManager 完成 Checkpoint。
- Client(客户端):用户提交作业的入口,负责将代码打包并提交给 JobManager,之后客户端可以断开连接(作业由集群接管)。

2.2 部署模式
2.2.1 本地部署
- 适用场景:开发调试、单元测试,无需集群资源。
- 核心特点
- JobManager 和 TaskManager 运行在同一个 JVM 进程中,资源由本地 JVM 分配。
- 无需配置集群,直接通过
StreamExecutionEnvironment.createLocalEnvironment()启动。 - 任务 Slot 数量可自定义(默认 CPU 核心数),适合快速验证代码逻辑。
- 启动方式:直接运行 Flink 应用的
main方法,或通过命令行./bin/flink run -m local [jar包]。
2.2.2 集群部署
集群模式是生产环境的主流选择,根据资源管理器类型分为以下 3 种,核心是将资源管理与计算任务解耦。
| 部署类型 | 资源管理器 | 核心特点 | 适用场景 |
|---|---|---|---|
| Standalone 集群 | Flink 自身 | 独立部署,不依赖第三方资源管理,配置简单;但资源隔离弱,无法动态扩缩容 | 小型生产集群、测试环境 |
| YARN 集群 | Hadoop YARN | Flink 作为 YARN 的应用运行,YARN 负责资源分配;支持 YARN Session 和 Per-Job 两种模式 | 与 Hadoop 生态集成的场景 |
| Kubernetes 集群 | Kubernetes (K8s) | 云原生部署,容器化管理;支持 Session Cluster、Job Cluster、Application Cluster 三种模式;支持动态扩缩容、滚动更新 | 大规模生产环境、云平台部署 |
(1) Standalone 集群(独立集群)
- 架构组成:
- 1 个 JobManager 节点(主节点):负责作业调度和监控。
- 多个 TaskManager 节点(从节点):负责执行计算任务。
- 所有节点通过网络通信,元数据存储在本地文件系统或分布式存储(如 HDFS)。
- 资源调度由 Flink 自己负责,所需要的所有 Flink 组件,都只是操作系统上运行的一个 JVM 进程。支持会话模式和应用模式,唯独不支持单作业模式,因为单作业模式需要其他资源调度器的参与。
(2) YARN 部署模式
YARN 是 Hadoop 生态的资源管理器,Flink 基于 YARN 部署可动态分配资源,支持会话模式、单作业模式和应用模式:

- YARN Session 模式(会话模式)
- 先启动一个长期运行的 Flink 集群(包含 JobManager 和 TaskManager),后续提交的所有作业竞争、共享该集群资源。
- 适合小作业、短作业的批量运行,缺点是资源易被占满,不适合大作业。
- 启动命令:
./bin/yarn-session.sh -n 2 -s 4(2 个 TaskManager,每个 4 个 Slot)

- YARN Per-Job 模式(单作业模式)
- 客户端运行程序,运行作业时再启动集群 为每个作业单独启动一个 Flink 集群,作业结束后集群自动销毁,资源完全隔离,需要借助第三方资源调度器。
- 适合大作业、长作业,是生产环境的推荐方式。
- 启动命令:
./bin/flink run -m yarn-cluster [jar包]。
(3) Kubernetes 部署模式(云原生主流)
Kubernetes 提供容器编排能力,Flink 与 K8s 深度集成,支持三种部署模式:
- Session Cluster:共享集群资源,适合多个小作业。
- Job Cluster:为单个作业创建专用集群,作业结束后销毁。
- Application Cluster:直接在 K8s 上运行 Flink 应用,支持多作业提交,是 Flink 1.15+ 推荐的云原生模式。
- 核心优势:支持动态扩缩容(基于负载自动调整 TaskManager 数量)、故障自动恢复、滚动升级,适配大规模云环境。
2.2.3 作业提交流程

(1)Standalone模式作业提交流程

(2)YARN会话模式作业提交流程

(3)YARN单作业模式任务提交流程

2.3 运行核心逻辑(数据处理层面)
Flink 处理数据的核心是 “数据流” 和 “状态管理”,步骤如下:

(1)数据输入(Source)
Flink 从外部系统读取数据(比如 Kafka、Redis、文件、MySQL),数据以 “流” 的形式进入 Flink,每个数据单元称为 “事件(Event)”。
(2)数据处理(Transformation)
这是 Flink 的核心环节,通过算子(Operator)处理数据,比如:
- 基础算子:map(转换)、filter(过滤)、flatMap(扁平化);
- 聚合算子:keyBy(按 Key 分组)、sum(求和)、reduce(归约);
- 窗口算子:TumblingWindow(滚动窗口)、SlidingWindow(滑动窗口)(处理无界流的核心)。
关键特性:分布式并行处理Flink 会将数据流拆分为多个 “子流”,分配到不同的 TaskManager 的 Slot 中并行处理,提升计算效率。
(3)状态管理(State)
Flink 是 “有状态” 的计算框架 —— 比如统计用户近 5 分钟的点击量,需要记住 “用户已点击的次数”,这个 “次数” 就是状态。
- 状态存储:默认存在 TaskManager 的内存中,也可配置到 RocksDB(持久化);
- 状态一致性:通过 Checkpoint 机制保证。
(4)Checkpoint(检查点)—— 容错核心
Flink 会定期(比如每 10 秒)给整个作业的状态做一个 “快照”,保存到外部存储(比如 HDFS)。如果集群故障,重启后可以从最近的 Checkpoint 恢复状态,保证数据处理的 “精确一次(Exactly Once)” 语义(即数据不会重复处理,也不会丢失)。
(5)数据输出(Sink)
处理完成的数据写入外部系统,比如 Kafka、Elasticsearch、MySQL、Hive 等。
更多推荐




所有评论(0)