目录

1、Flink基本概念

2、Flink 的核心工作机制

2.1 核心架构(组件层面)

2.2  部署模式

2.2.1 本地部署

2.2.2 集群部署

(1) Standalone 集群(独立集群)

(2) YARN 部署模式

(3) Kubernetes 部署模式(云原生主流)

2.2.3 作业提交流程

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

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

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

2.3 运行核心逻辑(数据处理层面)

(1)数据输入(Source)

(2)数据处理(Transformation)

(3)状态管理(State)

(4)Checkpoint(检查点)—— 容错核心

(5)数据输出(Sink)


之前提到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 SessionPer-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 等。

    Logo

    汇聚全球AI编程工具,助力开发者即刻编程。

    更多推荐