
上周在公司值班手机连着弹了好几条告警打开Flink的Web界面一看作业挂了好几个。再翻到JobManager日志最扎眼的一行就是Slot request bulk is not fulfillable!。做实时计算最怕遇到这种报错因为它不像业务数据异常那样有一条清晰的因果链你只知道资源调度在某个环节卡住了至于为什么卡住得顺着并行度、Slot、共享组一条线查下去。今天这篇就把这次排障的完整过程和背后的机制一起整理出来顺便把最近社区里经常提到的“Flink JDBC连接器异常”为什么最后也会引出这个报错说清楚。无论你是刚接触Flink的新手还是已经被资源类问题折磨过几轮的同学应该都能从里面找到对应自己情况的排查路径。1. 先搞清“Slot request bulk is not fulfillable!”到底是什么1.1 报错现场还原与触发时机先还原一下现场。这种报错通常长这样Caused by: org.apache.flink.runtime.JobException: Slot request bulk is not fulfillable! at org.apache.flink.runtime.scheduler.SchedulerBase.lambda... at org.apache.flink.runtime.scheduler.DefaultScheduler...具体堆栈在不同版本里会有差异但关键字一定是这句Slot request bulk is not fulfillable!。它不出现在业务算子日志里而是出现在JobManager的调度器日志里意思是作业向集群申请的一批slot集群在当前状态下无论如何都凑不齐。我整理了一下最容易触发这个报错的场景有三个。第一个是提交新作业。作业里配置的并行度总和超过了集群可用slot数JobManager一开始调度就会发现“这批slot我怎么给都给不齐”直接报错。最简单的情况就是集群一共8个slot你提交的作业需要10个slot。第二个是作业运行中手动调整并行度。比如通过命令行或者API触发rescaleFlink会尝试重新调度受影响的部分如果新的并行度超过了当时空闲slot的总量就会出现这个报错。第三个是failover恢复。这是最隐蔽也最容易踩的一种。作业因为某个算子OOM、网络异常、数据库连接超时等原因挂了Flink自动重启重启时JobManager需要把整个ExecutionGraph重新部署一遍一次性申请一整批slot。如果这时候其他作业还占着部分slot而旧作业的slot又没来得及释放一样会炸。所以看到这条报错先别急着改业务逻辑它本质上是“资源分配失败”不是“业务计算失败”。1.2 为什么叫“bulk”一次打包的Slot请求很多刚接触Flink的人看到“bulk”这个词会觉得奇怪为什么后面跟着“批量”两个字。这里就要说到Flink调度器的一个关键行为当一个作业需要启动或者重启时调度器会把这一轮所有需要部署的ExecutionVertex对应的slot请求打包成一个集合一次性提交给SlotManager去分配。这个“批量申请”的设计本身是为了效率。如果一个一个去申请slot每个请求都要经过一次资源匹配、反馈、确认的流程几十个task就是几十轮交互作业恢复速度会变得非常慢。批量提交可以让调度器一次性决定“这批位置是否整体可用”减少了大量的协商开销。但代价就是“整体满足”的约束。用食堂打饭来类比你端着餐盘需要同时取五道菜结果窗口今天只有四道菜你没法给自己端走四道然后剩下那道等会儿再补只能整盘退回去。Flink也一样如果这一批slot请求里的任何一个没办法落地整个部署就会失败接下来你会看到作业反复尝试分配、反复失败。这一点对排查极有帮助。你在看日志时要注意报错里常有“Not enough free slots”之类的前置日志那时候要立刻意识到是“总量不够”而不是某台机器调度出了问题。1.3 先排除两个干扰项不是OOM也不是连接异常我见过很多人在日志里同时看到几类报错就搞不清到底该先处理哪一个。这里有必要把三个容易混淆的问题放一起对比一下方便大家快速判断方向。现象典型日志关键字根因方向优先处理思路内存不足OutOfMemoryError / PMD Memory consumptionTM堆内存或直接内存不足调整内存配置减少单slot内存负载调度资源不够Slot request bulk is not fulfillable!并行度和可用slot不匹配或共享组限制扩容、降低并行度、检查slot sharing数据库连接异常SQLTransientConnectionException / Connection is not availableMySQL连接数或连接池被占满优化sink并行度、检查数据库连接上限注意这三类问题往往不是孤立出现的。完美的“帮凶”链条是连接异常导致task挂掉挂掉触发failoverfailover恢复需要重新申请大量slot而资源不足就爆出Slot request bulk。这个关系链我在第4章会用一个真实复盘展开。现在你只需要记住看到Slot相关报错要先把它的时间戳和前后日志串起来看判断它是“起因”还是“结果”。2. 并行度、Slot与TaskManager先把“工位模型”讲透2.1 工位模型一个Slot到底能放几个任务要彻底理解这个报错必须把Flink的运行模型掰开揉碎。我把Flink的资源调度比喻成公司租用办公区的过程。TaskManager就是一间办公室它是真正执行任务计算任务的进程承担CPU、内存、网络开销。每个TaskManager启动时可以配置taskmanager.numberOfTaskSlots也就是这间办公室物理上能放多少张工位。每个工位就是Slot同一时刻可以坐一个员工这里的员工就是作业中的一个并行子任务Subtask。一个并行度为P的算子会产生P个并行子任务也就是说需要P个Slot同时在线工作。集群总共能提供的Slot数很简单TaskManager数量 乘以 每台TaskManager的Slot数。比如3台TM、每台4个Slot总共就是12个任何作业申请的资源都不能超过这个数。但有一点必须说明Slot本身就是一段内存资源和执行线程的隔离单位它不是CPU隔离的单位。一个Slot在同一个时刻确实只跑一个Subtask但这张工位上的员工是可以换班的。一个Subtask执行完了、释放了Slot下一个Subtask可以立刻坐进来。所以Slot数本质上决定的是“同一时刻最多能并发运行多少个任务副本”而不是“一天下来能处理多少任务”。2.2 为什么执行计划上的并行度总和不等于Slot需求这一节是理解整个报错的重中之重也是大部分人栽跟头的地方。很多人潜意识里认为作业所有算子的并行度加起来是多少就需要多少Slot。这句话只对了一半甚至很多时候完全不对。Flink默认会把能串在一起的算子“链化”。比如一个简单的Source到Map再到Filter如果上下游并行度一致且没有shuffle边界它们会被合并成一条算子链所有算子像一个任务一样跑在同一个Slot里。这时候你表面上看到三个算子、每个并行度10实际上它们整体只占10个Slot而不是30个。但是一旦两个算子之间存在shuffle边界比如你调用了keyBy、rebalance、rescale或者两个算子的并行度设置不一样且需要重新分发数据上下游就必须同时运行才能一边产出数据一边消费数据。这种情况下两个算子不能复用同一批Slot计算总需求时需要把两边的并行度相加。我举个例子一个实时订单统计作业Kafka Source并行度是8中间做keyBy聚合的并行度是12最后JDBC Sink并行度也是12。Source和聚合之间是keyBy必然有shuffle边界聚合和Sink之间如果要按key写入也是shuffle边界。那么在调度阶段Source的8个Subtask、聚合算子的12个Subtask、Sink的12个Subtask都需要各自的Slot总共就是 8 12 12 32个Slot。看起来核心算子并行度只有12可一旦链路里有多个并行度不同的阶段Slot需求就是多个阶段的叠加。所以在看报错之前养成一个习惯打开执行计划图把每个存在shuffle边界的算子并行度分别列出来按阶段相加这才是你作业真实的Slot需求。2.3 总Slot数够了为什么还是失败这就是最让人恼火的地方明明集群总Slot数是够的可作业还是报 fulfillable。问题出在哪汇总下来有三类常见原因。第一个原因是可用Slot分散。集群总共20个Slot可用的有12个但这12个分布在四台不同的TaskManager上每台上只有3个。你作业某个并行阶段需要8个Slot调度器会给这8个Subtask挑选合适的TM结果有一台TM上只剩1个空位而需要放到这台TM上的Subtask可能就有2个分配关系对不上整体就失败。第二个原因是slot sharing group的设置。一个作业内部可以定义多个slot sharing group默认情况是所有算子都在一个名为“default”的共享组里。但有些人会在代码里给某些算子单独设置slotSharingGroup比如把Sink单独丢到一个自定义组里。这会导致这个分组内的算子独自占用一批Slot其他分组不能复用。表面上Slot总数没变但实际上集群资源被划分成了几个“隔间”某个隔间里的Slot不够照样报错。第三个原因就是集群里不止跑这一个作业。Flink的SlotManager是按整个集群维度管理的其他作业占用的Slot就是不可用。尤其是有些集群把几个作业混在一个Flink Session集群里跑别的作业高峰期把Slot全占了你这个作业大雪天提交上来自然分不到工位。理解这三点之后排查就有方向了。不管系统报告说得多隐晦你最终要回答的问题只有一个在这批申请发生的那个时间点可用的、且能放进对应共享组的Slot数量到底够不够。3. 从看到报错到改完参数一条完整排查路径3.1 第一步确认“申请量 vs 可用量”这两个数字很多人拿到这个报错就开始猜有的去调内存有的去改SQL运气好碰对了运气不好折腾到半夜。我现在的习惯是先把数字摆出来。打开Flink Web UI先看概览页Overview页会显示Registered Task Managers和Total Slots。比如显示Total Slots是12Available Slots是12那说明集群层面一个作业都没占多少资源。但具体作业需要多少还得看作业详情页里的Task Slots标签或者Task Deployment情况。在生产环境不方便开Web UI的情况下可以直接用REST API拉数据。curl -s http://jobmanager:8081/overview | python3 -m json.tool curl -s http://jobmanager:8081/jobs/overview | python3 -m json.tool curl -s http://jobmanager:8081/jobs/jobid | python3 -m json.tool | head -100第一个命令看集群总Slot数第二个命令看作业列表和状态第三个命令看指定作业的详细信息里面能找到每个ExecutionVertex的并行度和状态。很多版本还支持/jobs/jobid/exceptions可以拿到完整的异常列表。拿到数字后做减法。如果“请求需要的Slot数”大于“当前可用Slot数”根因基本明确直接走第3.4节的处置方案。如果“请求数”不大于“可用数”那就得继续往下看共享组和分布情况。3.2 第二步把执行计划和并行度设置从头捋一遍这一步的目的是搞清楚“为什么作业会申请这么多slot”。建议直接用Flink SQL里的EXPLAIN语法或者DataStream作业里把JobGraph导出来看一眼。对于SQL作业一个非常典型的坑是并行度不在你手里。Kafka Source的并行度默认会跟随Topic分区数。假设你Topic有24个分区Source并行度就是24。你在SQL里把parallelism.default设置成12那只能影响下游聚合和Sink的并行度Source仍然是24。如果Source和聚合之间存在shuffleSlot需求就变成了24加12比你以为的12高出一大截。DataStream里同理每个算子可以单独设置并行度优先级排序为算子单独设置 env设置 parallelism.default。建议把代码里的每个setParallelism都列出来对照执行计划逐个过。检查时重点看有没有“自动插入的重新分区”。当一个算子并行度与上下游都不一致时Flink为了保持数据分发正确会悄悄插入Forward或Rebalance的边这些边就是shuffle边界。每多一个shuffle边界Slot需求就可能多出一组并行度。3.3 第三步查Slot Sharing Group和TaskManager存活状况排除掉并行度设置的问题后接着看共享组。UI上查看作业的Execution Plan点任意节点可以看到它的Slot Sharing Group属性。如果所有节点都显示default那说明共享组没被改动可以排除这个原因。如果有些节点出现在自定义组里就要认真算一下自定义组的Slot需求是否超出了预期。同时还要确认TaskManager是不是真的都活着。尤其是在Kubernetes或YARN上跑Flink的集群某个TaskManager可能因为内存压力被内核杀了、容器被回收了、或者节点异常被JobManager判定超时移除了。看一眼kubectl get pod | grep taskmanager或者YARN的Node Manager列表确认实际存活的TM数量和你配置时假设的数量一致。一个很容易被忽略的细节是JobManager把TaskManager踢下线之后该TM上原本的Slot要经过一定的超时时间才会被回收释放。在恢复期内虽然UI上显示总Slot数减少了但旧的作业状态信息还挂在那边新的作业来申请资源时看到的可用数就会非常紧张。3.4 第四步处置三板斧扩资源、降并行度、并共享组排查到这里解决方案基本就三类我列成表格方便对号入座场景推荐操作说明作业并行度远超集群能力调低并行度到空闲Slot数以下能临时救急但可能牺牲吞吐集群总Slot不够业务又确实需要这么大并行度增加TaskManager数量或调大每台TM的Slot数扩容前先算单Slot内存是否满足Subtask需要自定义Slot Sharing Group导致资源浪费去除多余共享组把所有算子合并回default不花钱、不改业务还能立刻见效集群里其他作业长期抢占Slot给作业规划独立Flink集群或用资源队列做隔离会话集群混跑容易埋雷具体怎么改配置如果决定增加Slot可以调整TM的启动参数taskmanager.numberOfTaskSlots: 8如果决定降低并行度Flink SQL可以直接设置SET parallelism.default 16;DataStream作业则可以设置env.setParallelism(16)。这里要提醒一句不要一上来就把全局并行度调小那样等于把所有算子的并发度都拉低。最合理的是只把出问题的那个阶段并行度降下来或者调整Kafka Source的并行度上限保住核心计算环节的吞吐。4. 热词联动JDBC连接器异常怎么和Slot不足扯上关系4.1 复盘一次“写MySQL把作业写崩”的过程最近在社区里经常看到“flink的jdbc连接器异常”这个热词很多人跑来问“我的SQL没问题连接串没问题数据库也没满为什么JDBC Sink老报连接不可用”这里我想分享一次实际排障因为这看起来和Slot不足毫无关系最后却一起炸了。当时作业架构很简单Flink从Kafka读订单数据做完聚合后写入MySQL。高峰期一到先是日志里出现了一堆异常Caused by: java.sql.SQLTransientConnectionException: HikariPool-1 - Connection is not available, request timed out after 30000ms.连接池里拿不到连接Sink任务开始连接超时后面的写入数据部分堆积。紧接着一部分Sink的Subtask因为超时异常失败触发作业failover。就在恢复阶段JobManager日志里紧跟其后地出现了那句Caused by: org.apache.flink.runtime.JobException: Slot request bulk is not fulfillable!两个八竿子打不着的错在一条时间线上串成了因果链。梳理一下发生了什么高峰期数据库连接被Flink连接器本身和其他服务占满了Flink Sink想要新的连接拿不到任务反复失败。失败触发重启重启时要重新申请整条链路的Slot而这个时候旧作业的Slot还在等待清理、集群里还有别的作业占着资源一次性申请量瞬间超出了可用数。于是真正的病根——连接池打满——反而被Slot报错盖在了后面。4.2 判断“先有鸡还是先有蛋”先看时序再动手遇到这类混合报错最大的敌人是急躁。很多人一看到“Connection is not available”就开始调sink并行度或者连接池结果作业频繁重启第二天又炸一次而且每次炸出来的都是Slot报错。我的判断方法很简单看时间戳。先确认连接异常出现的时间点再确认Slot fulfillable报错出现的时间点。如果连接异常在前、Slot报错在后说明Slot报错只是恢复过程中的次生灾害主线应该是“为什么连接不够用”。如果Slot报错在前、连接异常是因为task一直起不来导致的那主线就是资源不足。然后看持续状态。打开Flink Web UI的BackPressure面板如果Sink上游一直处于High状态问题大概率出在下游写入速度跟不上数据流入速度。再回头看MySQL的max_connections、当前连接数、以及慢查询基本能锁死是数据库侧还是Flink侧的问题。实际操作中我习惯按这个顺序排查看作业是不是在“运行-失败-重启-再失败”的循环里看背压是不是先出现在Sink上游看数据库当前连接数和连接池可用连接数对比两类报错的时间戳先后这套顺序能保证在五分钟内找到主线而不是被日志里五花八门的Caused by带偏。4.3 可以从容调整的几个参数解决这类问题不是一味增加并行度不要误会。如果数据库已经被连接数打满了你增加Sink并行度只会让每个任务去抢连接情况更糟。合理的做法是把Sink的写入改成批量、低频让固定数量的连接能发挥更高的吞吐。Flink SQL的JDBC连接器一般支持批量写入参数可以这样设置CREATE TABLE order_sink ( order_id STRING, amount DECIMAL(10,2), ts TIMESTAMP(3) ) WITH ( connector jdbc, url jdbc:mysql://your-host:3306/test, table-name order_sink, driver com.mysql.cj.jdbc.Driver, username your-user, password your-password, sink.buffer-flush.max-rows 1000, sink.buffer-flush.interval 2s );max-rows1000表示攒够1000条才写一批interval2s表示最多等2秒就写一次这两个参数结合起来能大幅降低连接被频繁占用和释放的压力。DataStream里如果直接使用JdbcSink也可以配置批量执行JdbcExecutionOptions execOptions JdbcExecutionOptions.builder() .withBatchSize(1000) .withBatchInterval(Duration.ofSeconds(2)) .build();另外如果连接池用的是HikariCP不要只调大maximumPoolSize还要把connectionTimeout调短一点。宁可让任务快速失败触发重启也不要让几十个连接请求全部挂在等待队列里占着线程。数据库侧的max_connections建议预留30%到40%的缓冲给其他服务和运维操作而不是被Flink作业吃满。5. 经验复盘这类报错最容易忽悠人的地方5.1 “总量够”和“能分配”之间隔着一层分布做了这么多年运维和实时开发我发现这个报错最著名的“谎话”就是怎么看集群都够用。有时候总Slot数确实够但作业就是起不来原因在于分配并不只看总量而是看每个TaskManager上的分布是否合理。举一个很简单的例子你有4台TaskManager每台2个Slot总共8个。现在另一个作业占用了其中一台机器的2个Slot剩下6个。你的作业申请6个Slot数字上刚好够但如果要部署的6个Subtask里有部分逻辑要求尽量发送到同一台机器比如同一个Source的多个分片而那台机器的空余Slot已经不足调来调去都凑不齐最终就报fulfillable。Flink的SlotManager不是物理机上的资源池它不会把所有空闲Slot打散后重新切块给你。每个ExecutionVertex还是要落到具体某一台TaskManager的某个Slot上只要有一个落地不了整批申请就失败。所以在做集群容量规划时不要只看“总Slot / 峰值并行度”这两个数还要看单机Slot配比和作业内并行度关系。如果作业里单个阶段并行度是12而每台TaskManager上可用Slot数长期只有4或5那这个阶段天生就容易分配失败要么重新规划TM的Slot分配要么调整作业的资源请求。5.2 failover恢复期为什么平时看着没事一出问题就炸这个现象我见过太多次了作业稳定运行几周内存也没爆数据库也正常某天突然因为一条脏数据引发反压或者异常紧接着failover然后就报Slot不足。你第一反应是“之前不是一直够吗怎么恢复的时候反而不够了”。原因是恢复这个动作本身就比正常运行要更多瞬时资源。正常运行的时候部分算子可能因为数据水位不齐、或者某些Subtask处理得快而提前释放Slot但是failover恢复时JobManager要把整个ExecutionGraph的所有Subtask一下子全铺开而且旧Subtask占用的Slot还没完全释放新一批又冲上来瞬时需求直接翻倍。面对这个问题除了预留足够的Slot余量外还可以把恢复节奏放平稳。常用的参数组合是这样的restart-strategy: fixed-delay restart-strategy.fixed-delay.attempts: 3 restart-strategy.fixed-delay.delay: 10s把restart-strategy.fixed-delay.delay从默认的秒级改成10秒以上可以给旧Slot的清理释放留出缓冲时间。如果你用的是Flink 1.15以上的版本还可以调整slotmanager.taskmanager-timeout增加JobManager等待TaskManager注册或Slot释放的耐心。但这里有一个反向提醒如果物理上真的没有更多可用Slot把等待超时调得再大也没有意义作业只会多等几分钟然后照样失败。所以调参数的唯一前提是——你已经确认集群有Slot正在释放只是需要时间。5.3 我自己一直保留的三个排查习惯排这类资源调度问题排多了我给自己定了三条规矩谈不上多高级但每次都能让我快速冷静下来。第一任何作业上线前先按目标并行度跑一次空载任务。不要等生产出问题再发现资源不够。现在Flink作业启动很快空载跑5分钟能够看到Slot申请是否顺利、有没有因为共享组配置导致分配失败花不了多少成本。第二集群监控里一定要有Slot使用率指标。我自己的告警阈值是80%超过这个数就发预警给整个集群留出20%的余量。这20%不是浪费是给failover恢复、临时扩并行度、以及旁边作业的突发流量准备的缓刑空间。第三看到这条报错时第一反应不是改SQL也不是换连接器而是把阻塞点的时间戳和数字量化出来。打开JobManager日志把“什么时候申请的、申请了多少、当时可分配数量是多少”这三个问题回答清楚再动手调任何参数。很多时候问题一眼就能看出来根本不用碰业务代码。如果你也在排障时看到“Slot request bulk is not fulfillable!”建议顺着这条路径先稳住自己。实时计算里很多现象看着像业务问题根子却在资源调度这层。只有把底层这个模型彻底想明白遇到问题才不会慌。