做实时计算这行,跟Flink打交道是绕不开的。如果你刚接触Flink,或者正在纠结怎么把手里的作业跑起来,我建议你从StandAlone模式入手。很多人在学习阶段就直接上YARN或者K8s,结果被资源管理、容器调度这些概念绕晕了,反而把核心的作业提交逻辑给忽略了。StandAlone模式是Flink自带的独立部署方式,不需要依赖任何外部资源调度框架,一台机器或者几台机器就能搭出一个完整的集群。你可以用它来学习、测试,也可以在小规模生产环境里跑一些轻量级的实时任务。
这篇文章会从零开始,把StandAlone模式下提交Flink作业的完整流程走一遍。内容包含集群的架构设计、环境准备、配置文件调整、作业打包、命令行提交、Web UI提交、SQL Client提交,以及我在实际运维中踩过的坑和排查思路。不论你是刚入门的菜鸟,还是已经接触过Flink但没系统梳理过提交流程的同学,这篇文章都能给你一个可以直接照做的参考。
1. StandAlone模式整体架构与设计思路
1.1 为什么先选StandAlone模式
我见过不少刚接触Flink的开发者,一上来就直奔YARN或者Kubernetes,觉得那是生产环境的标准答案。这个想法本身没错,但问题在于:当你对Flink的作业生命周期、任务调度、检查点机制这些都还没有直观概念的时候,再叠加上一套外部资源管理框架,排查问题的时候你会分不清到底是Flink的问题还是资源框架的问题。
StandAlone模式最大的价值,是把Flink本身的功能边界画得清清楚楚。它只用Flink自带的组件来完成资源管理与任务调度,不需要额外的依赖。你在一台笔记本上就能启动一个最小集群,然后观察一个作业从提交、调度、执行到完成的全过程。对于学习来说,这比任何文档都直观。
这套模式在小规模生产环境里也有用武之地。比如你有一个数据量不大的实时ETL任务,每天处理几百GB数据,用两三台机器搭一个StandAlone集群完全够用。它的运维成本很低,不需要维护额外的资源调度服务,出问题的时候定位路径也比较短。
1.2 核心组件与角色划分
StandAlone集群由两大角色构成:JobManager和TaskManager。JobManager是集群的“大脑”,负责接收作业、调度任务、协调检查点、处理故障恢复;TaskManager是“工人”,负责真正执行算子逻辑、缓存数据、与上下游交互。
在Flink 1.11之后的版本里,JobManager内部进一步拆分成了Dispatcher、ResourceManager和JobMaster三个逻辑组件。Dispatcher负责接收用户提交的作业,并提供REST接口供Web UI和命令行调用;ResourceManager负责管理TaskManager的槽位资源,作业需要多少槽位它就协调多少;JobMaster则具体管理单个作业的整个生命周期。
你可以这样理解:Dispatcher是前台接待,ResourceManager是后勤调度,JobMaster是项目负责人,TaskManager是一线干活的人。这个类比虽然粗糙,但能帮助你快速建立整体印象:提交作业到Dispatcher,Dispatcher通知ResourceManager分配槽位,JobMaster负责把作业调度到具体的TaskManager上执行,过程中任何一个环节出问题,日志都会指向对应的组件。
1.3 StandAlone与YARN、K8s模式选型对比
选哪种部署模式,本质上是在回答一个问题:谁来管理计算资源?
StandAlone模式下,Flink自己管理资源。集群启动后TaskManager的槽位是固定好的,作业提交后只能在这些槽位里调度。优点是部署简单、概念清晰;缺点是资源利用率低,集群闲的时候浪费,忙的时候可能不够。YARN模式把Flink作业放在Hadoop的YARN集群里,由YARN动态分配容器,作业结束后资源自动释放,适合跟Hadoop生态深度绑定的场景。K8s模式则是通过容器编排平台来管理Flink集群和作业,弹性最好,适合云原生环境,但学习和运维成本也是最高的。
| 对比项 | StandAlone | YARN | Kubernetes |
|---|---|---|---|
| 部署复杂度 | 低 | 中 | 高 |
| 资源动态分配 | 不支持 | 支持 | 支持 |
| 学习门槛 | 低 | 中 | 高 |
| 适合场景 | 学习、测试、小规模生产 | Hadoop生态生产环境 | 云原生大规模生产 |
| 故障恢复 | 依赖Flink自身 | YARN重新拉起 | K8s重新调度 |
我当时第一次在生产环境用Flink,就是从StandAlone开始的。三台机器,每台跑一个TaskManager,跑了一个实时指标计算任务,一年下来几乎没有因为集群本身出过问题。它不像大家说的那么“玩具”,关键是你要把资源规划好,并且接受它没有动态伸缩这个事实。
2. 环境准备与集群启动
2.1 环境要求与JDK版本选择
StandAlone模式对硬件的要求不高,最低配置一台2核4G的机器就能跑起来。但如果你要跑真实的计算任务,我建议最少三台机器,一台跑JobManager,两台跑TaskManager,避免单点故障导致作业直接挂掉。机器操作系统选Linux最省心,CentOS 7或Ubuntu 20.04都行,Windows上也能跑,但遇到网络和权限问题的概率会高不少。
JDK版本是第一个容易踩坑的地方。Flink 1.13及之前的版本对JDK 8支持得最好,Flink 1.14开始支持JDK 11,Flink 1.18开始支持JDK 17。我的建议是:如果你用的是Flink 1.17及以下版本,老老实实用JDK 8;如果你用Flink 1.18以上版本,JDK 8和JDK 11都可以,但一定要保持所有节点JDK版本一致。公司里有同事曾经因为JobManager用的JDK 8、TaskManager用的JDK 11,导致作业提交后ClassNotFoundException,排查了一下午。
确保所有节点之间网络互通,hostname能互相解析。最简单的方式是配置/etc/hosts,把JobManager和TaskManager的主机名和IP映射都写进去。这个步骤看似琐碎,但漏掉它,后面TaskManager注册不上来的时候你就会很痛苦。
2.2 flink-conf.yaml核心配置项解析
Flink的配置文件在conf目录下,核心是flink-conf.yaml。这个文件决定了集群的资源行为和运行参数。下面这几个配置项,是你启动StandAlone集群之前必须搞清楚的。
# JobManager进程总内存,包含堆内、堆外、JVM自身开销 jobmanager.memory.process.size: 1600m # TaskManager进程总内存 taskmanager.memory.process.size: 2048m # 每个TaskManager可提供的任务槽位数量 taskmanager.numberOfTaskSlots: 4 # 默认并行度 parallelism.default: 1 # 故障恢复策略 jobmanager.execution.failover-strategy: regionjobmanager.memory.process.size和taskmanager.memory.process.size控制的是进程级别的总内存,Flink会在内部进一步拆分成堆内、托管内存、网络缓冲等。新手最容易犯的错误是只给TaskManager设置很小的内存,比如512M,结果作业一跑起来就频繁Full GC,性能惨不忍睹。对于普通的流式作业,我建议TaskManager至少给2G,如果你有窗口聚合、状态存储这类需求,4G起步。
taskmanager.numberOfTaskSlots决定了一个TaskManager能跑多少个并行子任务。槽位数量不等于物理核数,你可以把槽位理解成“执行位”。一个4核8G的机器,配置4个槽位是合理的选择。槽位配多了,每个任务分到的资源变少,反而容易出问题。如果你的作业并行度是2,两个TaskManager各4个槽位,那么Flink会把作业的两个并行实例尽量分散到不同的TaskManager上。
注意:parallelism.default只是默认值,作业提交时通过-p参数指定的并行度优先级更高。别以为配置文件里写死就万事大吉。
2.3 启动集群与Web UI验证
配置完成后,在JobManager节点上执行启动脚本:
# 进入Flink安装目录 cd /opt/flink # 启动StandAlone集群 bin/start-cluster.sh如果你之前修改过workers文件(在conf目录下,每行一个TaskManager主机名),脚本会自动在所有Worker节点上启动TaskManager进程。第一次启动时,我建议你把每个节点的日志都检查一遍。JobManager的日志在logs/flink-{user}-standalonesession-{id}-{host}.log,TaskManager的日志在logs/flink-{user}-taskexecutor-{id}-{host}.log。
启动后,打开浏览器访问http://jobmanager:8081(如果你在本地,就是http://localhost:8081),你会看到Flink的Web UI。重点看两个地方:左侧的Task Managers页面,确认所有TaskManager都成功注册,并显示你配置的槽位数量;Overview页面的Task Slots字段,确认总的可用槽位数符合预期。
如果TaskManager没有出现在Web UI上,最常见的两个原因:一是hostname解析不了,二是TaskManager的RPC端口(默认6122)被防火墙拦截。第一个问题去检查/etc/hosts,第二个问题检查安全组和firewalld规则。你还可以在JobManager节点上手动执行下面的命令,确认Worker到JobManager的网络是否通:
# 在TaskManager节点上执行,JobManagerHost换成实际主机名 telnet JobManagerHost 6123端口通了,再去看日志,别一上来就怀疑配置。集群起不来时,九成是网络,剩下才是配置。
3. 作业打包与提交方式详解
3.1 作业代码的核心思路
先写一段最简的流式作业作为演示。这个作业从Socket读取文本,按行做词频统计,输出到标准控制台。它不涉及外部存储,目的是让你把注意力完全放在提交流程上。
import org.apache.flink.api.common.functions.FlatMapFunction; import org.apache.flink.api.java.tuple.Tuple2; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.util.Collector; public class SocketWordCount { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); DataStream<String> text = env.socketTextStream("localhost", 9999); text.flatMap(new FlatMapFunction<String, Tuple2<String, Integer>>() { @Override public void flatMap(String line, Collector<Tuple2<String, Integer>> out) { for (String word : line.split("\\s+")) { out.collect(Tuple2.of(word, 1)); } } }) .keyBy(value -> value.f0) .sum(1) .print(); env.execute("Socket Word Count"); } }这里有个细节很多人会忽略:env.execute("Socket Word Count")这一行才是真正把作业提交给集群的入口。如果注释掉这行,程序会在本地执行算子链但不会触发作业的实际提交。在实际工程中,你的代码可能非常复杂,但提交的入口永远是这一句。另外,socketTextStream这种数据源只适合本地测试,生产环境请使用Kafka、Pulsar这类消息队列作为Source,否则作业重启后数据源就断了。
3.2 Maven打包与依赖处理
打包是整个提交流程里最能体现工程化水平的一步。很多人的作业在IDE里跑得好好的,打成Jar提交到集群就报ClassNotFoundException,十有八九是依赖没有打进去。
这句话请记住:Flink的核心依赖是集群自带的,你的作业Jar里不要包含flink-streaming-java、flink-core这些Flink核心库,否则会跟集群自身带的版本冲突。但第三方依赖,比如连接器、JSON解析库、HTTP Client,则必须打进Jar里,或者放到Flink的lib目录下。
推荐使用maven-shade-plugin来做打包,它会自动把依赖合并成一个Fat Jar,并且可以处理依赖冲突。
<build> <plugins> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-shade-plugin</artifactId> <version>3.4.1</version> <executions> <execution> <phase>package</phase> <goals> <goal>shade</goal> </goals> <configuration> <transformers> <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer"> <mainClass>com.example.SocketWordCount</mainClass> </transformer> </transformers> <filters> <filter> <artifact>*:*</artifact> <excludes> <exclude>META-INF/*.SF</exclude> <exclude>META-INF/*.DSA</exclude> <exclude>META-INF/*.RSA</exclude> </excludes> </filter> </filters> </configuration> </execution> </executions> </plugin> </plugins> </build>顺便说一句,排除META-INF下的签名文件(SF、DSA、RSA)这个细节,是很多人在打包环节踩坑的重灾区。如果不去排除,某些依赖带签名信息,Jar提交后可能出现SecurityException,提示Invalid signature file digest。
3.3 命令行提交与参数解析
集群启动、作业打包完成之后,就到了关键的提交环节。先启动一个Socket数据源用于测试:
# 在作业节点的终端上 nc -lk 9999然后执行提交命令:
# 在Flink安装目录下 bin/flink run \ -m jobmanager-host:8081 \ -c com.example.SocketWordCount \ -p 2 \ -d \ /path/to/your-job.jar一个一个参数来说。
-m指定JobManager的地址和Web UI端口,flink命令会把作业提交给JobManager的Dispatcher。如果不加这个参数,flink命令会读取conf/flink-conf.yaml里的rest.address配置,在本地模式下则默认提交到localhost:8081。
-c指定包含了main方法的入口类。当你的Jar里有多个main方法时,这个参数是必须的。虽然maven-shade-plugin在MANIFEST里写入了Main-Class,但显式指定更稳妥,也能避免别人拿错Jar包的时候一头雾水。
-p是并行度,这里指定2意味着作业的每个算子会启动两个并行实例。并行度设置的原则是:不超过集群的总槽位数。比如你有两个TaskManager各4个槽位,总槽位是8,那么作业并行度最大可以设到8。超过的话作业会一直处于调度等待状态,直到有槽位空出来。
-d是detached模式,提交后命令行立即返回,作业在后台运行。如果不加这个参数,命令行会一直跟着作业的运行状态,作业结束命令才返回。实际生产环境推荐加-d,把作业的运行交给集群管理。
提交成功后,命令行会返回一个Job ID,同时提示作业已经提交到集群。你可以在Web UI的Job Manager页面看到作业由RUNNING到FINISHED的状态变化。
如果你想用REST API来做提交,Flink也提供了相应的接口。实际工作中,很多自动化平台就是通过POST /jar/upload上传Jar,再通过POST /jars/{id}/run来触发作业运行。这个方式比命令行更适合集成到调度系统里,但细节比较多,这里不展开。
3.4 Web UI提交与SQL Client提交
除了命令行,Web UI也提供了直观的提交入口。在Web UI页面右侧,有一个“Submit New Job”按钮,点击后可以上传Jar包,填写Main Class、并行度、Program Arguments等参数,然后点击Submit提交。这种方式适合临时测试场景,或者当你想快速观察作业在集群上的资源分配情况时使用。缺点是当你需要频繁提交、停止作业时,手动点页面效率太低,建议还是走命令行或者REST API。
如果你的作业是基于Flink SQL开发的,那就要用到SQL Client或者SQL Gateway。SQL Client是Flink自带的交互式SQL终端,启动方式如下:
bin/sql-client.sh进入SQL Client之后,你可以用标准的SQL语法来创建Source、Sink和计算逻辑,比如:
CREATE TABLE source_table ( user_id STRING, page_id STRING, view_time TIMESTAMP(3), WATERMARK FOR view_time AS view_time - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'user_views', 'properties.bootstrap.servers' = 'localhost:9092', 'properties.group.id' = 'flink-sql-group', 'format' = 'json', 'scan.startup.mode' = 'latest-offset' ); CREATE TABLE sink_table ( user_id STRING, cnt BIGINT ) WITH ( 'connector' = 'print' ); INSERT INTO sink_table SELECT user_id, COUNT(*) AS cnt FROM source_table GROUP BY user_id;SQL Client的好处是开发效率高,不需要写Java代码就能完成流式ETL;缺点是不太适合复杂的业务逻辑和需要精细调优的场景。如果你想把SQL作业做成工程化、自动化的交付方式,建议用SQL Gateway把SQL查询暴露成REST接口,让上游调度系统直接调用。这两年Flink社区在SQL Gateway上投入很大,如果你所在团队有多个业务线都要用Flink做实时计算,SQL Gateway是个值得关注的方向。
我在实际项目中还遇到过数据血缘的审计需求。用SQL方式开发时,一张目标表的数据来自哪些源表、经过了哪些计算,在SQL里天然是能追溯的;但如果用DataStream API写,这项能力就要自己实现。所以如果你们公司对数据治理有要求,SQL优先是一个合理的选择。
4. 常见问题与排查技巧实录
4.1 资源不足与并行度设置
作业卡在“ResourceManager请求Slot超时”,这是我在群里被问得最多的问题。通常有两种情况:一是你设置的并行度超过了集群可用槽位总数,二是TaskManager内存分配不合理,导致TaskManager启动后反复崩溃。
先说第一种。假设集群有两个TaskManager,每个4个槽位,总共8个。你提交作业时并行度设成10,那么有两个并行任务永远等不到槽位,作业会一直处于SCHEDULED状态。解决办法很简单:要么把并行度降下来,要么在集群里增加TaskManager或槽位。
第二种情况更隐蔽。TaskManager配置了2G内存,但机器本身只有4G,还同时跑了JobManager、Namenode、Kafka等一堆服务,结果TaskManager进程一启动就被操作系统OOM杀掉。此时你去日志里看不到Flink报错,只看到进程消失了。排查的时候一定要先看机器当前的可用内存。
free -h如果内存确实比较紧张,可以临时回收一些不需要的服务,或者把TaskManager内存调低。但调低的时候要注意:Flink分配给堆外内存(managed memory、network memory)是固定比例的,你不能把总内存压得太小,否则作业一跑就OOM。
4.2 依赖冲突与ClassNotFound
ClassNotFoundException是Flink作业提交后最经典的报错。我总结下来主要分两类。
第一类:作业代码里用了某个第三方库,但打包时没打进去,或者打进去了但版本不对。解决办法是检查Jar包。你可以用jar tf命令查看Jar内容,确认对应类是否存在,或者用jar -tf target/your-job.jar | grep 关键字来快速定位。
第二类:Flink核心库的版本冲突。比如你本地用Flink 1.17写的代码,但集群是1.15,提交后可能出现NoSuchMethodError。这种问题的根源在于Flink的核心类在集群端,你的代码编译期用的版本和集群运行期加载的版本不一致。解决办法是保证开发环境和集群版本一致,并在pom里把Flink依赖的scope设置为provided。
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java</artifactId> <version>${flink.version}</version> <scope>provided</scope> </dependency>还有一种情况是Jar包里出现了多个版本的同一个类。maven-shade-plugin在默认情况下不做依赖去重,你需要通过relocation把有冲突的包重命名掉。这个技巧在处理guava、netty这类重依赖冲突时特别管用,但配置比较复杂,新手可以先避开,等真遇到了再去研究。
4.3 网络通信与端口问题
StandAlone模式下,JobManager和TaskManager之间的通信走的是Akka和Netty,涉及多个端口。默认情况下,JobManager的RPC端口是6123,TaskManager的RPC端口是6122,Blob服务端口是6124,Web UI端口是8081。TaskManager启动时还会向JobManager注册,注册用的就是RPC端口。
我在实操时遇到过TaskManager日志一直报“Failed to connect to JobManager”的情况。检查下来发现是TaskManager节点上的防火墙把6123端口拦了。StandAlone模式下你不用每个端口都精确放行,Flink支持配置端口范围,比如:
taskmanager.rpc.port: 6122-6132这样TaskManager会在这个范围内选择可用的端口,方便统一配置安全组。这个技巧在节点数量较多的时候特别有用,省得一台一台去开放端口。
4.4 任务运行中失败与恢复
作业提交成功、运行一段时间后突然失败,是另一种让人头疼的情况。如果是StandAlone模式下JobManager进程没有挂,只是作业失败,优先看JobManager日志和TaskManager日志。日志里如果出现“Checkpoint failed”字样,说明检查点持久化有问题。StandAlone模式默认的检查点存储路径是JobManager的本地文件系统,一旦JobManager重启,检查点就丢了。生产环境建议把检查点配置到HDFS或者对象存储上:
state.backend: hashmap state.checkpoints.dir: hdfs:///flink/checkpoints state.savepoints.dir: hdfs:///flink/savepoints execution.checkpointing.interval: 60s execution.checkpointing.exactly-once: true这四行配置的含义分别是指定状态后端、检查点存储路径、Savepoint路径、检查点间隔和语义。其中exactly-once语义是Flink的核心卖点,它保证故障恢复后数据不丢不重(在特定条件下)。配置完之后,如果作业再次失败,你可以在Web UI的Checkpoints页面看到最近一次Checkpoint的状态,失败原因也会列出来,排查起来会直观很多。
5. 实操心得与优化建议
5.1 我在实际项目中积累的经验
做了一段时间Flink之后,我有个体会:提交方式本身不复杂,复杂的是把整个流程想清楚。很多人在本地敲一行flink run就把作业提交上去了,但没有想过这个作业需要多少个槽位、状态后端放在哪里、失败之后怎么恢复、日志怎么采集。
我个人比较推荐的做法是这样的:把提交命令固化成脚本或者接入调度平台,在脚本里预设好并行度、内存参数、检查点路径、保存点路径,最重要的是把Job ID记录下来。为什么要记录Job ID?因为后续你想对作业做Stop、Cancel、从Savepoint恢复,都需要这个ID。别等到要操作的时候再去Web UI里翻,那效率太低了。
还有一个容易被忽视的点:提交作业的客户端机器必须要能访问到JobManager的REST端口,但不需要能直接访问TaskManager。因为作业提交完成后,Jar包是传到JobManager的Blob服务上,再由JobManager分发给各个TaskManager的。你在远程客户端上提交时,确保网络策略允许客户端访问8081和6124端口就行。
5.2 从StandAlone到更复杂的部署形态
如果你把StandAlone模式的流程完全搞懂了,接下来可以尝试两个方向。第一个方向是把Flink作业容器化,用Docker把JobManager和TaskManager包起来,再结合容器编排平台做资源调度。第二个方向是提升作业的可观测性,接入Prometheus监控指标、配置Flink的Metrics Reporter,同时在作业里埋点记录处理延迟和数据量。
我为什么建议先跑通StandAlone再说?因为Flink的作业提交、状态管理、检查点恢复这些核心概念,跟部署模式无关。你在StandAlone模式下把这些概念吃透了,换到任何部署模式都只是换个环境变量和启动脚本的事。
最后分享一个我自己一直在用的小技巧:维护一个标准的作业部署清单。每次提交新作业之前,按照清单检查一遍——并行度是否匹配集群槽位、检查点路径是否可写、日志级别是否合理、Jar包里的依赖是否干净。这个清单看起来简单,但它帮我避免了至少五次线上事故。做实时计算的,最怕的就是作业跑着跑着因为资源或者依赖问题挂了,而你连它为什么会挂都不知道。把这些基本功做扎实,比追求花哨的架构要实在得多。