云计算百科
云计算领域专业知识百科平台

十五、Flink 核心原理与架构详解

概述

Apache Flink 是一个开源的分布式流处理框架,被设计用于在有界和无界数据流上进行有状态的计算。与传统的微批处理(Micro-Batch)引擎不同,Flink 从底层就将流处理作为一等公民——批处理仅仅是有界流的特例。这种"流批一体"的设计哲学使 Flink 能够在低延迟、高吞吐和精确一次(Exactly-Once)语义之间取得卓越的平衡。

Flink 的核心能力包括:有状态计算、事件时间处理、轻量级容错、高吞吐低延迟、以及灵活的窗口和状态管理。目前已被阿里巴巴、Uber、Netflix、LinkedIn 等公司在生产环境中大规模使用,是实时计算领域的事实标准。


一、整体架构分层

Flink 的软件架构是一个分层、分布式、主从式的设计,从上到下可以分为四层:

1.1 部署层(Deployment Layer)

Flink 可以运行在多种环境中:

  • Standalone 模式:Flink 自带资源管理,不依赖外部系统
  • Flink on YARN:利用 Hadoop YARN 作为资源调度器
  • Flink on Kubernetes:原生云原生部署,支持弹性伸缩
  • 本地模式:单进程运行,用于开发和调试

1.2 编程接口层(API Layer)

Flink 提供三层 API,抽象程度从高到低:

层级API适用人群特点
最高层 SQL & Table API 数据分析师 声明式,自动优化,流批统一语义
中间层 DataStream API 应用开发者 精细控制流处理逻辑,核心编程模型
最底层 Stateful ProcessFunction 高级开发者 直接操作状态和定时器,最灵活

SQL & Table API 是最常用的接口,支持标准 SQL 语法,Flink 内部将其编译优化为 DataStream 程序执行。DataStream API 提供了 map、filter、keyBy、window、join、process 等丰富的算子,是构建复杂流处理应用的基础。

1.3 运行时层(Runtime Layer)

运行时层是 Flink 的核心引擎,负责作业调度、资源管理和分布式执行。

1.4 存储层(Storage Layer)

包括状态后端(RocksDB/Heap)、Checkpoint 存储(HDFS/S3/OSS)和高可用存储(ZooKeeper/K8s)。


二、核心组件详解

2.1 JobManager(Master 节点)

JobManager 是 Flink 集群的"大脑",负责整个作业的生命周期管理。其内部包含四个核心子组件:

子组件职责
Dispatcher 接收用户提交的作业,持久化 JobGraph,为每个作业创建独立的 JobMaster
JobMaster 负责单个作业的全生命周期管理,将 JobGraph 展开为 ExecutionGraph 并调度执行
CheckpointCoordinator 周期性触发分布式快照,协调所有 TaskManager 完成 Checkpoint
ResourceManager 管理 TaskManager 及其 Task Slot,与外部资源系统(YARN/K8s)交互申请资源

在高可用(HA)模式下,集群可以部署多个 JobManager,通过 ZooKeeper 或 Kubernetes 进行 Leader 选举,确保只有一个 Active JobManager,其余为 Standby。

2.2 TaskManager(Worker 节点)

TaskManager 是 Flink 的"工人",所有实际计算在此发生。每个 TaskManager 管理若干 Task Slot(任务槽):

  • Task Slot 是资源调度的基本单位,代表 TaskManager 内存的一个固定子集
  • 每个 Slot 运行一个 Task(可包含多个通过算子链融合的 Subtask)
  • 一个 Subtask 是一个算子的一个并行实例

TaskManager 的关键模块:

  • Operator Chain:将多个算子融合在同一线程中执行,减少序列化/网络开销
  • Network Stack:基于 Netty + Credit-based 流控,管理 Task 间数据传输
  • Memory Manager:管理 Managed Memory(用于排序、哈希、状态存储等)
  • State Backend:存储和管理算子状态

2.3 Client(客户端)

Client 不是运行时的组成部分,它负责:

  • 解析用户代码,构建 StreamGraph
  • 优化为 JobGraph(算子链合并)
  • 通过 HTTP REST 接口提交给 Dispatcher
  • 支持 CLI / REST API / Web UI 三种操作方式

三、作业执行流程

Flink 作业从提交到运行,经历四层图转化:

用户代码

StreamGraph(逻辑拓扑,包含所有算子节点和边)
↓ 算子链合并优化
JobGraph(减少不必要的序列化/反序列化)
↓ 提交到 Dispatcher
ExecutionGraph(物理执行计划,包含并行实例和网络拓扑)
↓ 调度执行
Task 在 TaskManager 的 Slot 中运行

