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

Flink on YARN 部署全指南:从 Session/Application 模式到高可用与资源调优

发布时间:2026/9/24 8:37:52

资讯中心
01
ARTICLE

Flink on YARN 部署全指南:从 Session/Application 模式到高可用与资源调优

Flink on YARN 部署全指南:从 Session/Application 模式到高可用与资源调优
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载Apache Flink 原生支持将作业集群部署到 Apache Hadoop YARN 之上利用 YARN 的资源管理能力实现 JobManager / TaskManager 容器的动态分配与失败恢复。本文基于 Flink 官方 YARN 部署文档结合仓库内flink-yarn模块的源码实现完整讲解环境准备、三种部署模式的使用命令、YARN 专属配置项、资源分配行为、高可用机制以及用户 Jar 与 Classpath 的控制方式帮助你在真实 YARN 集群上快速跑通并稳定运维 Flink 作业。Getting Started在 YARN 上启动你的第一个 Flink 集群背景Flink 与 YARN 如何协作YARN 是众多数据处理框架普遍使用的资源提供者。Flink 服务被提交到 YARN 的 ResourceManager由它在一台台运行着 NodeManager 的机器上拉起容器ContainerFlink 将自身的 JobManager 与 TaskManager 实例部署进这些容器中。在此基础上Flink 会根据 JobManager 上运行作业所需的处理槽位slot数量动态地申请和释放 TaskManager 资源从而实现资源的弹性伸缩。准备环境本节假设你已有一个可用的 YARN 环境建议版本不低于 2.10.2。生产环境中最常见的方式是通过 Amazon EMR、Google Cloud DataProc 或 Cloudera 等产品直接获得托管 YARN 环境本地单机或自建集群的 YARN 搭建方式可参考 Hadoop 官方文档SingleCluster / ClusterSetup但不建议仅为了完成本教程而手动搭建。开始之前请完成两项检查运行yarn top确认 YARN 集群已准备好接收 Flink 应用输出不应出现错误信息。从 Flink 下载页获取一个最新的 Flink 发行包并解压。务必确认HADOOP_CLASSPATH环境变量已设置可通过echo $HADOOP_CLASSPATH检查。若未设置执行export HADOOP_CLASSPATHhadoop classpath这一步至关重要Flink on YARN 运行时依赖的 Hadoop 相关类需要从HADOOP_CLASSPATH中加载详见后文支持的 Hadoop 版本一节。启动 YARN Session 并提交示例作业确认HADOOP_CLASSPATH就绪后进入解压后的 Flink 发行目录依次执行# 假设当前位于解压后的 Flink 发行包根目录 # (0) 导出 HADOOP_CLASSPATH export HADOOP_CLASSPATHhadoop classpath # (1) 启动 YARN Session分离模式 ./bin/yarn-session.sh --detached # (2) 通过命令输出末尾打印的 URL 访问 Flink Web UI # 也可以从 YARN ResourceManager 的 Web UI 中找到入口 # (3) 提交示例作业 ./bin/flink run ./examples/streaming/TopSpeedWindowing.jar # (4) 停止 YARN Sessionapplication id 以 yarn-session.sh 命令输出为准 echo stop | ./bin/yarn-session.sh -id application_XXXXX_XXX完成以上步骤你就在 YARN 上成功运行了一个 Flink 应用。Flink on YARN 支持的部署模式生产环境建议优先使用 Application Mode 部署 Flink 应用因为它在应用之间提供了更好的隔离性。Application Mode应用模式Application Mode 会在 YARN 上启动一个 Flink 集群应用 Jar 中的main()方法直接在 YARN 上的 JobManager 中执行。集群会在应用结束后自动关闭你也可以通过yarn application -kill ApplicationId或取消 Flink 作业来手动停止集群。./bin/flink run-application -t yarn-application ./examples/streaming/TopSpeedWindowing.jarApplication Mode 集群部署完成后仍可与之交互执行取消、触发 savepoint 等操作# 列出集群上运行的作业 ./bin/flink list -t yarn-application -Dyarn.application.idapplication_XXXX_YY # 取消运行中的作业 ./bin/flink cancel -t yarn-application -Dyarn.application.idapplication_XXXX_YY jobId注意在 Application Cluster 上取消作业会同时停止整个集群。要充分发挥应用模式的优势可以结合yarn.provided.lib.dirs配置项将应用 Jar 预先上传到集群所有节点都能访问的位置如 HDFS此时提交命令形如./bin/flink run-application -t yarn-application \ -Dyarn.provided.lib.dirshdfs://myhdfs/my-remote-flink-dist-dir \ hdfs://myhdfs/jars/my-application.jar上述方式让作业提交变得非常轻量所需的 Flink Jar 与应用 Jar 直接从指定的远程位置获取而不是由客户端打包上传到集群。从源码看该配置在 YarnConfigOptions.java 中定义为分号分隔的 provided lib 目录列表要求预先上传且对集群内所有节点可读Flink 会据此跳过本地 Flink 发行文件如flink-dist、lib/、plugins/的上传以加速提交同时 YARN 会在节点上缓存这些文件避免每个应用重复下载。相关逻辑由 YarnApplicationFileUploader.java 与 Utils.java 处理。Session Mode会话模式Session Mode 已在本文开头的 Getting Started 中演示过它有以下两种运行形态attached 模式默认yarn-session.sh客户端将 Flink 集群提交到 YARN 后保持运行持续跟踪集群状态。若集群失败客户端会展示错误若客户端被终止也会通知集群随之关闭。detached 模式-d或--detachedyarn-session.sh客户端提交集群后立即返回后续需要再次调用客户端或使用 YARN 工具来停止集群。Session 模式会在/tmp/.yarn-properties-username下生成一个隐藏的 YARN properties 文件命令行接口在提交作业时会自动读取它来完成集群发现。你也可以在提交 Flink 作业时手动指定目标 YARN 集群./bin/flink run -t yarn-session \ -Dyarn.application.idapplication_XXXX_YY \ ./examples/streaming/TopSpeedWindowing.jar重新挂接到 YARN Session使用以下命令./bin/yarn-session.sh -id application_XXXX_YY除了通过 Flink 配置文件 传递配置外你还可以在提交时用-Dkeyvalue参数向./bin/yarn-session.sh客户端传入任意配置。./bin/yarn-session.sh -h可查看客户端内置的常用设置快捷参数。Per-Job Mode已废弃危险提示Per-job 模式仅由 YARN 支持且自 Flink 1.15 起已被弃用跟踪 issue 为 FLINK-26000未来版本将移除。如需在 YARN 上为每个作业启动独立集群请改用 Application Mode。Per-Job Cluster 模式会先在 YARN 上启动一个 Flink 集群然后在本地运行提供的应用 Jar最后将 JobGraph 提交给 YARN 上的 JobManager。若传入--detached提交被接受后客户端即退出。YARN 集群会在作业停止后随之关闭。./bin/flink run -t yarn-per-job --detached ./examples/streaming/TopSpeedWindowing.jarPer-Job 集群部署完成后同样支持交互操作# 列出集群上运行的作业 ./bin/flink list -t yarn-per-job -Dyarn.application.idapplication_XXXX_YY # 取消运行中的作业 ./bin/flink cancel -t yarn-per-job -Dyarn.application.idapplication_XXXX_YY jobId同样注意取消 Per-Job Cluster 上的作业会停止整个集群。Flink on YARN 配置参考配置方式与运行时托管参数YARN 专属的全部配置项统一列在配置页中。以下参数由 Flink on YARN 在运行时托管可能被框架动态覆盖请勿手动指定jobmanager.rpc.address由 Flink on YARN 动态设置为 JobManager 容器所在地址io.tmp.dirs若未设置Flink 会使用 YARN 定义的临时目录high-availability.cluster-id由 Flink 自动生成用于在 HA 服务中区分多个集群。如需向 Flink 传递额外的 Hadoop 配置文件可通过HADOOP_CONF_DIR环境变量指定一个包含 Hadoop 配置文件的目录。默认情况下所需 Hadoop 配置文件都通过HADOOP_CLASSPATH环境变量从 classpath 加载。YARN 核心配置项一览源码级在 YarnConfigOptions.java 中集中定义了全部yarn.*配置项以下是部署与运维最常用的一组配置项默认值说明yarn.containers.vcores-1每个 YARN 容器报告的虚拟核数。默认等于每个 TaskManager 配置的 slot 数未配置 slot 时为 1。需要 YARN 集群启用 CPU 调度如FairScheduler才生效yarn.appmaster.vcores1YARN Application Master即 JobManager 容器使用的虚拟核数yarn.application-attempts无默认 1HA 时为 2ApplicationMaster 重启次数同时受 YARN 的yarn.resourcemanager.am.max-attempts限制yarn.classpath.include-user-jarORDER用户 Jar 是否进入系统 classpath 及其位置取值DISABLED/FIRST/LAST/ORDERyarn.provided.lib.dirs无分号分隔的远程 provided lib 目录跳过本地 Flink 发行文件上传以加速提交YARN 会缓存到节点yarn.provided.usrlib.dir无远程 provided usrlib 目录用于排除本地usrlib/上传与yarn.provided.lib.dirs不同YARN 不会在节点上缓存它yarn.heartbeat.interval5秒Application Master 与 YARN ResourceManager 之间的心跳间隔yarn.heartbeat.container-request-interval500ms请求容器时的心跳间隔值越小容器分配通知越快但过度分配可能给 YARN 造成压力yarn.application.id无指定运行中 YARN 集群的 application id用于flink run/cancel/list等命令定位集群yarn.application.queue无作业提交到的 YARN 队列yarn.application.name/yarn.application.type无自定义 YARN 应用名称 / 类型yarn.application.priority-1提交优先级需 YARN 开启优先级调度-1 表示使用集群默认优先级yarn.properties-file.location无自定义.yarn-properties-username文件位置适用于多用户共享 Flink 安装的场景yarn.application-master.port0Application Master / JobManager 的 RPC 端口可指定端口、端口范围如50100-50200或列表0 表示由操作系统分配建议保持默认yarn.ship-files无分号分隔的需要随作业一起 ship 到 YARN 集群的文件/目录支持本地路径与 HDFS 路径yarn.ship-archives无分号分隔的需随作业 ship 的归档文件.tar.gz、.tar、.tgz、.dst、.jar、.zip会在本地化时解压yarn.staging-directory无提交应用时存放 YARN 文件的暂存目录默认使用所配置文件系统的 home 目录yarn.tags空逗号分隔的 YARN 应用标签yarn.file-replication-1每个本地资源文件在 HDFS/S3 上的副本数不配置则使用 Hadoop 配置的默认副本数适合容器数超百的加速场景yarn.application.node-label/yarn.taskmanager.node-label无YARN 节点标签后者可单独为 TaskManager 指定并覆盖前者yarn.container-start-command-template%java% %jvmmem% %jvmopts% %logging% %class% %args% %redirects%容器启动命令模板占位符分别对应 Java 路径、JVM 内存、JVM 参数、日志、主类、参数与输出重定向flink.hadoop.key/flink.yarn.key无通用探测配置去掉前缀后写入 Hadoop / YARN 配置。例如flink.hadoop.dfs.replication5会转换为 Hadoop 的dfs.replication5此外Kerberos 安全场景下可关注yarn.security.kerberos.ship-local-keytab默认true将 keytab 作为 YARN 本地资源 ship 出去与yarn.security.kerberos.localized-keytab-path默认krb5.keytab。资源分配行为运行在 YARN 上的 JobManager 在现有资源不足以运行全部已提交作业时会主动申请额外的 TaskManager。特别是在 Session Mode 下随着更多作业被提交JobManager 会按需分配更多 TaskManager闲置的 TaskManager 会在超时后被释放。JobManager 与 TaskManager 进程的内存配置会被 YARN 实现严格尊重。默认情况下容器上报的 VCores 数量等于每个 TaskManager 配置的 slot 数如需自定义可通过 yarn.containers.vcores 覆盖但该参数生效的前提是 YARN 集群已启用 CPU 调度。失败的容器包括 JobManager会由 YARN 自动替换。JobManager 容器的最大重启次数由 yarn.application-attempts 控制默认 1一旦所有尝试次数耗尽YARN 应用即宣告失败。从源码的 YarnConfigOptions.java 注释可以看到该值在 standalone 场景默认 1、在高可用场景默认 2且受到 YARN 侧yarn.resourcemanager.am.max-attempts的约束同时该参数返回 String 类型是因为 Integer 类型的配置项必须拥有静态默认值。在 YARN 上实现高可用High-AvailabilityYARN 上的高可用由 YARN 与一个高可用服务共同完成一旦配置了 HA 服务它会持久化 JobManager 元数据并执行 leader 选举YARN 则负责重启失败的 JobManager。JobManager 的最大重启次数由两个参数共同决定Flink 侧的yarn.application-attempts高可用场景下默认值为 2YARN 侧的yarn.resourcemanager.am.max-attempts默认值同样为 2是 Flink 该参数的上限。特别注意high-availability.cluster-id在 YARN 上部署时由 Flink 托管默认设置为 YARN application id用于在 HA 后端如 ZooKeeper中区分不同的 HA 集群。不要覆盖该参数否则多个 YARN 集群可能互相干扰。容器关闭行为按 YARN 版本2.3.0 version 2.4.0Application Master 失败时所有容器都会重启。2.4.0 version 2.6.0TaskManager 容器在 Application Master 失败后保持存活。好处是启动更快用户无需重新等待容器资源。2.6.0 version将 attempt failure validity interval尝试失败有效性窗口设置为 Flink 的 Pekko 超时值。该窗口的意思是只有在同一个时间窗口内达到最大应用尝试次数应用才会被杀死从而避免长运行作业过早耗尽尝试次数。危险提示Hadoop YARN 2.4.0 存在一个严重 bug已在 2.5.0 修复会阻止容器从重启的 Application Master / JobManager 容器中恢复详见 FLINK-4142。在 YARN 上搭建 HA 场景时建议至少使用 Hadoop 2.5.0。支持的 Hadoop 版本Flink on YARN 针对 Hadoop 2.10.2 编译所有 2.10.2的 Hadoop 版本包括 Hadoop 3.x均受支持。向 Flink 提供所需 Hadoop 依赖首选方式是在 Preparation 一节中介绍的HADOOP_CLASSPATH环境变量如果无法设置该变量也可以把依赖直接放入 Flink 发行包的lib/目录。此外Flink 官网的 Downloads / Additional Components 区域还提供预打包的 Hadoop fat jar 用于放入lib/目录这些 fat jar 经过 shading 处理以避免与常用库产生依赖冲突。需要注意Flink 社区并未针对这些预打包 jar 测试 YARN 集成。在防火墙环境中运行 Flink on YARN部分 YARN 集群会用防火墙控制集群内外网络流量。在这种环境中Flink 作业通常只能在集群网络内部防火墙之后提交到 YARN session。若生产上无法满足这一条件Flink 允许为 REST 端点用于客户端与集群通信配置端口范围从而支持跨防火墙提交作业。配置参数为 rest.bind-port它既接受单个端口如50010也接受范围50000-50025或两者组合。与之配套yarn.application-master.port也可用于为 Application Master 指定端口或端口范围以适配限制性防火墙环境。用户 Jar 与 ClasspathSession Mode在 YARN 上以 Session Mode 部署时只有启动命令中指定的那个 JAR 会被识别为用户 Jar 并纳入用户 classpath。Per-Job Mode 与 Application Mode以这两种模式部署时启动命令中指定的 JAR 与 Flinkusrlib目录下的所有 JAR 都会被识别为用户 Jar。默认情况下 Flink 将用户 Jar 放入系统 classpath该行为可用 yarn.classpath.include-user-jar 参数控制ORDER默认按字典序将 Jar 加入系统 classpathFIRST将 Jar 加到系统 classpath 开头LAST将 Jar 加到系统 classpath 末尾DISABLED不将用户 Jar 放入系统 classpath改为放入用户 classpath。从 YarnConfigOptions.UserJarInclusion 枚举的源码注释可以看出DISABLED意味着用户 Jar 会由用户 classloader 加载从而实现与系统类的隔离。更深入的类加载机制可参考 类加载调试文档。附源码级部署入口速览Flink on YARN 的部署逻辑集中在flink-yarn模块中可作为理解上述流程的源码锚点YarnClusterDescriptor.java集群部署的入口类其中 deploySessionCluster 负责部署 Session 集群deployApplicationCluster 负责部署 Application 集群并在内部校验deployment.target是否为yarn-applicationYarnConfigOptions.java上文全部yarn.*配置项的权威定义YarnApplicationFileUploader.java 与 Utils.java负责本地资源Jar、配置文件等的上传与 provided lib 目录的解析YarnResourceManagerDriver.java 与 YarnTaskExecutorRunner.java分别对应 YARN 侧资源管理驱动与 TaskExecutor 容器的启动入口。围绕上述部署与配置能力仓库中的flink-yarn模块还配套了完整的单元测试与端到端测试脚本如 test_ha_datastream.sh、test_resume_externalized_checkpoints.sh 等可用于验证 Session / Application 模式部署、HA 切换、savepoint 恢复等关键路径。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Apache Flink on YARN 部署实战Session、Per-Job、Application 三种模式与配置详解Apache Flink on YARN 部署实战Session、Per Job、Application 三种模式与配置详解 Flink 可以借助 Apach大数据流处理批处理数据工程Flink CDC on YARN 部署实战Session 与 Application 两种模式的完整落地流程Flink CDC on YARN 部署实战Session 与 Application 两种模式的完整落地流程 Flink CDC 生产环境中最常见的托管方式后端数据集成大数据流处理变更数据捕获数据同步在 YARN 上部署 Flink CDCSession 模式与 YARN Application 模式完整指南在 YARN 上部署 Flink CDCSession 模式与 YARN Application 模式完整指南 Flink CDC 是流式数据集成工具当需要后端数据集成大数据流处理变更数据捕获数据同步创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
02
RELATED NEWS

相关资讯

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

03
WHY YAOTU

想打造同款高转化官网?

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

场景化定制

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

营销型架构

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

全周期服务

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

免费获取你的建站方案

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