Flink on YARN:大数据集群部署最佳实践
Flink on YARN:大数据集群部署最佳实践
关键词:Flink、YARN、资源管理、流处理、集群部署、高可用性、动态扩缩容
摘要:本文从大数据集群的实际需求出发,结合Flink流处理框架与YARN资源管理系统的特性,系统讲解Flink on YARN的部署原理、最佳实践与调优技巧。通过生活类比、代码示例和实战案例,帮助读者理解Flink与YARN的协作机制,掌握从环境搭建到故障排查的全流程操作,最终实现高效、稳定的大数据实时处理集群部署。
背景介绍
目的和范围
在大数据领域,实时处理需求日益增长。Flink作为Apache顶级流处理框架,以低延迟、高吞吐和精准一次处理语义著称;而YARN(Yet Another Resource Negotiator)作为Hadoop生态的资源管理核心,负责集群资源的统一分配与调度。两者结合(Flink on YARN)能充分利用YARN的资源管理能力,实现Flink任务的动态扩缩容、高可用和多租户隔离。本文覆盖从基础概念到实战部署的全流程,适用于中大型企业的生产环境部署场景。
预期读者
- 大数据工程师(负责Flink任务开发与集群运维)
- 集群管理员(负责YARN资源调度与集群稳定性)
- 技术经理(需要理解Flink on YARN的技术价值与成本)
文档结构概述
本文先通过“食堂打饭”的生活案例引出核心概念,再拆解Flink与YARN的协作原理;接着通过实战步骤演示部署过程,包括环境搭建、配置调优和故障排查;最后结合实际场景总结最佳实践,展望未来趋势。
术语表
核心术语定义
- Flink JobManager(JM):Flink的“指挥官”,负责任务调度、资源申请和故障恢复(类比食堂经理)。
- Flink TaskManager(TM):Flink的“执行者”,运行具体的数据流处理任务(类比打饭窗口的工作人员)。
- YARN ResourceManager(RM):YARN的“总调度”,管理集群所有资源(类比食堂的“后勤主任”,分配餐桌和餐具)。
- YARN NodeManager(NM):YARN的“现场管理员”,负责单个节点的资源监控与容器管理(类比食堂各楼层的值班员)。
- Container(容器):YARN分配的资源单元(CPU、内存),是TaskManager/JM的运行载体(类比食堂的“专用打饭窗口”)。
缩略词列表
- YARN:Yet Another Resource Negotiator(Hadoop资源管理器)
- JM:JobManager(Flink主节点)
- TM:TaskManager(Flink工作节点)
- RM:ResourceManager(YARN主节点)
- NM:NodeManager(YARN工作节点)
核心概念与联系
故事引入:用“食堂打饭”理解Flink on YARN
假设你是一个学校食堂的“实时打饭系统”设计者:
- 学生(数据):源源不断涌入食堂,需要快速打饭(处理)。
- 打饭窗口(TaskManager):每个窗口一次能服务多个学生(并行处理),但需要占用餐桌和餐具(CPU/内存资源)。
- 食堂经理(JobManager):负责规划开多少窗口、分配窗口位置,并监控窗口是否故障(任务调度)。
- 后勤主任(YARN RM):掌管食堂所有餐桌和餐具(集群资源),决定给“实时打饭系统”分配多少资源(Container)。
- 楼层值班员(YARN NM):在每层楼检查窗口是否正常使用资源,防止某个窗口占用过多餐具(资源隔离)。
Flink on YARN的本质,就是“食堂经理”(JM)向“后勤主任”(RM)申请“专用打饭窗口”(Container),然后“打饭窗口”(TM)在“楼层值班员”(NM)的监管下,高效处理“学生”(数据)。
核心概念解释(像给小学生讲故事一样)
核心概念一:YARN——集群资源的“大管家”
YARN就像学校的“后勤大管家”,它掌管着整个集群的“资源仓库”(CPU、内存、磁盘)。无论你是运行Flink、Spark还是Hive任务,都需要向YARN申请“资源包”(Container)。每个Container是一个“资源盒子”,比如包含4核CPU和8GB内存,YARN会根据任务需求分配这些盒子,并监控盒子是否被合理使用。
核心概念二:Flink JobManager(JM)——任务的“指挥官”
Flink任务启动时,首先需要一个“指挥官”(JM)。它的工作是:
- 规划任务需要多少“执行者”(TM);
- 向YARN“大管家”(RM)申请对应的“资源盒子”(Container);
- 监控每个“执行者”的状态,一旦某个执行者罢工(故障),就重新申请资源恢复任务。
核心概念三:Flink TaskManager(TM)——数据的“处理员”
TM是真正干活的“处理员”。每个TM运行在YARN分配的Container里,它的任务是:
- 从数据源(比如Kafka)接收数据;
- 按照用户编写的Flink程序(比如实时统计订单量)处理数据;
- 将结果输出到存储系统(比如HBase、ClickHouse)。
核心概念之间的关系(用小学生能理解的比喻)
YARN与Flink JM的关系:资源申请与分配
JM就像“项目负责人”,需要向YARN(后勤大管家)申请“办公场地”(Container)。负责人(JM)会说:“我需要3间办公室,每间4张桌子(4核CPU)和80GB文件柜(8GB内存)。” 后勤大管家(YARN RM)检查仓库(集群资源),如果有足够资源,就分配对应的办公室(Container),并通知各楼层值班员(NM)准备场地。
YARN NM与Flink TM的关系:资源使用的“监管员”
当TM(办公人员)进入分配的办公室(Container)后,楼层值班员(NM)会盯着他们:
- 不能偷偷占用额外的桌子(超过分配的CPU);
- 不能把文件柜塞得太满(超过分配的内存);
- 如果办公人员罢工(TM进程崩溃),值班员会通知后勤大管家(RM),大管家再告诉项目负责人(JM)重新申请资源。
Flink JM与TM的关系:指挥官与执行者
JM(指挥官)给每个TM(执行者)分配具体任务:“TM1负责处理北京的订单数据,TM2负责上海的订单数据。” 如果TM1处理太慢,指挥官(JM)会调整任务分配(比如给TM1多分配一个处理线程);如果TM1挂了,指挥官会立即向YARN申请新的Container,启动新的TM继续工作。
核心概念原理和架构的文本示意图
Flink on YARN的协作流程可总结为“三步曲”:
- 启动JM:用户提交Flink任务时,YARN先为JM分配一个Container(类似“项目负责人办公室”)。
- JM申请TM资源:JM启动后,根据任务并行度(需要多少执行者)向YARN RM申请多个TM Container。
- TM运行任务:YARN NM在各个节点启动TM进程,TM连接JM后开始处理数据。
Mermaid 流程图
核心算法原理 & 具体操作步骤
Flink on YARN的核心是资源动态分配机制。Flink JM会根据任务的并行度(parallelism)和每个TM的任务槽数量(taskmanager.numberOfTaskSlots),计算需要的TM数量。例如:
- 任务并行度=8(需要8个任务槽同时工作);
- 每个TM有4个任务槽(
taskmanager.numberOfTaskSlots=4); - 则需要的TM数量=8/4=2个。
JM会向YARN申请2个TM Container,每个Container的资源由taskmanager.memory.process.size(TM总内存)和taskmanager.cpu.cores(TM CPU核数)决定。
关键配置参数(以Flink 1.17为例)
| 参数名 | 含义 | 默认值 | 推荐值(生产环境) |
|---|---|---|---|
yarn.application.name | YARN中显示的应用名称(方便监控) | Flink job | 业务名称(如“实时订单统计”) |
yarn.containers.max | JM可申请的最大TM Container数量(防止资源耗尽) | 无限制 | 根据集群资源调整(如20) |
taskmanager.numberOfTaskSlots | 每个TM的任务槽数量(决定并行度上限) | 1 | CPU核数(如4核设4) |
taskmanager.memory.process.size | TM进程总内存(包括JVM堆内存、堆外内存) | 1600MB | 16GB(根据任务复杂度调整) |
jobmanager.memory.process.size | JM进程总内存(太小会导致调度失败,太大浪费资源) | 1600MB | 4GB(生产环境至少2GB) |
yarn.resourcemanager.address | YARN RM的地址(必须与Hadoop集群一致) | localhost | 集群RM节点IP:8032 |
数学模型和公式 & 详细讲解 & 举例说明
TM数量计算模型
TM数量计算公式:
T
M
n
u
m
=
⌈
T
o
t
a
l
s
l
o
t
s
T
M
s
l
o
t
s
⌉
TM_{num} = \lceil \frac{Total_{slots}}{TM_{slots}} \rceil
TMnum=⌈TMslotsTotalslots⌉
其中:
-
T
o
t
a
l
s
l
o
t
s
Total_{slots}
Totalslots:任务总并行度(用户设置的
parallelism); -
T
M
s
l
o
t
s
TM_{slots}
TMslots:每个TM的任务槽数(
taskmanager.numberOfTaskSlots)。
举例:任务并行度=10,每个TM有3个任务槽,则需要TM数量= ⌈ 10 / 3 ⌉ = 4 \lceil 10/3 \rceil=4 ⌈10/3⌉=4(因为3×3=9不够,需要4个TM提供12个槽)。
内存分配模型
每个TM Container的内存由taskmanager.memory.process.size决定,需满足:
T
M
m
e
m
o
r
y
=
J
V
M
h
e
a
p
+
J
V
M
o
f
f
h
e
a
p
+
N
a
t
i
v
e
m
e
m
o
r
y
TM_{memory} = JVM_{heap} + JVM_{off_heap} + Native_{memory}
TMmemory=JVMheap+JVMoffheap+Nativememory
其中:
-
J
V
M
h
e
a
p
JVM_{heap}
JVMheap:JVM堆内存(用于存储对象,通过
taskmanager.memory.heap.size设置); -
J
V
M
o
f
f
h
e
a
p
JVM_{off_heap}
JVMoffheap:JVM堆外内存(用于网络缓冲区、 RocksDB状态存储等,通过
taskmanager.memory.off-heap.size设置); - N a t i v e m e m o r y Native_{memory} Nativememory:本地内存(用于进程本身,如线程栈、元数据等)。
举例:若taskmanager.memory.process.size=16GB,可分配:
- 堆内存=8GB(
taskmanager.memory.heap.size=8192m); - 堆外内存=4GB(
taskmanager.memory.off-heap.size=4096m); - 本地内存=4GB(自动计算)。
项目实战:代码实际案例和详细解释说明
开发环境搭建
前置条件
- Hadoop集群(Hadoop 3.3.6+)已部署,YARN服务正常(
yarn --version验证); - Flink二进制包(flink-1.17.1-bin-scala_2.12.tgz)下载到集群节点;
- Java 11+环境(
java -version验证)。
步骤1:配置Flink与YARN的集成
修改flink-conf.yaml(Flink配置文件),添加YARN相关配置:
# YARN配置
yarn.application.name: realtime_order_analysis # 应用名称
yarn.resourcemanager.address: hadoop-rm:8032 # YARN RM的地址(hadoop-rm是RM节点的主机名)
yarn.provided.lib.dirs: hdfs:///flink/lib # Flink依赖库的HDFS路径(提前上传flink-shaded-hadoop-3-uber-...等JAR包)
# JM配置
jobmanager.memory.process.size: 4096m # JM内存4GB
jobmanager.rpc.port: 6123 # JM RPC端口(默认即可)
# TM配置
taskmanager.memory.process.size: 16384m # TM内存16GB
taskmanager.numberOfTaskSlots: 4 # 每个TM 4个任务槽(假设节点是8核,留4核给其他进程)
taskmanager.cpu.cores: 4 # TM使用4核CPU
# 高可用配置(ZooKeeper)
high-availability: zookeeper
high-availability.zookeeper.quorum: hadoop-zk1:2181,hadoop-zk2:2181,hadoop-zk3:2181
high-availability.storageDir: hdfs:///flink/ha # JM元数据存储在HDFS
步骤2:上传Flink依赖到HDFS(关键优化)
为避免每次提交任务都上传Flink依赖(耗时),可将Flink的opt目录下的JAR包上传到HDFS:
hadoop fs -mkdir -p /flink/lib
hadoop fs -put /path/to/flink-1.17.1/opt/* /flink/lib
步骤3:提交Flink任务到YARN
使用Flink的yarn-session模式(长期运行的会话)或per-job模式(每个任务独立)。这里以per-job模式为例(更推荐,资源隔离更好):
./bin/flink run \
-t yarn-per-job \
-D yarn.application.name=realtime_order_analysis \
-D jobmanager.memory.process.size=4096m \
-D taskmanager.memory.process.size=16384m \
-D taskmanager.numberOfTaskSlots=4 \
/path/to/your-flink-job.jar
源代码详细实现和代码解读
Flink任务的主类通常继承StreamExecutionEnvironment,示例代码(统计每分钟订单量):
public class OrderCountJob {
public static void main(String[] args) throws Exception {
// 初始化流执行环境(自动检测YARN环境)
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 从Kafka读取订单数据流
DataStream<Order> orderStream = env.addSource(
KafkaSource.<Order>builder()
.setBootstrapServers("kafka-broker:9092")
.setTopics("order_topic")
.setGroupId("order_count_group")
.setDeserializer(new OrderDeserializer())
.build()
);
// 按窗口统计每分钟订单量
DataStream<OrderCount> countStream = orderStream
.keyBy(Order::getShopId) // 按店铺分组
.window(TumblingProcessingTimeWindows.of(Time.minutes(1))) // 1分钟滚动窗口
.process(new OrderCountProcessFunction()); // 自定义处理函数
// 输出到ClickHouse
countStream.addSink(
ClickHouseSink.<OrderCount>builder()
.setUrl("jdbc:clickhouse://clickhouse-node:8123")
.setTableName("order_count")
.setUsername("user")
.setPassword("password")
.build()
);
// 执行任务(提交到YARN)
env.execute("Realtime Order Count");
}
}
代码解读:
StreamExecutionEnvironment.getExecutionEnvironment():自动检测运行环境(本地或YARN),提交到YARN时会触发JM和TM的资源申请。KafkaSource:从Kafka读取实时订单数据(数据源)。keyBy+window:按店铺分组,每分钟统计一次订单量(核心计算逻辑)。ClickHouseSink:将结果写入ClickHouse(数据输出)。env.execute():触发任务执行,Flink会自动与YARN交互申请资源。
实际应用场景
场景1:电商实时销售大屏
某电商双11期间需要实时展示各品类的销售额、订单量。通过Flink on YARN部署实时计算任务,YARN可根据流量动态扩缩容TM数量(比如流量峰值时自动增加TM Container),确保延迟低于1秒。
场景2:日志实时监控
某银行需要监控服务器日志,实时检测异常请求(如5xx错误)。Flink on YARN的高可用特性(通过ZooKeeper实现JM故障转移)保证了任务7×24小时运行,YARN的资源隔离防止了日志任务与其他任务(如离线ETL)的资源竞争。
场景3:物联网设备数据聚合
某工厂的传感器每秒钟产生10万条数据,需要实时计算设备温度、转速的平均值。Flink on YARN的动态资源分配(taskmanager.numberOfTaskSlots灵活调整)支持高并发数据处理,YARN的容器化管理确保了不同设备的数据任务互不干扰。
工具和资源推荐
监控工具
- Flink Web UI:默认端口8081,查看任务拓扑、并行度、延迟等指标(
http://jm-node:8081)。 - Prometheus+Grafana:通过Flink的Prometheus reporter导出指标(如
numRecordsInPerSecond、checkpointDuration),用Grafana可视化(推荐Flink官方Dashboard)。 - YARN ResourceManager Web UI:默认端口8088,查看Container使用情况、队列资源占用(
http://rm-node:8088)。
调试工具
yarn logs -applicationId <appId>:查看YARN应用的日志(包括JM和TM的日志)。flink yarn info <appId>:查看Flink on YARN应用的状态(如Running/Pending/Failed)。jmap/jstack:用于分析TM/JM的内存溢出或线程阻塞(需登录节点执行)。
文档资源
- Flink官方YARN文档
- Hadoop YARN官方文档
- 《Flink基础与实践》(作者:张利兵)—— 第5章详细讲解Flink资源管理。
未来发展趋势与挑战
趋势1:更细粒度的资源分配
Flink正在推进“细粒度资源管理”(Fine-Grained Resource Management),未来TM的任务槽(Task Slot)可以独立申请YARN Container,避免资源浪费(例如:一个任务槽需要2GB内存,无需为整个TM分配16GB)。
趋势2:与YARN的深度集成优化
Flink社区计划优化与YARN的心跳机制,减少JM与RM的通信延迟;同时支持YARN的“抢占式调度”,当高优先级任务需要资源时,自动释放低优先级任务的Container。
挑战1:多租户资源隔离
大型集群中多个团队共享YARN资源,需通过YARN的容量调度器(Capacity Scheduler)或公平调度器(Fair Scheduler)配置队列权重,防止某团队的Flink任务占用过多资源(推荐设置yarn.scheduler.capacity.maximum-am-resource-percent=0.3,限制所有JM的总资源不超过集群的30%)。
挑战2:故障排查复杂性
Flink on YARN的故障可能由Flink本身(如Checkpoint失败)、YARN(如Container被OOM Kill)或底层节点(如磁盘故障)引起。需要结合Flink日志(log/flink-*.log)、YARN日志(yarn logs)和节点监控(如Prometheus的node_exporter)综合分析。
总结:学到了什么?
核心概念回顾
- YARN:集群资源的“大管家”,负责分配Container(资源盒子)。
- Flink JM:任务的“指挥官”,向YARN申请资源并调度任务。
- Flink TM:数据的“处理员”,运行在YARN的Container中,处理具体数据流。
概念关系回顾
- YARN为Flink提供资源载体(Container),Flink通过JM与YARN RM交互申请资源。
- TM在YARN NM的监管下运行,确保资源使用不超限。
- JM监控TM状态,故障时通过YARN重新申请资源,保证任务高可用。
思考题:动动小脑筋
-
假设你的集群有10台节点,每台节点32核、64GB内存,YARN预留10%资源给系统进程。如果要部署一个Flink任务,并行度=40,每个TM需要4个任务槽,你会如何配置
taskmanager.numberOfTaskSlots和taskmanager.memory.process.size?(提示:计算每台节点可运行的TM数量,避免资源过载) -
如果Flink任务的延迟突然升高,你会从哪些方面排查?(提示:YARN资源使用率、TM的任务槽负载、Kafka消费速率、ClickHouse写入速度)
-
如何利用YARN的队列功能,让生产环境的Flink任务(高优先级)优先于测试任务获取资源?(提示:参考YARN Capacity Scheduler的队列配置)
附录:常见问题与解答
Q1:提交Flink任务到YARN时,提示“Container launch failed”怎么办?
A:常见原因:
- YARN资源不足(查看YARN RM UI的可用资源);
- TM内存超过节点可用内存(检查
taskmanager.memory.process.size是否超过节点内存×90%); - Flink依赖未正确上传到HDFS(重新执行
hadoop fs -put上传opt目录下的JAR包)。
Q2:Flink任务运行一段时间后,TM频繁重启,YARN日志显示“Container killed by YARN for exceeding memory limits”
A:TM内存溢出。解决方法:
- 增大
taskmanager.memory.process.size; - 调整堆外内存比例(如增加
taskmanager.memory.off-heap.size,减少堆内存); - 检查Flink任务是否有内存泄漏(通过
jmap -dump:format=b,file=heap.bin <pid>导出堆快照,用MAT分析)。
Q3:Flink JM故障后,如何快速恢复?
A:启用ZooKeeper高可用(high-availability: zookeeper),JM故障时,YARN会自动启动新的JM,从HDFS(high-availability.storageDir)加载元数据,恢复任务状态(需确保Checkpoint正常)。
扩展阅读 & 参考资料
- Flink官方YARN部署指南
- Hadoop YARN Capacity Scheduler配置
- 《Hadoop权威指南(第4版)》—— 第8章YARN核心原理。
- Flink细粒度资源管理提案
更多推荐


所有评论(0)