详细流程:

  • 客户端构建 StreamGraph:用户通过 DataStream API 或 SQL 编写作业逻辑,Client 解析算子 DAG,构建纯逻辑的 StreamGraph
  • 优化为 JobGraph:Client 对 StreamGraph 进行算子链(Operator Chain)合并优化,生成可提交的 JobGraph
  • 提交给 Dispatcher:JobGraph + 用户 Jar 包通过 HTTP REST 接口发送至 Dispatcher
  • 创建 JobMaster:Dispatcher 将 JobGraph 持久化到 HA 存储,创建专属 JobMaster
  • 展开为 ExecutionGraph:JobMaster 将 JobGraph 展开为 ExecutionGraph,包含所有并行实例和网络连接
  • 调度执行:JobMaster 向 ResourceManager 申请 Slot 资源,将 Task 调度到 TaskManager 执行
  • 数据流处理:Task 之间通过 Pipeline 数据交换模式进行实时数据传输
  • 状态监控与容错:CheckpointCoordinator 定期触发快照,TaskManager 持续上报心跳

  • 四、核心机制

    4.1 算子链(Operator Chain)

    Flink 为了优化性能,会将多个满足条件的算子融合为一个 Task,在同一个线程中执行。这称为算子链(Operator Chaining)。

    融合条件(必须同时满足):

    • 算子之间是 One-to-One 连接(如 map → filter → keyBy 中的 map → filter)
    • 算子具有相同的并行度
    • 算子属于相同的 Slot 共享组
    • 没有被用户通过 disableChaining() 显式禁用

    好处: 避免不必要的序列化/反序列化和网络传输开销,显著降低延迟。

    4.2 并行度(Parallelism)

    Flink 的每个算子都可以独立设置并行度。一个算子的多个并行实例称为 Subtask(子任务),它们可以运行在不同的 TaskManager 上。

    Source(并行度=3)→ Map(并行度=3)→ KeyBy/Window(并行度=2)→ Sink(并行度=1)

    Source-0 ──┐
    Source-1 ──┼── Map-0 ──┐
    Source-2 ──┘ Map-1 ──┼── Window-0 ── Sink-0
    Map-2 ──┘ Window-1 ──┘

    • 并行度可以在代码中设置,也可以在全局配置中指定
    • Flink 2.0 进一步支持了动态并行度调整(Adaptive Scheduler),根据运行时负载自动优化

    4.3 数据交换模式

    Flink 在 Task 之间采用 Pipeline 数据交换模式:

    • One-to-One(Forwarding):数据在上下游 Subtask 之间一一对应传输,保持分区不变
    • Redistributing(Shuffle):数据经过重分区后传输,包括:
      • Hash 分区(keyBy 后的默认方式)
      • Broadcast(广播到所有下游)
      • Global(全部发往下游第一个 Subtask)
      • Rescale(在局部范围内轮询分发)

    Pipeline 交换的核心优势:数据产生后立即推送给下游,实现真正的流式处理,延迟极低。


    五、状态管理

    状态管理是 Flink 最核心的能力之一,使流式应用能够"记住过去的事件"并影响未来处理。

    5.1 状态类型

    类型作用域说明典型场景
    Keyed State 按 Key 分区 每个 Key 维护独立的状态实例 用户累计消费金额
    Operator State 算子级别 整个算子共享一份状态 Kafka Consumer 的 offset

    Keyed State 的子类型:

    • ValueState:存储单个值
    • ListState:存储一个列表
    • MapState:存储键值对映射
    • ReducingState:通过 ReduceFunction 聚合
    • AggregatingState:通过 AggregateFunction 聚合

    5.2 状态后端

    后端存储位置适用场景特点
    HashMapStateBackend JVM 堆内存 状态较小、对延迟敏感 访问快,但受内存限制
    EmbeddedRocksDBStateBackend(Flink 1.x) 本地磁盘(嵌入式 RocksDB) 状态较大(GB~TB 级) 不受内存限制,但访问有磁盘 I/O 开销

    5.3 状态 TTL

    可以为状态设置过期时间(Time-To-Live),自动清理不再使用的状态数据,防止状态无限膨胀:

    StateTtlConfig ttlConfig = StateTtlConfig
    .newBuilder(Time.hours(24))
    .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
    .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
    .cleanupInRocksdbCompactFilter(1000)
    .build();


    六、时间语义与窗口

    6.1 三种时间语义

    时间类型含义特点
    Event Time 事件发生的时间(数据自带) 最精确,能处理乱序数据,依赖 Watermark
    Processing Time 数据到达算子时的系统时间 最简单,延迟最低,但结果不确定

    6.2 Watermark(水位线)机制

    Watermark 是 Flink 处理乱序数据的核心机制,用于衡量事件时间的进展:

    数据流: [t=10] [t=12] [t=11] [t=15] [W(t=13)] [t=14] [t=16] [W(t=15)]
    ↑ ↑
    Watermark=13 Watermark=15

    • Watermark = 当前最大事件时间 – 允许的最大乱序时间
    • 当 Watermark 达到或超过窗口结束时间时,触发窗口计算
    • 晚于 Watermark 的数据被视为迟到数据,可通过侧输出流(Side Output)处理

    Watermark 生成策略:

    • BoundedOutOfOrderness:允许固定时间的乱序(最常用)
    • Monotonous Timestamps:假设时间戳单调递增(无乱序)
    • 自定义策略:根据业务特点自定义 Watermark 生成逻辑

    6.3 窗口类型

    窗口类型特点示例
    滚动窗口(Tumbling) 固定大小、无重叠 每5分钟统计一次
    滑动窗口(Sliding) 固定大小、有重叠 每1分钟统计最近5分钟的数据
    会话窗口(Session) 基于活动间隙动态划分 用户无操作超过30分钟则关闭窗口
    全局窗口(Global) 默认不触发,需自定义 Trigger 自定义触发逻辑

    窗口生命周期:

    数据进入 → 分配到窗口 → 触发器(Trigger)决定是否触发计算
    → 窗口函数执行 → 驱逐器(Evictor)可选清理数据 → 输出结果


    七、容错机制(Checkpoint)

    7.1 Checkpoint 原理

    Flink 的容错基于 Chandy-Lamport 分布式快照算法的变体——异步屏障快照(Asynchronous Barrier Snapshotting)。

    核心流程:

  • JobManager 触发 Checkpoint:CheckpointCoordinator 周期性地向所有 Source 算子注入 Barrier(屏障标记)
  • Barrier 随数据流传递:Barrier 作为数据流中的特殊标记,随数据一起流向下游
  • 算子收到 Barrier 后持久化状态:每个算子收到 Barrier 后,将当前状态快照异步写入后端存储(HDFS/S3)
  • Barrier 对齐(Exactly-Once 模式):当算子有多个输入时,需等待所有输入的 Barrier 到达后才执行快照
  • 确认完成:所有算子完成快照后,JobManager 收到确认,该 Checkpoint 完成
  • 7.2 Checkpoint vs Savepoint

    特性CheckpointSavepoint
    触发方式 自动周期性触发 用户手动触发
    用途 故障恢复 作业升级/迁移/版本变更
    存储位置 由 Flink 管理 由用户管理(外部存储)
    格式 轻量级,优化性能 标准化,支持跨版本兼容
    生命周期 旧 Checkpoint 自动清理 永久保留直到手动删除

    7.3 非对齐 Checkpoint(Unaligned Checkpoint)

    传统对齐 Checkpoint 的问题:当存在**反压(Backpressure)**时,Barrier 在反压节点堆积,导致 Checkpoint 超时。

    非对齐 Checkpoint 的解决方案:

    • Barrier 不再等待所有输入对齐
    • 未对齐的数据也被包含在状态快照中
    • 代价:状态快照更大,但 Checkpoint 时间与反压无关

    八、窗口函数与聚合

    8.1 增量聚合函数

    函数说明
    ReduceFunction 每次输入与当前累积值合并,输出新的累积值
    AggregateFunction 更灵活的聚合,支持累加器(Accumulator)类型
    FoldFunction 类似 ReduceFunction,但输入和累加器类型可以不同(已废弃)

    8.2 全量窗口函数

    函数说明
    ProcessWindowFunction 接收窗口内所有数据,提供全局上下文(窗口起止时间、状态、定时器等)
    ApplyWindowFunction 简化版的 ProcessWindowFunction

    8.3 组合使用示例

    dataStream
    .keyBy(event -> event.getUserId())
    .window(TumblingEventTimeWindows.of(Time.minutes(5)))
    .aggregate(
    new CountAgg(), // 增量聚合:快速计算中间结果
    new WindowResultFunction() // 全量函数:附加窗口元信息
    );


    九、Connector 与数据源

    Flink 通过 Connector 与外部系统集成:

    Connector方向说明
    Kafka Source + Sink 最常用的消息队列 Connector,支持 Exactly-Once
    JDBC Source + Sink 关系型数据库读写
    Elasticsearch Sink 写入 ES 索引
    HDFS Source + Sink 文件系统读写,支持 Parquet/ORC 等格式
    Redis Source + Sink 键值存储交互
    CDC(Change Data Capture) Source 实时捕获数据库变更(MySQL/PostgreSQL/Oracle)
    Paimon Source + Sink 流式湖存储,与 Flink 深度集成

    Flink CDC 是一个重要的 Connector,基于 Debezium 实现,支持:

    • 全量 + 增量一体化读取数据库变更
    • Schema Evolution(DDL 变更自动同步)
    • 多表整库同步
    • 支持写入 Paimon、Iceberg 等湖格式

    十、典型应用场景

    10.1 实时数仓(Real-time Data Warehouse)

    Flink 是实时数仓的核心计算引擎:

    • 实时 ETL:从 Kafka 消费原始数据,清洗转换后写入 DWD 层
    • 实时聚合:在 DWD 层基础上计算 DWS 汇总指标
    • 实时报表:将计算结果写入 Doris/ClickHouse,供 BI 看板查询

    10.2 实时风控

    • 实时交易监控:毫秒级检测异常交易模式
    • 反欺诈:基于规则引擎和 CEP(复杂事件处理)实时识别欺诈行为
    • 信用评估:结合用户历史行为和实时交易进行动态评分

    10.3 实时监控与告警

    • 系统监控:实时采集和分析系统指标
    • 业务监控:实时监控 GMV、DAU 等核心业务指标
    • 异常检测:基于滑动窗口和 CEP 识别异常模式

    10.4 实时推荐

    • 用户行为实时分析:基于用户最近的行为序列更新推荐模型
    • 特征工程:实时计算用户画像特征
    • A/B 测试:实时分流和效果评估

    10.5 数据集成与 CDC

    • 数据库实时同步:通过 Flink CDC 将 MySQL 变更实时同步到数仓
    • 异构系统数据迁移:在不同存储系统之间实时同步数据
    • Schema Evolution:源表 DDL 变更时自动同步到目标系统

    十一、安装部署

    11.1 环境要求

    组件版本要求
    Java JDK 8 或 JDK 11(Flink 2.0 要求 JDK 11+)
    Hadoop 2.x 或 3.x(使用 HDFS 时需要)

    11.2 下载与解压

    下载地址:https://archive.apache.org/dist/flink/flink-1.17.1/ 在这里插入图片描述

    下载完后上传到 hadoop1 的 /opt/software/ 目录下,然后使用 tar 命令解压。 在这里插入图片描述

    修改配置文件 flink-conf.yaml

    cd /opt/module/flink-1.17.1/conf
    vim flink-conf.yaml

    修改下面的条目

    # JobManager节点地址.
    jobmanager.rpc.address: hadoop1
    jobmanager.bind-host: 0.0.0.0
    rest.address: hadoop1
    rest.bind-address: 0.0.0.0
    # TaskManager节点地址.需要配置为当前机器名
    taskmanager.bind-host: 0.0.0.0
    taskmanager.host: hadoop1

    修改配置文件 workers 为如下内容:

    hadoop1
    hadoop2
    hadoop3

    修改配置文件 masters

    hadoop1:8081

    11.3 分发到其它机器

    cd /opt/module/
    xsync flink-1.17.1

    在其它机器上修改 flink-conf.yaml 文件中的 taskmanager.host 为自己的主机名:hadoop2、hadoop3

    11.4 启动集群

    在 hadoop1 上执行命令 Flink 集群

    cd /opt/module/flink-1.17.1
    bin/start-cluster.sh

    在这里插入图片描述

    11.5 访问 Web UI

    启动成功后,可以访问 http://hadoop1:8081 对flink集群和任务进行监控管理

    在这里插入图片描述


    十二、与 Spark Streaming 对比

    维度Apache FlinkSpark Streaming
    处理模型 真正的流处理(逐条处理) 微批处理(Micro-Batch)
    延迟 毫秒级 秒~分钟级
    容错机制 Checkpoint + Barrier(Chandy-Lamport) 基于 RDD 血缘重算
    状态管理 原生支持,多种后端 依赖外部存储(如 HBase)
    时间语义 Event Time + Watermark 支持但实现较复杂
    窗口模型 丰富的窗口类型,原生支持 有限的窗口支持
    批处理 有界流(统一引擎) 原生批处理引擎
    SQL 能力 完善,流批统一 成熟,生态丰富
    生态成熟度 快速成长中 非常成熟,社区庞大
    适用场景 实时计算、低延迟、有状态流处理 批处理、ETL、对延迟要求不高的场景

    关键要点总结

  • 架构设计:分层架构(API → Runtime → Storage),主从式分布式集群(JobManager + TaskManager)
  • 流批一体:批处理是有界流的特例,统一引擎、统一语义、统一 API
  • 状态管理:核心能力,支持 Keyed State 和 Operator State,多种状态后端可选
  • 容错机制:基于 Chandy-Lamport 算法的 Checkpoint,支持 Exactly-Once 语义
  • 时间语义:Event Time + Watermark 机制精确处理乱序数据
  • 窗口模型:滚动、滑动、会话、全局窗口,灵活覆盖各类时序计算场景
  • 算子链优化:融合算子减少序列化/网络开销,降低延迟
  • Connector 生态:丰富的数据源连接,Kafka/CDC/Paimon 是数仓场景的核心
  • 赞(0)
    未经允许不得转载:网硕互联帮助中心 » 十五、Flink 核心原理与架构详解
    分享到: 更多 (0)

    评论 抢沙发

    评论前必须登录!