ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Flink StandAlone模式提交作业全流程:从集群搭建到实战排障

2026/9/15 13:35:40 拓冰建站 浏览量
Flink StandAlone模式提交作业全流程:从集群搭建到实战排障 做了这么多年 Flink 排查系列笔记整理到第 11 篇正好落到大家问得最多的 StandAlone 模式提交作业。很多人学 Flink 不是从 YARN 也不是从 K8S 开始的而是先在本地或一台测试机上把 StandAlone 集群拉起来然后手写一句flink run把任务交上去。这个流程看起来简单实际走一遍会发现里面有不少细节比如配置文件怎么写、端口绑到哪个地址、提交后任务为什么看不到、日志到底去哪里翻。这篇文章我就用一篇完整笔记把 StandAlone 模式提交 Flink 作业的全部流程拆开讲。内容包括集群搭建、核心配置、提交 JAR、跑 SQL 作业、Savepoint 管理以及我在真实环境里踩过的坑和排查思路。想快速搭一套本地验证环境的朋友或者刚接触 Flink、想搞清楚作业在集群里到底怎么流转的同学都能在里面找到可以直接抄作业的内容。1. 为什么先从 StandAlone 聊起1.1 什么是 StandAlone 模式StandAlone 是 Flink 官方提供的最基础部署模式之一。它的核心思想是集群里的组件只有 JobManager 和 TaskManager没有第三方资源调度层参与。TaskManager 负责执行具体的计算任务JobManager 负责接收作业、调度任务、协调 Checkpoint同时对外提供 Web UI 和 REST 接口。用生活里的场景打个比方YARN 模式相当于你租房资源是房东YARN分配的房子大小、租期都由房东说了算K8S 模式相当于住酒店容器由编排平台统一调度想要什么规格随时申请。而 StandAlone 模式相当于你自己买房、自己当房东硬件、内存、CPU 全部独占没有中间层来干预。也正因为这样StandAlone 的结构非常透明出了问题容易定位最适合理清 Flink 作业提交过程。从客户端提交作业的角度看StandAlone 模式下用户的flink run命令会先把用户代码编译成 JobGraph然后发给 JobManager。JobManager 拿到 JobGraph 后生成 ExecutionGraph再分发给可用的 TaskManager 去执行。这个链路里没有资源队列、没有动态分配所以只要你集群起来了提交作业的行为非常直观。1.2 什么时候该用 StandAlone我自己的判断标准是学习 Flink 原理、调试用户代码、做功能验证用 StandAlone 最合适。跑一个几十并发的小型实时任务且不需要很强的弹性伸缩StandAlone 完全能胜任。公司里已经有大数据集群但不想绑定特定 YARN 版本也不打算引入容器平台StandAlone 是一个低成本的独立运行方案。如果任务量大、需求弹性伸缩、多团队共享资源那就不推荐 StandAlone应该优先考虑 YARN 或 K8S。从成本和维护角度讲StandAlone 的缺点也很明显。TaskManager 的 Slot 是静态分配的你规划多少个就是多少个任务高峰期没办法自动扩容低峰期资源也不能自动释放。所以它更贴近“固定资源池”的思路。用人话讲StandAlone 适合“单机玩具”和“小规模生产”不适合“弹性平台”。1.3 StandAlone 模式下作业是怎么流转的要搞清楚完整流程必须先建立整体概念。一次 StandAlone 模式提交作业会经历下面几个阶段客户端执行flink run或通过 SQL Client 提交作业。客户端解析用户代码生成 JobGraph并把依赖 Jar 或 SQL 相关的 DAG 提交给 JobManager。JobManager 收到作业后将其转换为 ExecutionGraph根据作业要求的并行度和当前集群的 Slot 资源进行调度。JobManager 将任务部署到有可用 Slot 的 TaskManager 上。TaskManager 启动任务线程开始消费数据、执行计算、周期性地做 Checkpoint。作业运行过程中客户端可以断开因为作业已经托管在 JobManager 上重新连接 Web UI 或再次执行flink list就能看到状态。理解这条链路之后再去看配置文件、命令参数和报错日志思路会清晰很多。后面所有实操都是围绕这条链路展开的。2. 环境准备与集群搭建实操2.1 版本选择与基础环境Flink 版本差异会直接影响下载链接、配置项名称和命令用法。我推荐先想清楚自己主要写 Java 还是 SQL再选版本。以目前常用的 1.18 和 1.19 为例两个版本对 JDK 8、JDK 11、JDK 17 都有支持但如果你要用较新的 Flink CDC 组件建议用 1.18 以上版本接口兼容性更好。搭建 StandAlone 集群最基础的依赖如下Linux 或 macOS 系统Windows 也能跑但网络配置和脚本兼容性会让你多踩不少坑。JDK 8 或 11安装后确认java -version可用。不需要 Hadoop 环境除非你的作业要读写 HDFS。这是 StandAlone 和 YARN 模式的明显区别。本地准备好 Flink 安装包建议用 tar.gz 而不是源码包前者解压即用。下载 Flink 后解压到固定目录可以看到bin、conf、lib、log、plugins等目录。bin目录下有很多管理脚本conf目录下是核心配置文件lib目录放 Flink 自带依赖和需要补齐的连接器 Jar。这些结构在后续配置里都会用到。2.2 修改核心配置进入解压目录后最先要改的是conf/flink-conf.yaml。这个文件的格式是key: value容易理解但有几个配置项必须认真对待。第一项是jobmanager.rpc.address。默认是localhost如果你只是本机提交、本机跑保持不变可以但如果你是远程访问或者客户端机器和集群不在同一台机器必须改成 JobManager 所在机器的实际 IP 或主机名。很多“提交作业连不上 JobManager”的问题都是因为这个配置没改。第二项是rest.bind-address和rest.port。默认rest.bind-address是0.0.0.0Web UI 和 REST API 监听所有网卡一般不用动。rest.port默认 8081如果端口被占用需要手动换掉同时记住这个端口后面提交作业、查看 Web UI 都靠它。第三项是 TaskManager 的 Slot 数量。taskmanager.numberOfTaskSlots默认是 1建议根据机器 CPU 核数和内存改大一些比如 4 或 8。简单理解每个 Slot 相当于 TaskManager 能同时运行几个任务线程的资源配额。如果机器是 8 核 16G我一般会配 4 个 Slot既留出系统余量也方便并行度测试。第四项是内存参数。Flink 在配置内存时按组件拆分比如jobmanager.memory.process.size和taskmanager.memory.process.size。这两个是进程总内存写清楚点避免 JVM 被系统杀掉。我的经验值JobManager 给 1G 到 2G 就够TaskManager 需要根据任务复杂度给轻量任务 2G重状态任务建议 8G 以上但不能超过机器物理内存。除了flink-conf.yaml还可以关注conf/workers文件。这个文件每一行写一个 TaskManager 所在主机名或 IP。默认只写localhost表示只启动本机的 TaskManager。如果要多节点组成集群就在这个文件里逐行添加节点地址并且确保节点之间能通过 SSH 免密登录。2.3 启动集群并验证配置完成后启动命令非常简单。回到 Flink 根目录执行bin/start-cluster.sh这个脚本会按照workers文件里的节点清单在本地启动 JobManager 和对应的 TaskManager。启动完成后终端会输出进程启动信息。如果想分开启动也可以手动执行bin/jobmanager.sh start bin/taskmanager.sh start分开启动的好处是便于排查。比如 JobManager 启动失败不会影响 TaskManager 的启动过程日志也会分别写到对应目录。启动之后我习惯先做三件事验证集群状态。第一件事查看进程是否存在jps正常情况下应该看到至少两个进程其中包含StandaloneResourceManager、StandaloneDispatcher这些名字的进程都属于 JobManager带TaskManagerExecutor字样的就是 TaskManager。第二件事访问 Web UI。浏览器打开http://jobmanager-ip:8081如果页面正常显示并且 Task Managers 列表里有节点说明集群基本正常。这一步能同时验证 REST 绑定地址是否正确。第三件事通过命令行查看可用资源bin/flink list如果没有运行中的作业会显示“No running jobs”之类的内容同时输出了 JobManager 地址。出现这个结果说明客户端能正常连上 JobManager为后面提交作业扫清了最大障碍。2.4 一个容易被忽略的细节IP 与主机名StandAlone 模式远没有容器环境那么自动。很多情况下集群搭好了但客户端提交作业时反复失败最后发现是主机名解析出了问题。举个真实例子。我有一台测试机主机名是node01/etc/hosts里写的映射是127.0.0.1 node01然后我把jobmanager.rpc.address配成了node01。结果客户端在远程执行flink run时它去解析node01解析到了本机回环地址根本连不上 JobManager。处理方式有两种第一种是让/etc/hosts里把node01映射到真实内网 IP192.168.1.10 node01第二种是干脆把jobmanager.rpc.address配成 IP不依赖主机名。在我的实操经验里更推荐第二种配置越简单后续越不容易出问题。如果你要搭建多节点集群还需要在每台机器的/etc/hosts里写上其他节点的 IP 和主机名映射否则 TaskManager 注册不到 JobManager。3. 提交作业的完整流程详解3.1 JAR 作业flink run 的一生集群正常后最常见的提交方式是提交一个 JAR 包。假设你有一个打好的作业包主类名是com.example.MyFlinkJob提交命令可以这样写bin/flink run \ -m node01:8081 \ -c com.example.MyFlinkJob \ -p 4 \ ./examples/my-flink-job.jar \ --input /data/input.txt --output /data/result这里面的参数有必要逐一说明-m node01:8081指定 JobManager 的 REST 地址。如果没有这个参数Flink 会读取flink-conf.yaml里的rest.address配置。显式指定时优先级更高适合命令行临时切换环境。-c指定主类。如果 Jar 包的 MANIFEST 里已经写好了 Main-Class不写也能跑但很多代码打包不规范显式指定更稳妥。-p指定作业并行度。这个并行度会覆盖代码里设置的并行度但不会超过集群可用 Slot 数量。后面的--input和--output是传给用户主类的参数Flink 不解析它们会原样透传给用户代码。执行这条命令后终端会显示类似“Job has been submitted successfully with JobID xxx”的信息。此时作业已经交给 JobManager终端可以关掉不影响作业运行。至于命令行输出的 JobID是后面做flink cancel、查看日志、恢复 Savepoint 的关键标识建议复制保存下来。如果不小心丢了可以用bin/flink list再次查看。提交作业时还有一个细节值得注意flink run默认是同步等待作业结束还是提交后立刻返回答案取决于作业类型。对于流式作业命令提交后很快返回因为流作业一般会一直运行对于批作业命令可能会一直卡在那里直到批作业跑完才返回这属于正常现象。如果不想等批作业结束可以加-d参数让命令提交后立刻分离方便脚本化调用。3.2 SQL 作业sql-client 与 SQL Gateway很多朋友学了 Flink 之后第一感觉是 Java API 太重于是转向 SQL。StandAlone 集群同样可以跑 SQL 作业方式比 JAR 更简单。Flink 自带一个 SQL Client命令行启动bin/sql-client.sh embedded进入 SQL 交互模式后可以用标准 DDL 建表、用 DML 提交计算逻辑。比如本地有一个 CSV 文件想做一个简单统计CREATE TABLE source_table ( user_id STRING, event_time TIMESTAMP(3), page_id STRING, WATERMARK FOR event_time AS event_time - INTERVAL 5 SECOND ) WITH ( connector filesystem, path /data/events.csv, format csv ); CREATE TABLE sink_table ( page_id STRING, cnt BIGINT ) WITH ( connector print ); INSERT INTO sink_table SELECT page_id, COUNT(*) FROM source_table GROUP BY page_id;这里特意加了 WATERMARK 语句因为实时流场景里处理迟到的数据基本离不开事件时间和水位线。需要注意的是SQL Client 创建的会话默认使用planner和execution.runtime-mode等配置如果想以流模式跑要确保没有设置成批模式。SQL 作业提交后在 Web UI 的 Jobs 页面同样能看到一个正在运行的作业。它的 Job 名称通常是 SQL 语句里 INSERT 的目标表名或者是一个系统生成的字符串。你在 SQL Client 里按 Ctrl C 退出客户端不会影响已经提交的作业这个跟 JAR 作业是一样的。SQL Client 适合人工调试。如果想沉淀成服务接口让多个客户端并发提交 SQL可以用 SQL Gateway。启动命令是bin/sql-gateway.sh start启动后默认监听 8083 端口。SQL Gateway 的核心价值是把“SQL 会话”服务化不同客户端可以连上来复用 Session底层共享一个 Flink 集群不需要每个客户端都起一个 Flink 环境。直接连命令行执行 SQL 的话可以用bin/sql-client.sh gateway --gateway-address node01 --port 8083这种方式适合团队内统一提交入口权限控制、语法校验都可以前置处理。3.3 Savepoint 与作业停止作业跑起来之后另外一个绕不开的操作是停止作业。对于流作业直接 kill 进程会导致状态丢失下次恢复只能从最早的 Checkpoint 或者外部存储重新算。所以生产环境里更规范的做法是用flink stop或者flink cancel配合 Savepoint。先看怎么触发一个 Savepointbin/flink savepoint jobId file:///data/savepoints这条命令会让 JobManager 为指定作业生成一个一致性快照保存到本地或其他分布式文件系统。生成成功后命令行会输出 Savepoint 的路径。之后如果要停止作业并保留状态可以执行bin/flink cancel -s file:///data/savepoints/savepoint-xxxx jobId这里的-s表示停止前先生成 Savepoint然后再取消作业。想从指定 Savepoint 恢复作业时提交命令需要加-s参数bin/flink run \ -s file:///data/savepoints/savepoint-xxxx \ -c com.example.MyFlinkJob \ ./examples/my-flink-job.jar这个流程在状态算子比较重的作业里特别重要。举个例子一个窗口聚合作业统计用户最近一小时行为状态里可能缓存了大量中间结果没有 Savepoint 停机下次启动就要重算一小时数据代价非常大。有了 Savepoint恢复后能接着上次的状态继续算。我处理作业升级时标准顺序一般是先触发 Savepoint确认生成成功再停掉旧作业用新 JAR 包加上-s参数恢复启动。这样可以做到基本不丢状态、不重算历史。3.4 从 Web UI 确认作业状态提交作业后别急着关页面Web UI 上有几个关键信息值得盯一下。第一个是 Job Manager 首页右上角的“Jobs”列表。如果刚提交的作业出现在这里状态是 RUNNING基本说明调度成功。如果状态是 FAILED点进去看原因通常最有用的是 Exception 信息会明确告诉你错误发生在哪个算子、什么异常。第二个是 Task Managers 页面。这里能看到每个 TaskManager 的 Slot 使用情况。如果作业提交成功但一直处于 SCHEDULED 或 DEPLOYING先看 Slot 够不够。很多新手把并行度设成 8但集群只有 4 个 Slot作业会一直排队等资源很容易被误判成集群死锁。第三个是 Checkpoints 页面。流作业开启了 Checkpoint 后这个页面会展示最近一次 Checkpoint 是否成功、时延多少、状态大小等。如果 Checkpoint 持续失败通常意味着状态太大、存储不稳定或 Barrier 对齐超时这些内容后面排障会反复看到。4. 常见问题与排查技巧实录4.1 TaskManager 无法注册这是 StandAlone 模式里最常见的问题。现象是 JobManager 启动正常Web UI 能打开但 Task Managers 页面一直为空或者 TaskManager 启动后很快退出。排查思路按顺序走看 TaskManager 日志。日志文件在log/目录下文件名类似flink-user-taskexecutor-host.log。看有没有类似TaskManager did not connect within 30 seconds的报错。如果有重点检查taskmanager.hostname是否被解析成了回环地址。检查conf/workers文件里的节点是否都能被 SSH 访问以及各节点时间是否一致。时钟漂移会导致握手失败。检查防火墙是否放行了jobmanager.rpc.port默认端口 6123 和taskmanager.data.port默认 6121、6122。我在测试环境里遇到过最奇葩的一次是 TaskManager 因为日志目录磁盘满了启动到一半直接崩溃。所以如果 TaskManager 反复重启除了看日志也别忘记df -h一下。4.2 flink run 报连接错误提交作业时如果出现类似cannot connect to JobManager或Connection refused先检查-m参数指向的地址是不是当前 JobManager 的实际 REST 端口。有时候 Web UI 能打开但flink run还是连不上。原因是 Web UI 走的是rest.port而客户端提交作业走的是 RPC 端口jobmanager.rpc.port。两者是不同的通道。如果防火墙只放开了 8081RPC 无法通信就会出现“UI 能打开但提交失败”的诡异现象。快速定位方法是在客户端机器上执行telnet jobmanager-ip 6123如果连接失败基本就是网络隔离或防火墙问题。如果通了再去翻flink-conf.yaml确认jobmanager.rpc.address配置的是可达 IP。4.3 依赖冲突与 NoClassDefFoundError提交 JAR 作业时经常看到NoClassDefFoundError或LinkageError。这类问题多半是用户 Jar 里的依赖和 Flink 自带的依赖冲突或者缺少连接器。以 Flink CDC 为例很多同学直接把 flink-connector-mysql-cdc 的 Jar 放到用户作业包里提交结果报错找不到com.ververica.cdc.connectors.mysql.source.MySqlSource。原因通常是打包时没有将 CDC 依赖打进去或者打进去的版本和集群 Flink 版本不兼容。我推荐两种解决办法。第一种是下载对应 Flink 版本的连接器 Jar放到lib/目录下重启集群适用于 SQL 和 CDC 这类重量级连接器。第二种是在用户工程里把连接器依赖设为provided再用 Maven Shade 插件把其他第三方依赖打包进作业 Jar避免把 Flink 核心依赖也打进去。对于 SQL Client 提交的作业连接器 Jar 可以直接用-j参数指定bin/sql-client.sh embedded -j /path/to/flink-sql-connector-kafka-3.0-1.18.jar这样不用重启集群也能让 SQL Client 所在的会话加载到对应连接器。4.4 内存与 Slot 不足StandAlone 的槽位是固定的。如果你配置了 4 个 Slot但提交作业的并行度是 6作业会卡在等待资源的状态。很多新手以为这是配置错误反复重启集群实际上只需要把并行度调小或者给 TaskManager 增加 Slot。调整方式很简单taskmanager.numberOfTaskSlots: 8但这里有一个容易误解的点Slot 数目只是资源槽位不代表可以无限塞任务。每个 Slot 能跑多少个任务取决于任务的实际 CPU 和内存占用。如果 TaskManager 内存给得不大Slot 数量倒是配了 8结果部署 6 个高压任务时照样 OOM。真正常态是Slot 数量按 CPU 核数设定内存按单个任务预期乘以 Slot 数再乘以安全系数来估算。遇到 OOM日志里会有java.lang.OutOfMemoryError。这时先检查taskmanager.memory.process.size是否够再检查taskmanager.memory.managed.size是否给得过大。托管内存用不完的部分可以适当调低腾给堆内存或者反过来。具体要看你的作业是消耗 CPU 多、状态多还是网络数据量大。4.5 作业提交成功但并行度跟预期不符有时候作业明明提交成功了Web UI 里看到的并行度不对。或者并行度很高但实际数据延迟很大。这背后通常有两个原因。一是代码里指定了env.setParallelism(N)而命令行-p又设置了另一个值。不同优先级需要记清楚算子级别并行度最高环境级别并行度其次命令行-p再次最后是配置文件默认并行度。所以如果你在代码里给某个算子加了 2 并行度命令行-p 8也不会改变这个算子的并行度。二是 KeyBy 之后的数据倾斜。并行度没问题但某些 Key 数据量特别大导致个别子任务处理不过来。这种问题从并行度表面看不出来要去看每个 Subtask 的 Input/Output 指标。如果某个 Subtask 的积压数据明显高于其他就需要考虑改 Key 设计或者增加局部预聚合。5. 运维层面的几个补充建议5.1 日志别只盯一个文件很多同事在排查问题时只打开一个日志文件然后说“没有报错”。实际上 StandAlone 模式有两个层面的日志要分开看。JobManager 的日志记录的是集群管理事件比如作业提交、资源分配、Checkpoint 协调。TaskManager 的日志记录的才是任务执行细节比如业务里log.info、异常堆栈、以及每个算子的处理情况。用户代码里的输出基本只在 TaskManager 日志里出现。如果你的作业代码里打印了日志但找不到内容先确认日志对应的 TaskManager 是哪个再打开那个文件。Flink 默认每个 TaskManager 一个日志文件文件里会包含任务的名字和作业 ID可以直接搜索。5.2 Slot 规划与并行度既然 StandAlone 资源是静态的那么一开始规划好 Slot 和内存比事后频繁调参更省心。我给出一个朴素的经验公式。假设单机是 16 核 32G 内存给系统预留 4GFlink 进程预留 28G。如果目标并行度是 10那每个 Slot 大概能分到 2.8G。如果单个任务状态很大需要 8G 堆内存那 10 个 Slot 显然不现实资源规划时要以状态最大的任务为准。另外要注意多个作业共用一个 StandAlone 集群时Slot 资源是全局共享的。假设两个作业并行度都是 6但集群只有 8 个 Slot后提交的作业会一直等待。这个不算故障但会造成资源利用率忽高忽低。我的建议是如果一个集群要跑多个作业给每个作业预留好 Slot别全都塞满。5.3 HA 与多 JobManager如果 StandAlone 集群需要高可用就不能只启动一个 JobManager。Flink 原生支持 StandAlone HA思路是让多个 JobManager 节点通过 ZooKeeper 做 Leader 选举。主要配置项包括high-availability: zookeeper high-availability.zookeeper.quorum: node01:2181,node02:2181,node03:2181 high-availability.storageDir: file:///data/flink/ha/这里要注意high-availability.storageDir需要所有 JobManager 节点都能访问本地路径在多机场景下不合适生产环境一般用 HDFS 或共享存储。配了 HA 之后客户端提交作业时如果连接的是备用节点Flink 会自动重定向到当前 Leader。这套配置在测试环境稍微复杂但对线上稳定性很重要。我自己的项目经验是测试环境没必要上 HA开发调试时单 JobManager 足够但一旦作业要 7x24 小时运行至少要考虑 HA 和定期备份 Savepoint否则一次进程挂掉就需要手动恢复。最后再分享一个习惯。每次启动集群或提交作业前我会固定执行一次flink list确认客户端和 JobManager 的通信链路是通的再继续往下操作。这个动作看起来多余但能帮你把“集群问题”和“作业问题”快速切分开。代码写得再稳环境不稳也白搭希望这篇 StandAlone 提交作业的完整流程能帮你少走几步弯路。