尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

Apache Flink 弹性伸缩完全指南:Adaptive Scheduler、Reactive Mode 与 Adaptive Batch Scheduler 实战解析

发布时间:2026/9/24 5:05:41

资讯中心
01
ARTICLE

Apache Flink 弹性伸缩完全指南:Adaptive Scheduler、Reactive Mode 与 Adaptive Batch Scheduler 实战解析

Apache Flink 弹性伸缩完全指南:Adaptive Scheduler、Reactive Mode 与 Adaptive Batch Scheduler 实战解析
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载Flink 作业的并行度历来是静态的批作业根本无法在线调整流作业也只能停作业 → 打 savepoint → 改并行度重启。本文基于 Apache Flink 官方文档 docs/content/docs/deployment/elastic_scaling.md系统讲解三类让 Flink 在运行时动态调整并行度的调度器——流式场景的Adaptive Scheduler、其特化形态Reactive Mode反应式模式以及批式场景的Adaptive Batch Scheduler。读完本文你将掌握它们的启用方式、核心配置参数含源码级默认值、底层资源管理原理与各自的适用边界能够在实际集群中按需选择并落地运行时弹性伸缩。为什么需要运行时弹性伸缩在传统模型下作业的并行度在其生命周期内是固定不变的并在提交时一次性确定批作业完全无法缩放流作业只能通过停止 → 保存点savepoint→ 以不同并行度重启的方式完成一次性的、需要外部编排的调整。这带来两个痛点一是集群资源波动如 TaskManager 故障、提交时资源不足时作业无法自动适配二是资源盈余时无法自动扩容以充分利用集群。为此Flink 引入了两个新的调度器Adaptive Scheduler流式根据可用 slot 动态调整并行度Adaptive Batch Scheduler批式根据输入数据量自动为算子推导并行度。从调度器的选择逻辑可以看出 Flink 的取舍在 DefaultSlotPoolServiceSchedulerFactory.java 的getSchedulerType()中如果作业是批类型JobType.BATCH即使显式配置了adaptive或reactive也会被强制改写为AdaptiveBatch日志会输出 Adaptive Scheduler configured, but Batch job detected. Changing scheduler type to AdaptiveBatch只有scheduler-modereactive才会强制使用流式 Adaptive Scheduler。这从源码层面印证了文档中Adaptive Scheduler 仅支持流作业、批作业走 Adaptive Batch Scheduler的约束。Adaptive Scheduler按可用资源自动伸缩的流式调度器Adaptive Scheduler 的核心能力是根据当前可用的 slot 数量调整作业并行度提交时资源不足、或运行中发生 TaskManager 故障导致 slot 减少时它会自动降低并行度继续运行当新 slot 可用如新的 TaskManager 加入时它会向上扩容直至达到配置的并行度上限在 Reactive Mode 下配置的并行度被忽略、视为无穷大作业永远尽可能多地占用资源。相比默认调度器Adaptive Scheduler 的一大优势是能优雅地处理 TaskManager 丢失——资源减少时直接缩容而不是让整个作业失败。底层原理声明式资源管理Declarative Resource ManagementAdaptive Scheduler 构建在声明式资源管理FLIP-138 引入之上。与传统请求精确数量 slot的交互式模型不同其核心交互如下见下图Dispatcher 提交作业并启动 JobMasterJobMaster 内的 Adaptive Scheduler 向 ResourceManager声明资源需求如上图所示声明 min:1、max:1 的并行度区间Reactive 模式下 max 为无穷大TaskManager 向 ResourceManager 注册自身包含地址与 slot 数量ResourceManager 按声明需求向 JobMaster 的 Slot Pool 提供 slotJobMaster 将任务提交到 TaskManager 上执行。也就是说JobMaster 不再点名要某个槽位而是声明我期望多少资源由 ResourceManager 尽力满足。运行时的自动缩放Rescale当 JobMaster 在运行期间获得更多资源时它会自动基于最新的可用 savepoint 对作业进行缩放全程无需外部编排。上图展示了这一过程声明资源需求为 min:1、max:INF 后多个 TaskManager 加入并向 Slot Pool 提供更多 slotAdaptive Scheduler 据此从检查点/保存点重启任务并调整并行度把任务分发到新增的 TaskManager 上执行。需要特别说明从 Flink 1.18.x 起可以通过下文的外部化声明式资源管理Externalized Declarative Resource Management重新声明运行中作业的资源需求否则 Adaptive Scheduler 无法应对输入速率变化或负载性能变化这类需要主动调整资源诉求的场景。外部化声明式资源管理Externalized Declarative Resource Management注意该功能目前是 MVP最小可行产品形态Flink 社区正在通过邮件列表收集用户反馈请务必关注下文列出的限制。同时它也可以与 Apache Flink Kubernetes Operator 的 autoscaler 配合获得开箱即用的完整自动伸缩体验。该特性旨在解决两个典型部署场景Session 集群上的 Adaptive Scheduler多个作业竞争同一批资源时需要对作业间的资源分配做更细粒度的控制Application 集群 主动式资源管理器如 Native Kubernetes依赖 Flink 贪婪地拉起新的 TaskManager同时仍想保留类似 Reactive Mode 的缩放能力。它通过新增的 REST API 端点/jobs/job-id/resource-requirements实现——允许你按顶点vertex设置并行度上下界从而重新声明运行中作业的资源需求PUT /jobs/job-id/resource-requirements REQUEST BODY: { first-vertex-id: { parallelism: { lowerBound: 3, upperBound: 5 } }, second-vertex-id: { parallelism: { lowerBound: 2, upperBound: 3 } } }从一定程度上该端点可以理解为再缩放端点re-scaling endpoint是构建 Flink 自动伸缩体验的重要基石。在源码层面该端点由 JobResourcesRequirementsUpdateHeaders.java 定义URL 模板为/jobs/:jobid/resource-requirements。你也可以在 Flink Web UI 的 Job 总览页中通过任务列表里的 up-scale / down-scale 按钮手动体验该功能。启用方式与配置项提示如果在 session 集群上使用 Adaptive Scheduler当集群资源不足时多个作业之间的 slot 分配没有任何保证外部化声明式资源管理只能部分缓解该问题因此官方仍建议在 application 集群 上使用 Adaptive Scheduler。要在集群级别启用 Adaptive Scheduler需在配置中设置jobmanager.scheduler: adaptiveAdaptive Scheduler 的全部行为由所有以jobmanager.adaptive-scheduler前缀命名的配置项控制相关配置定义见 JobManagerOptions.java。核心参数汇总如下默认值均来自源码配置项默认值说明jobmanager.adaptive-scheduler.resource-wait-timeout5 分钟Job 提交或重启后JobManager 等待获取全部所需资源的最长时间超时后将以较低并行度运行若连最低资源都无法获取则失败。设为负值可禁用超时无限等待。Reactive 模式下默认改为负值见下文jobmanager.adaptive-scheduler.resource-stabilization-timeout10 秒资源稳定超时当可用资源不足但已足以运行时JobManager 会等待该时长再开始执行避免频繁重启。Reactive 模式下默认改为 0见下文jobmanager.adaptive-scheduler.min-parallelism-increase1触发扩容所需的最小聚合并行度增量。例如 source并行度 2 sink并行度 2聚合并行度为 4默认 1 意味着任何聚合并行度提升都会触发一次重启jobmanager.adaptive-scheduler.scaling-interval.min30 秒两次缩放操作之间的最小间隔防止缩放过于频繁jobmanager.adaptive-scheduler.scaling-interval.max禁用无默认值强制缩放的最大间隔若设置了该值即使聚合并行度增量未达到min-parallelism-increase的要求新资源加入后也会在间隔时间到达时安排一次缩放资源未变化时缩放会被忽略jobmanager.adaptive-scheduler.scale-on-failed-checkpoints-count2连续失败的 checkpoint 数量达到该值时即使没有已完成的 checkpoint 也会触发缩放jobmanager.adaptive-scheduler.max-delay-for-scale-trigger无默认值JobManager 推迟评估已观察到的缩放事件的最大时间默认情况下禁用 checkpoint 时为 0ms启用 checkpoint 时约为checkpoint 间隔 ×该计数 1此外还有判断是否应该缩放的决策组件 RescalingController.javashouldRescale(currentParallelism, newParallelism)它结合上述阈值决定是否真正执行缩放。局限仅限流作业提交批作业时 Flink 会自动改用批作业的默认调度器Adaptive Batch Scheduler不支持部分故障恢复partial failover默认调度器可以只重启失败的部分Flink 内部的 region而 Adaptive Scheduler 会重启整个作业。该限制只会影响易并行embarrassingly parallel作业的恢复时间缩放事件会引发作业与任务重启从而增加任务尝试Task attempt次数。Reactive Mode始终用满集群资源的流式调度模式Reactive Mode 是 Adaptive Scheduler 的一种特殊模式其前提是每个集群只运行一个作业由 Application Mode 强制保证。它的行为非常直观作业永远使用集群中的所有资源——增加一个 TaskManager 作业就扩容移除资源作业就缩容Flink 总是把并行度设置到当前可达到的最高值。与手动缩放相比Reactive Mode 有两个关键差异从最新的已完成 checkpoint 恢复而不是创建 savepoint因此没有 savepoint 的开销缩放后重放的数据量取决于 checkpoint 间隔恢复耗时取决于状态大小。典型的自动伸缩架构Reactive Mode 让外部服务只需关注资源分配与回收作业的存活性完全交给 Flink外部服务监控 consumer lag、聚合 CPU 利用率、吞吐量或延迟等指标指标超过/低于阈值时通过修改 Kubernetes Deployment 的replica 因子、或调整 AWS 的Auto Scaling Group来增删 TaskManagerFlink 负责让作业始终在新资源规模下正常运行。快速上手Getting Started以下步骤假设你在单机上部署 Flink 发行版并位于发行版根目录# 将示例作业放入 lib/ 目录 cp ./examples/streaming/TopSpeedWindowing.jar lib/ # 以 Reactive Mode 提交作业 ./bin/standalone-job.sh start -Dscheduler-modereactive -Dexecution.checkpointing.interval10s -j org.apache.flink.streaming.examples.windowing.TopSpeedWindowing # 启动第一个 TaskManager ./bin/taskmanager.sh start逐一解读上面的提交命令./bin/standalone-job.sh start以 Application Mode 部署 Flink对应发行版脚本 standalone-job.sh-Dscheduler-modereactive启用 Reactive Mode-Dexecution.checkpointing.interval10s配置 checkpoint 间隔与重启策略最后一个参数传入作业的主类名示例为 TopSpeedWindowing.java。启动后访问 Web 界面 可以看到作业运行在一个 TaskManager 上。扩容只需再加入一个 TaskManager./bin/taskmanager.sh start缩容则移除一个 TaskManager./bin/taskmanager.sh stop配置详解启用与并行度规则启用 Reactive Mode 只需把scheduler-mode配置为reactive。作业中单个算子的并行度完全由调度器决定不可配置——即使显式设置了算子级或作业级并行度也会被忽略。影响并行度的唯一途径是为算子设置 max parallelism调度器会尊重该上限它被限制在 2^1532768以内。如果你不为算子或整个作业设置 max parallelism将应用默认并行度规则其推导出的下限可能低于最大值。与默认调度模式一样请参考并行度最佳实践建议显式设置 max parallelism。注意过高的 max parallelism 可能影响作业性能因为 Flink 需要维护更多内部结构来支持可缩放状态。关键配置项与 Reactive 模式下的默认值变化启用 Reactive Mode 时源码 JobManagerOptions.java 中如下配置的默认值会发生变化jobmanager.adaptive-scheduler.resource-wait-timeout默认变为-1即 JobManager 会永远等待足够资源出现。若希望资源不足一段时间后停止作业请显式配置该超时jobmanager.adaptive-scheduler.resource-stabilization-timeout默认变为0一旦资源足够就立即开始运行。注意如果 TaskManager 不是同时接入、而是陆续接入这会导致每次接入都触发作业重启如需等待资源稳定再调度请调大该值jobmanager.adaptive-scheduler.min-parallelism-increase默认 1即任何聚合并行度提升都会触发重启。若不想让微小资源变化引发频繁重启可调大该值jobmanager.adaptive-scheduler.scaling-interval.max默认禁用。设置后即使聚合并行度增量不满足min-parallelism-increase也会在间隔到达时强制安排缩放jobmanager.adaptive-scheduler.scaling-interval.min默认 30 秒用于限制两次缩放的最小间隔避免缩放过于频繁。推荐实践有状态作业务必配置周期性 checkpointReactive Mode 在缩放事件中从最近一次完成的 checkpoint恢复若不启用周期性 checkpoint程序将丢失状态。同时 checkpoint 配置会顺带定义重启策略——Reactive Mode 尊重配置的重启策略若未配置任何重启策略作业将直接失败而不是缩放缩容耗时与心跳超时如果 TaskManager 没有优雅关闭例如使用 SIGKILL 而非 SIGTERMFlink 需要等待 JobManager 与被停 TaskManager 之间的心跳超时作业会卡住约 50 秒后才以较低并行度重新部署。默认心跳超时为 50 秒基础设施允许时可调低heartbeat.timeout但注意过低的超时可能因网络拥塞或长 GC 停顿导致误判失败且heartbeat.interval必须始终小于超时值。局限Reactive Mode 仍是实验性特性并非默认调度器的全部能力都可用仅支持 Standalone 形态的 Application 部署主动式资源提供者Native Kubernetes、YARN明确不支持Standalone session 集群也不支持且 Application 部署仅限单作业应用。目前支持的部署方式包括Standalone Application Mode即上文快速上手的方式、Docker Application Mode 与 Standalone Kubernetes Application ClusterAdaptive Scheduler 的全部局限同样适用于 Reactive Mode。Adaptive Batch Scheduler自动推导并行度的批式调度器Adaptive Batch Scheduler 是批作业调度器能够自动调整执行计划目前的核心能力是为批作业算子自动决定并行度如果某个算子未显式设置并行度调度器会根据其消费的数据集大小决定并行度。收益体现在三方面批作业用户免于手动调并行度自动调优的并行度能更好适配每天变化的数据量SQL 批作业中的不同算子可以获得各自经过自动调优的并行度。目前 Adaptive Batch Scheduler 是 Flink 批作业的默认调度器无需额外配置除非显式指定了其他调度器如jobmanager.scheduler: default。但注意你需要保持execution.batch-shuffle-mode不设置、或显式设置为ALL_EXCHANGES_BLOCKING默认值、ALL_EXCHANGES_HYBRID_FULL或ALL_EXCHANGES_HYBRID_SELECTIVE——因为目前它只支持 BLOCKING/HYBRID 混洗模式见下文局限。在源码 AdaptiveBatchScheduler.java 中它直接继承自DefaultScheduler并扩展了并行度与输入信息决策逻辑。自动并行度决策使用方式要让 Adaptive Batch Scheduler 自动为算子决定并行度需要1. 开启该特性自动并行度推导默认开启可通过execution.batch.adaptive.auto-parallelism.enabled开关控制源码默认true见 BatchExecutionOptions.java。此外还有几个关联配置默认值均来自源码配置项默认值说明execution.batch.adaptive.auto-parallelism.min-parallelism1自适应设置的并行度下限execution.batch.adaptive.auto-parallelism.max-parallelism128自适应设置的并行度上限未配置时用parallelism.default或StreamExecutionEnvironment#setParallelism()设置的默认并行度作为上限execution.batch.adaptive.auto-parallelism.avg-data-volume-per-task16 MiB期望每个任务实例处理的数据量均值。注意发生数据倾斜、或并行度已达上限数据过多时部分任务实际处理的数据可能远超该值execution.batch.adaptive.auto-parallelism.default-source-parallelism无默认值source 的默认并行度或自适应设置的 source 并行度上限。未配置时使用max-parallelismmax-parallelism也未配置时回退到parallelism.default或setParallelism()设置的默认并行度这些配置项的旧名称如jobmanager.adaptive-batch-scheduler.min-parallelism、jobmanager.adaptive-batch-scheduler.max-parallelism、jobmanager.adaptive-batch-scheduler.avg-data-volume-per-task、jobmanager.adaptive-batch-scheduler.default-source-parallelism已标记为废弃源码中以withDeprecatedKeys(...)兼容。2. 避免为算子显式设置并行度Adaptive Batch Scheduler只为未设置并行度的算子决定并行度。因此若希望算子并行度被自动决定就不要通过setParallelism()设置它。对 DataSet 作业还有额外要求设置parallelism.default: -1不要在ExecutionEnvironment上调用setParallelism()。为 Source 启用动态并行度推断新的 Source 可以通过实现接口 DynamicParallelismInference.java 来启用动态并行度推断public interface DynamicParallelismInference { int inferParallelism(Context context); }Context会向实现提供三类辅助信息见 DynamicParallelismInference.javagetParallelismInferenceUpperBound()推断并行度的上限getDataVolumePerTask()期望每个任务处理的数据量字节getDynamicFilteringInfo()动态过滤信息辅助并行度推断。Adaptive Batch Scheduler 会在调度 source 顶点之前调用该接口因此实现应尽量避免耗时操作。如果 Source 未实现该接口则使用execution.batch.adaptive.auto-parallelism.default-source-parallelism作为 source 顶点的并行度。与算子同理动态 source 并行度推断只作用于未显式设置并行度的 source 算子。性能调优建议推荐使用 Sort Shuffle并设置taskmanager.network.memory.buffers-per-channel为0这可以把所需网络内存与并行度解耦大型作业中更不容易出现 Insufficient number of network buffers 错误建议把execution.batch.adaptive.auto-parallelism.max-parallelism设置为最坏情况下你预期的并行度过大的值会影响性能——该选项会影响上游任务产生的 subpartition 数量大量 subpartition 会因小包问题降低 hash shuffle 与网络传输性能。局限仅限批作业提交流作业会抛出异常仅支持 BLOCKING 或 HYBRID 作业只支持 shuffle 模式为ALL_EXCHANGES_BLOCKING/ALL_EXCHANGES_HYBRID_FULL/ALL_EXCHANGES_HYBRID_SELECTIVE的作业。对于不识别上述 shuffle 模式的 DataSet 作业需要把 ExecutionMode 设为BATCH_FORCED以强制使用 BLOCKING shuffle不支持 FileInputFormat source包括StreamExecutionEnvironment#readFile(...)、readTextFile(...)与createInput(FileInputFormat, ...)这些方法定义于 StreamExecutionEnvironment.java 附近。使用 Adaptive Batch Scheduler 读取文件时应改用新式 sourceFileSystem DataStream Connector 或 FileSystem SQL ConnectorWebUI 上 broadcast 结果指标不一致自动决定并行度时broadcast 结果中上游任务发送的字节数/记录数与下游任务接收到的指标计数不一致可能在 Web UI 上造成困惑详见 FLIP-187Adaptive Batch Job Scheduler。如何选择三种调度形态的适用场景形态作业类型并行度如何决定典型场景Adaptive Scheduler流式按可用 slot 在配置上限内自动伸缩Session/Application 集群中资源波动大、需要优雅应对 TaskManager 故障Reactive Mode流式单作业/集群始终用满集群所有资源Standalone 部署 外部监控指标驱动的自动扩缩容Adaptive Batch Scheduler批式按消费数据量自动推导每日数据量波动的 SQL/DataSet 批作业免手动调并行度三者共享的底层理念是把并行度管理从用户手中交还给调度器让 Flink 更接近真正的云原生流批一体处理器Adaptive Scheduler 与 Reactive Mode 负责流作业在资源波动下的自愈与弹性Adaptive Batch Scheduler 负责批作业在数据规模波动下的自动调优。实际选型时请依据作业类型、部署形态Standalone / 主动式资源管理器 / Session / Application、以及对部分故障恢复与文件 source 的兼容性要求综合判断。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink 弹性扩缩容完全指南Adaptive Scheduler、Reactive 模式与 Adaptive Batch Scheduler 深度解析Flink 弹性扩缩容完全指南Adaptive Scheduler、Reactive 模式与 Adaptive Batch Scheduler 深度解析 导读大数据流处理批处理数据工程Flink Standalone 模式在 Kubernetes 上的部署实战Session / Application 集群、Kubernetes HA 与 Reactive 弹性伸缩Flink Standalone 模式在 Kubernetes 上的部署实战Session / Application 集群、Kubernetes HA 与大数据流处理批处理数据工程Dask 自适应部署Adaptive Deployments完全指南让集群规模随计算需求动态伸缩Dask 自适应部署Adaptive Deployments完全指南让集群规模随计算需求动态伸缩 导读 本文讲解 Dask 中自适应部署Adapti大数据数据分析任务调度创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。