ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Flink资源配置优先级实测:5类来源与部署模式对比

Flink资源配置优先级实测:5类来源与部署模式对比 1. 先聊聊这次验证的起因先说个事儿。上个月我们组的一个同步作业莫名其妙卡在pending状态整整一个小时没有分配到任何slot。当时集群里同时跑着十几个Flink任务有高配置的实时数仓作业也有低并行度的数据上报任务。我们怀疑是资源不够但看监控明明还有空闲TaskManager群里一度吵到要重启集群。后来查了半天才发现不是资源不够是一个配置项没有按我预期的优先级生效。我本想用命令行参数覆盖集群配置文件里的内存设置结果实际生效的是另一个来源的配置导致作业申请的内存和集群实际提供的slot不匹配就一直等着。这个经历让我意识到Flink资源配置的“优先级”虽然是个老话题但在1.19.3这种新版本上很多人包括我都没有真正较真过。所以这篇博文不聊架构也不聊源码只干一件事把Flink 1.19.3里各种资源配置来源的优先级实测一遍记录每种来源在什么情况下生效、什么情况下不生效以及我在验证过程中踩到的坑。本文实测环境主要是Standalone和YARN Application两种模式覆盖了日常开发和生产部署最常用的两条路径后面会附上K8s和Spring Boot整合场景的注意事项。如果你正在被“明明改了内存配置却没生效”“同一套配置在不同环境表现不一致”“任务卡在pending不知道什么原因”这类问题困扰这篇报告应该能帮你省掉不少排查时间。1.1 引发排查的一个“中间优先级”事故上面说的那个卡住的任务我先还原一下现场。作业是MySQL到ClickHouse的同步任务用的是Flink CDC JDBC Sink并行度设的是8TaskManager内存通过flink run -D taskmanager.memory.process.size4096m手动指定。按理说4096m内存、8个并行度跑个同步作业很轻松。但实际提交之后作业执行图里显示每个算子的并行度是8可TaskManager的slot却怎么都不够用。为什么因为集群配置文件flink-conf.yaml里写的是taskmanager.numberOfTaskSlots2而作业申请时把每个TM的slot当成了可用的计算单元来估算。我本意是想让作业并行度跑8但每个TM只有2个slot内存又是4G相当于每个slot要吃2G内存却只能放2个并发最后3个TM全被占满作业还是差2个slot一直pending。真正让我困惑的不是pending本身而是我明明加了-D参数为什么实际生效的配置好像不是我想的那样。后来把TM日志调出来才发现-D参数确实生效了但它只改了内存没有改slot数而slot数读的是flink-conf.yaml的值。这种“多个配置来源混在一起生效”的情况就是典型的多级优先级问题——中间层的任务最容易被卡住因为它既不像低并行度任务那样容易挤进现有资源又不像高并行度任务那样有大内存加持资源请求卡在一个尴尬的区间。1.2 配置优先级到底是什么意思这里我要强调一个容易混淆的点Flink的“资源配置优先级”不是说Flink内部有一套类似YARN队列的优先级调度机制有高优先级任务先跑、低优先级任务后跑这种概念。这里说的优先级是指在配置合并过程中当同一个配置项在多个位置都有值时最终以哪个为准。用Java位运算符来类比会更好理解Java里的优先级高于^^又高于|如果表达式里同时出现这三个运算符又不加括号结果经常和你想的不一样。Flink的配置合并顺序也是类似的固定层级不是你写在后面的配置就覆盖前面的而是Flink按照自己定义的顺序把各个来源的配置合并成一个Configuration对象。搞清楚这个合并顺序比背一堆配置项要实用得多。1.3 验证范围和约定本次验证覆盖以下配置来源flink-conf.yaml 集群配置文件flink run 命令行专用参数-p、-ys、-ytm、-yjm 等-D 动态参数SQL Client 的 SET 命令REST API 提交时传入的 flinkConfiguration验证的配置项主要聚焦在资源相关进程总内存taskmanager.memory.process.size、堆内存taskmanager.memory.heap.size、slot数taskmanager.numberOfTaskSlots、默认并行度parallelism.default以及两个和资源强相关的连接器参数JDBC Sink 的批量提交参数和 CDC 的并发读取参数。2. 验证环境与前置准备2.1 软件版本与环境拓扑先说环境方便大家对照。我用的是一台4核16G的测试机装了Flink 1.19.3跑的是Standalone Session模式验证大部分用例另外用一台YARN集群CDH 7.3专门验证Application模式下的优先级差异。JDK用的1.8.0_202Flink官方推荐Java 8/11/17但生产上很多团队还在用JDK8我特意用JDK8跑了一遍确保结论在JDK8场景下也成立。验证用的作业是自己打的包一个很小的DataStream程序包含两个算子map和keyBysum这样能通过Web UI直观看到每个算子的并行度。另外准备了一个WordCount的SQL脚本用来验证SQL Client场景下的优先级行为。2.2 验证作业和脚本作业本身要满足两个要求一是能稳定跑起来不涉及外部依赖二是能通过每5秒输出一次的Running状态观察到资源配置。我的选择是直接用Flink自带的examples目录里的TopSpeedWindowing.jar这个作业足够干净也不会干扰资源判断。SQL脚本则是为了验证SQL Client的SET行为内容很简单SET parallelism.default 4; SET pipeline.name priority_test_job; CREATE TABLE source_table ( user_id BIGINT, event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time ) WITH ( connector datagen, rows-per-second 100 ); CREATE TABLE sink_table ( user_id BIGINT, cnt BIGINT ) WITH ( connector print ); INSERT INTO sink_table SELECT user_id, COUNT(*) FROM source_table GROUP BY user_id;这里用了datagen连接器不需要任何外部依赖跑起来就能看到作业并行度后面讲SQL Client优先级时会用到这个脚本。2.3 怎么判断配置是否真的生效这是整个验证的关键——不能靠“感觉”来判断配置生效了没有。我用了三条路径交叉确认第一条是Web UI。Flink的Web UI在“Task Managers”页面会显示每个TM的内存、slot数在“Job Details”页面会显示每个算子的并行度。这是最直观的路径适合快速判断。第二条是日志。Flink在启动时会打印一份生效的内存配置摘要包括JVM堆大小、托管内存、直接内存等。Standalone模式下看TaskExecutor日志YARN模式下看ApplicationMaster日志里的TaskExecutor启动日志能看到类似这样的输出Memory configured: - process memory: 4096 MB - JVM heap: 2684 MB - managed memory: 2048 MB第三条是REST API。Flink的监控接口可以直接拉取作业的执行图信息用curl就能拿到每个算子的并行度适合写脚本批量验证curl -s http://localhost:8081/jobs/overview | python3 -m json.tool curl -s http://localhost:8081/jobs/{job_id} | grep -E parallelism|taskmanager三条路径互相印证如果Web UI和日志显示的值一致那基本可以确定实际生效的配置就是它。3. 核心验证五类配置来源优先级实测3.1 第一轮flink-conf.yaml 是基准盘第一轮验证最简单先把flink-conf.yaml里的资源项改成明确的值然后不带任何额外参数直接提交作业看生效配置是否等于flink-conf.yaml。我的flink-conf.yaml关键配置如下taskmanager.memory.process.size: 2048m taskmanager.numberOfTaskSlots: 2 parallelism.default: 1 jobmanager.memory.process.size: 1600m然后执行flink run -d ./examples/streaming/TopSpeedWindowing.jar提交后看Web UITaskManager内存显示2048mslot数2作业并行度1与配置文件完全一致。这个结论看似废话但它是后续所有验证的基准。这里特别提醒一点Parallelism.default在SQL作业里和DataStream作业里的行为略有不同SQL作业如果用了GROUP BY这种需要shuffle的算子最终并行度会跟随默认并行度但某些Source算子会根据并行度合并所以判断并行度时要看作业执行图而不是看源码里写的算子个数。3.2 第二轮CLI专用参数确实能压过配置文件轮到我最初踩坑的场景了。同样还是那份flink-conf.yaml我改用命令行专用参数提交flink run -d -p 4 -D taskmanager.memory.process.size4096m ./examples/streaming/TopSpeedWindowing.jar这里同时用了两个来源-p是CLI专用参数-D是动态参数。实测结果是作业并行度4TM内存4096m但slot数仍然是2。这说明两条信息是分开生效的——-p覆盖了parallelism.default-D覆盖了taskmanager.memory.process.size而flink-conf.yaml里的taskmanager.numberOfTaskSlots没有被覆盖所以保持2。换一个验证如果我只用-p 4不用-D指定内存TM内存就会回到2048m。这证明了专用参数和-D参数各自覆盖自己对应的配置项并不会“连带”覆盖其他资源项。接下来要回答一个关键问题-p和-D parallelism.default同时出现时谁生效我执行了一次对照flink run -d -p 6 -D parallelism.default2 ./examples/streaming/TopSpeedWindowing.jar结果作业并行度是6不是2。也就是说CLI专用参数的处理时机在-D之后-p会覆盖掉-D设置的parallelism.default。这一点如果写反了很容易出现“我明明在命令行设置了并行度跑起来却不是这个数”的怪现象。3.3 第三轮-D动态参数是不是真的“最高优先级”既然专用参数能覆盖-D那-D是不是就没有胜出的场景了不是。关键要看配置项的类型。CLI专用参数只覆盖固定的几个字段-p对应并行度-ys对应slot数-ytm对应TM内存等而-D可以指定任意配置项。所以只要某个配置项没有对应的CLI专用参数-D就是最优先的那个来源。比如我要调整pipeline.name、execution.checkpointing.interval这类没有专用参数的配置只能通过-D或配置文件控制此时-D的优先级最高。还有一个容易踩的坑在SQL Client里SET命令和启动参数同时出现时的优先级。我用上面的SQL脚本验证了一下sql-client.sh -D parallelism.default2 -f test.sql而test.sql里写着SET parallelism.default 4;结果作业最终并行度是4。说明SQL脚本里的SET是在运行时执行的会覆盖启动时通过-D设置的默认并行度。这个行为和DataStream CLI的“专用参数覆盖-D”正好相反因为SET命令本身就是一个运行时配置修改操作优先级自然更高。所以严格来说我不能简单说“-D就是最高优先级”更准确的表述是在DataStream CLI场景下-D会先于专用参数被合并因此专用参数会覆盖-D中的对应项在SQL Client场景下脚本中的SET会覆盖启动时的-D参数。这两条结论在不同提交入口之间是矛盾的但都符合Flink的实现逻辑。写脚本的朋友如果没注意到这个差异很容易把并行度配拧。3.4 第四轮REST API提交的配置优先级Standalone Session模式下除了命令行还可以通过REST API提交JAR。Flink的REST API在提交作业时允许在请求体里传flinkConfiguration我验证了这个字段的优先级curl -X POST http://localhost:8081/jars/{jar_id}/run \ -H Content-Type: application/json \ -d { entryClass: org.apache.flink.streaming.examples.windowing.TopSpeedWindowing, parallelism: 8, flinkConfiguration: { taskmanager.memory.process.size: 4096m, parallelism.default: 6 } }实测结果是作业并行度8TM内存2048m。注意这里出现了一个表面矛盾的组合并行度8但parallelism.default传了6而TM内存没有变化。原因是REST API的body里parallelism这个字段是作业级并行度它优先于flinkConfiguration里的parallelism.default而taskmanager.memory.process.size虽然传了但Session集群的TM资源已经由集群启动时决定作业配置无法改变运行中TM的内存。这也是Session模式的一个核心特点资源维度上Session集群的JM/TM内存和slot数由集群启动时决定作业只能影响并行度、状态后端、checkpoint等作业参数改内存配置是无效的。所以如果你在Session模式下遇到了“改了TM内存不生效”不是优先级顺序错了而是这个配置项在当前模式下根本不受作业控制。3.5 第五轮YARN Application模式才是资源优先级的“主战场”Standalone Session测完我在YARN Application模式下重跑了同样的配置组合。这个模式下的行为差异很大因为整个集群是按作业启动的所有资源相关配置都会在集群启动阶段生效。flink run-application -t yarn-application \ -D taskmanager.memory.process.size4096m \ -D taskmanager.numberOfTaskSlots4 \ -D parallelism.default8 \ -d ./examples/streaming/TopSpeedWindowing.jar此时TM内存4096m、slot数4、并行度8全都生效。这是因为Application模式下conf里的flink-conf.yaml是客户端本地的而命令行-D会覆盖它然后这些配置直接用于启动整个集群。这个结论很重要如果你要在生产环境精准控制资源优先用Application模式加-D参数提交逻辑最清晰不会被Session集群的固定资源干扰。但注意-D也有陷阱如果同时设置了taskmanager.memory.heap.size和taskmanager.memory.process.size且进程总内存太小Flink会报配置冲突错误因为堆内存占比超出了进程总内存的合理范围。这个我后面在坑位章节详细说。4. 部署模式对优先级的影响4.1 Standalone Session集群资源定死作业只能“戴着镣铐跳舞”Standalone模式是很多人本地学习和开发用的模式它的优先级逻辑有一个显著特点作业级配置和集群级配置是两个层面。启动集群时读取的flink-conf.yaml决定了JM和TM的进程资源集群起来之后无论你提交作业时怎么-D、怎么设并行度TM的内存、数量和slot数都不会因为作业而改变。那作业级配置能影响什么主要是并行度、算子链、slot共享组、checkpoint间隔、状态后端等。如果你的作业并行度超过了集群总slot数作业就会持续pending直到其他作业释放slot或者你去手动扩容集群。我们一开始遇到的那个“中间优先级的任务无法运行”本质上就是这种场景——作业请求的并行度大于集群可用的slot数但又不是大到离谱所以一直吊在那里。4.2 YARN Application资源跟着作业走优先级最干净YARN Application模式是生产环境里我最推荐的一种提交方式。它把Flink集群当作一个YARN应用来启动每个作业有自己独立的JM/TM资源配置完全由作业提交时的-D参数决定。这个模式下优先级顺序最清晰命令行专用参数-ys、-ytm、-yjm -D动态参数 flink-conf.yaml 默认值我在CDH集群上用下面这组配置验证过TM内存4096m、slot数4、并行度8全部按预期生效flink run-application -t yarn-application \ -D yarn.application.namemysql_to_clickhouse \ -D taskmanager.memory.process.size4096m \ -D taskmanager.numberOfTaskSlots4 \ -D parallelism.default8 \ -d /data/flink/jobs/mysql_sync.jar这也是这个模式被称为“资源优先级最干净”的原因没有Session集群的固定资源干扰所有配置来源都在同一个作业生命周期内合并。代价是启动慢每次提交都要拉起一个新的YARN集群从提交到作业Running通常是40到90秒实时性要求高的场景要提前评估。4.3 K8s与Spring Boot整合场景的差异K8s模式和YARN模式类似也分Session和Application两种资源优先级逻辑基本一致。但多了一个要注意的点K8s模式下Pod规格是通过flink-conf.yaml或-D生成而且JVM参数对容器资源敏感如果你设置了taskmanager.memory.process.size但没保留足够的JVM OverheadPod运行一段时间后会被K8s的OOMKilled。Spring Boot整合Flink的场景比较特殊。很多团队用Spring Boot起一个服务在代码里用StreamExecutionEnvironment.createLocalEnvironment()或者提交到远程集群。代码里通过setParallelism()、setTaskManagerMemory()等方式设置的配置优先级是最高的因为它直接写入当前Environment的Configuration对象比任何外部文件都优先。但要注意这仅限于代码内配置如果你用configuration.setString(taskmanager.memory.process.size, 4096m)设了配置然后提交到Session集群TM内存依旧不生效——因为Session的TM资源不由作业决定我在3.4节说过这个逻辑。4.4 一张总表把优先级理清楚把上面所有验证结果汇总成一张表方便大家直接参考。配置来源DataStream CLISQL ClientREST API生效范围flink-conf.yaml基准配置基准配置基准配置集群级作业级命令行专用参数-p/-ys/-ytm等高于-D对应项不适用不适用覆盖对应配置项-D动态参数高除专用参数对应项启动阶段生效低于SET高作业级作业代码/Environment最高如设置不适用不适用作业级SQL脚本SET不适用运行时最高不适用作业级REST请求体并行度不适用不适用最高作业级再看两张按部署模式区分的资源生效表Standalone SessionTM/JM内存由集群启动时flink-conf.yaml决定作业配置无效slot数由集群启动时决定作业配置无效并行度由提交参数、-D、代码共同决定状态后端/checkpoint由提交参数或代码决定YARN/K8s ApplicationTM/JM内存由提交时-D或本机flink-conf.yaml决定slot数由提交时-D或本机flink-conf.yaml决定并行度由提交参数、-D、代码共同决定TM数量由并行度/slot数自动推导也可用-D taskmanager.numberOfTaskSlots显式控制5. 实践中的三个坑与排查实录5.1 案例一JDBC连接器异常锅却在资源优先级这是我在生产环境排查过的一个典型问题。作业是Flink JDBC Sink写入MySQL运行半小时后开始频繁报Communications link failure偶尔还有Aborted connection告警。一开始以为是MySQL的连接数不够把max_connections调到2000结果问题依旧。后来看TaskManager日志才发现JDBC连接器的实际并发数和我配置的不一致代码里设置了sink并行度8作业实际却按12的并行度在跑。原因在于我通过-D设置了parallelism.default12但代码里没有显式覆盖SQL作业的Sink算子继承了默认并行度。12个并发任务同时打开数据库连接超出了MySQL和驱动在低内存环境下的承受能力通信频繁中断。这个问题的根因就是配置优先级没理清我以为代码里的并行度设置是全局的但代码里只设置了Source并行度Sink没设于是从默认并行度继承。从这里我总结出的教训是JDBC连接器相关参数sink.buffer-flush.max-rows、sink.buffer-flush.interval、jdbc.batch.size等的解读依赖实际算子并行度资源配置优先级一旦混乱连接器调参也全乱了。排查JDBC连接器异常时先确认Sink算子的实际并行度再谈MySQL配置顺序不能反。5.2 案例二中间优先级任务一直pending问题出在slot计算回到开头那个事故。任务卡pending的时候我看了一下集群里有12个可用slot而卡住的作业并行度8按理说8个slot完全够用。但为什么一直pending后来在Web UI的TaskManagers页面看到12个slot平均分布在4个TM上每个TM配置的内存是4096m、slot数3。这个配置很特殊并行度8的任务需要8个slot但Flink分配slot时优先在同一个TM上连续分配导致4个TM里前两个TM能各塞3个slot再想分配第7、第8个时剩余TM的内存余量已经不够因为每个slot需要的内存超过了TM能提供的比例。也就是说slot数和内存之间的配比不合理即便总量够也会发生资源碎片。这就是“中间优先级”任务卡住的一个典型成因低并行度任务只需1-2个slot随便塞高并行度任务可以申请很多TM资源而中间并行度任务需要的slot数恰好跨越了TM边界又没配够内存最终pending。这个问题的解法不是调优先级而是调资源配比——把taskmanager.numberOfTaskSlots从3改成4保证并行度8正好分配2个TM内存充足任务立刻跑起来。5.3 案例三位运算符同款“想当然”踩坑说个配置合并顺序的经典错例和C语言里不加括号写位运算符一样容易“看起来对”。我见过有人在flink-conf.yaml里写了taskmanager.memory.managed.size: 1024m taskmanager.memory.managed.fraction: 0.4又在-D里设置了taskmanager.memory.managed.fraction0.2。他以为“既然fraction和size都设了那size优先fraction无所谓”结果作业启动时报了配置校验错误Flink检测到显式size和fraction同时存在且两者计算出的托管内存不一致直接拒绝启动。Flink的托管内存设置规则是这样的如果显式配置了managed.size则fraction被忽略如果同时配置且有冲突Flink会按校验逻辑拒绝或采用其中一个更保守的值。这个行为在不同小版本上有细微差异1.19.3的实测结果是报错并提示你删除其中一个配置。和位运算符的道理一样优先级是确定的但“确定的优先级”不等于“按我的直觉生效”靠猜一定会翻车。5.4 排查资源配置问题的四板斧最后把我的排查思路整理成一个固定流程遇到优先级相关的问题照着走一遍基本能定位第一步先确认部署模式和提交方式。不同模式下“生效范围”完全不同Session和Application不能混为一谈。第二步跑一次最简配置的基线提交看Web UI的日志和REST API确认默认值不做任何额外参数先把基准盘摸清。第三步用-D单变量逐个覆盖每次只改一个配置项通过Web UI和日志确认生效值是不是预期值。测的时候不要同时改多个项否则无法定位是哪个来源的问题。第四步查看错误日志和配置校验信息。Flink启动异常时日志里会明确指出哪个配置项冲突、哪个配置被忽略这比看任何文档都准确。6. 延伸场景MySQL同步到ClickHouse的资源配置实战6.1 场景画像与资源需求估算最近团队里用Flink做MySQL到ClickHouse的同步任务特别多正好结合资源配置优先级聊一下这个场景的配置模板。先给出一个标准的需求画像假设源MySQL有4张需要同步的表单表峰值写入量约5万行/分钟ClickHouse的目标表以MergeTree引擎为主写入需要按批提交以避免大量小文件。按照这个画像一个合理的资源配置是并行度4到8每个TM分配4096m内存其中JVM堆按2048m、托管内存按1024m估算slot数2到4。Source端用Flink CDCSink端用JDBC或ClickHouse专用连接器。这里的思路是源端并行度决定下表的读取能力Sink端并行度决定写入ClickHouse的batch大小和写入频率两者如果都跟随parallelism.default很可能出现Source压力大而Sink写入跟不上。6.2 一份实测稳定的资源模板我实际跑下来的模板如下三个关键点Source并行度显式指定、Sink并行度单独指定、内存配置不混用size和fraction。flink run-application -t yarn-application \ -D taskmanager.memory.process.size4096m \ -D taskmanager.memory.jvm-overhead.fraction0.15 \ -D taskmanager.numberOfTaskSlots4 \ -D parallelism.default8 \ -d /data/flink/jobs/mysql_to_clickhouse.jar作业代码里DataStreamString source env.addSource(mySqlCdcSource) .setParallelism(4); DataStreamRow transformed source .map(new TransformFunction()) .setParallelism(8); transformed.addSink(clickhouseJdbcSink) .setParallelism(6);这个配置下Source并行度4保证不会对源库产生太大压力Sink并行度6可以分摊聚合写入的压力默认并行度8作为中间算子的兜底。需要注意SQL Client方式写同步任务的话建议直接在SQL里为每个Source/Sink表显式设置scan.parallelism和sink.parallelism否则它们都会继承parallelism.default可能会导致MySQL连接数被错误放大。6.3 配合Spring Boot整合Flink的注意事项最后说一个Spring Boot整合Flink的常见误区。很多团队把StreamExecutionEnvironment放在Spring容器里管理想要在启动时动态调整资源配置。这个做法可行但要注意两点一是Spring Boot里通过configuration.setString()设置的资源配置只在本地Environment内生效。如果你最后是通过env.execute()提交到远程集群那么作业运行时的TM内存等资源还是由集群决定本地设置的这些资源项不会传递给远程集群。真正的做法是本地只设置并行度、checkpoint等作业参数资源类配置交给提交时的-D或集群配置管理。二是Spring Boot的配置文件application.yml和Flink的flink-conf.yaml是两个独立的配置体系不要试图在application.yml里维护Flink资源参数。我见过一个项目把taskmanager.memory配置写进application.yml然后疑惑“为什么没生效”——因为Spring Boot根本不会读取那些配置它只负责把参数传给Flink的底层对象而作业提交时的配置合并完全走Flink自己的规则。搞清这个边界整合起来会顺畅很多。最后再说一句我的个人体会。Flink的资源配置优先级不是一道需要背的题目而是一套需要“现场确认”的规则。不同部署模式、不同提交入口、不同连接器都会对最终生效值产生影响。我自己每次换集群或升级版本都会先用一个最小作业跑一遍基线验证把默认配置、生效日志和Web UI对照一遍再开始调优。这套方法虽然看起来笨但确实帮我避掉了很多“配置没生效”的坑。如果你也在新版本上遇到资源配置问题建议也先跑一遍我这个验证流程比对着文档猜要快得多。
返回列表