温馨提示×

温馨提示×

您好,登录后才能下订单哦!

密码登录×
登录注册×
其他方式登录
点击 登录注册 即表示同意《亿速云用户服务条款》

Flink框架的任务调度策略有哪些

发布时间:2026-07-01 18:28:17 来源:亿速云 阅读:91 作者:小樊 栏目:软件技术

Flink 的任务调度策略主要围绕 JobManager 对作业(JobGraph → ExecutionGraph)的调度、TaskManager 对任务的分配与执行、以及运行时对资源与反压的动态调整 展开。下面从 核心调度流程、调度模式、Slot 分配策略、任务部署与执行策略、以及运行时动态调整机制 几个方面系统说明。


一、Flink 任务调度整体架构

1. 调度相关核心组件

  • JobManager
    • 负责将 JobGraph 转换为 ExecutionGraph
    • 生成调度计划(Scheduling)
    • 协调 Checkpoint、故障恢复
  • TaskManager
    • 提供 Slot(资源单位)
    • 执行具体的 Task(算子子任务)
  • ResourceManager
    • 管理 Slot 资源
    • 与 JobManager / TaskManager 交互

二、任务调度流程(宏观)

  1. JobGraph
    • 用户 API 生成的逻辑执行图
  2. ExecutionGraph
    • JobManager 生成的可执行并行图
  3. 调度阶段
    • 决定哪些任务先启动
    • 决定任务部署到哪些 TaskManager
  4. Slot 分配
    • Slot 是 Flink 最小的调度资源单位
  5. 任务部署与执行
    • Task 在 Slot 中运行

三、调度模式(Scheduling Mode)

Flink 提供了多种调度模式,决定 任务何时启动、如何并行启动

1. Eager 调度(默认,流作业)

  • 特点
    • 所有任务一次性调度
    • 所有 ExecutionVertex 同时部署
  • 适用场景
    • 流式作业(Streaming)
  • 优点
    • 启动快
    • 拓扑结构稳定
  • 缺点
    • 资源需求一次性较大
Streaming Job 默认使用 Eager 调度

2. Lazy From Sources 调度(批作业)

  • 特点
    • 从 Source 开始,按需调度下游任务
    • 上游执行完,才调度下游
  • 适用场景
    • 批处理(Batch)
  • 优点
    • 节省资源
    • 支持流水线式执行
  • 缺点
    • 调度逻辑复杂
Batch Job 默认使用 Lazy From Sources

3. Lazy 调度(更细粒度)

  • 可用于特殊批场景
  • 更严格控制执行阶段
  • 在 Flink 1.15+ 中已被更统一的调度器整合

四、Slot 分配与共享策略

Slot 是 Flink 最核心的调度资源单位。

1. Slot 基本概念

  • 一个 TaskManager 有多个 Slot
  • 每个 Slot 可运行一个或多个 Task(取决于共享策略)

2. Slot Sharing(槽共享,核心调度策略)

原则

同一个 Slot Sharing Group 中的任务可以共享一个 Slot

  • 默认情况下:
    • 所有算子属于同一个 Slot Sharing Group
  • 一个 Slot 内:
    • 可以运行一个 Pipeline 中的多个 SubTask

示例

Source → Map → Sink

在一个 Slot 中可同时运行这三个算子的子任务

优点

  • 减少跨 Slot 通信
  • 提高资源利用率
  • 降低线程切换

3. Slot 隔离(Disable Slot Sharing)

  • 通过 API 设置:
.map(...).disableChaining().slotSharingGroup("group1");
  • 特点:
    • 强制不同算子使用不同 Slot
  • 适用场景:
    • 资源敏感算子
    • 避免资源抢占

五、任务链(Operator Chain)策略

任务链是 调度层面的优化策略,影响任务部署方式。

1. 默认 Chain 条件

  • 一对一数据传输
  • 并行度相同
  • 属于同一 Slot Sharing Group

2. Chain 的好处

  • 减少序列化 / 反序列化
  • 减少线程切换
  • 提升性能

3. 控制 Chain

env.disableOperatorChaining();

.map(...).startNewChain();

六、任务部署策略(Deployment)

1. 本地部署优先

  • 同 TaskManager 内任务优先部署在一起
  • 减少网络传输

2. 位置偏好(Locality Preference)

  • 优先将任务调度到:
    • 数据所在节点(如 Kafka 分区本地性)
  • 减少数据移动

七、故障恢复与重启调度策略

1. 重启策略(Restart Strategy)

  • Fixed Delay(固定延迟重启)
  • Failure Rate(失败率重启)
  • No Restart
  • Exponential Delay(指数退避)
env.setRestartStrategy(RestartStrategies.fixedDelayRestart(
    3,
    Time.of(10, TimeUnit.SECONDS)
));

2. Failover 策略

  • Restart All
    • 整个作业重启(默认)
  • Restart Pipelined Region
    • 只重启失败区域(流作业优化)

八、动态扩缩容与自适应调度(新特性)

1. Adaptive Scheduler(Flink 1.15+)

  • 根据可用资源动态调整并行度
  • 适合云原生 / K8s 场景

2. Reactive Mode

  • 根据 TaskManager 数量自动调整并行度
  • 无需手动设置 parallelism

九、总结对比表

调度维度 策略 适用场景
调度模式 Eager 流作业
调度模式 Lazy 批作业
Slot 策略 Slot Sharing 默认推荐
Slot 策略 隔离 Slot 特殊资源需求
任务链 Operator Chain 性能优化
故障恢复 重启策略 稳定性
自适应 Reactive / Adaptive 云原生

如果你愿意,我也可以:

  • 结合 Flink 版本(1.14 / 1.17 / 1.20) 深入说明
  • 源码级调度流程(ExecutionGraph → Scheduler) 解释
  • 或结合 Kafka / 实时数仓 / 批流一体 的实际调度案例讲解
向AI问一下细节

免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。

AI