
做大数据平台这几年我见过太多团队把存算分离当成一个“语法糖”存储统一迁到对象存储计算集群挂上外部表就觉得架构升级完成了。结果生产环境一跑离线作业的输入扫描全部变成网络拉取作业时间从半小时涨到三个小时最后只能灰溜溜地把热数据导回本地。存算分离本身没有错错的是只拆不调度——存储和计算分开之后谁来决定每个任务到底跑在哪台节点上、什么时候扩缩容就成了一件决定成败的事。这篇文章就把这件事拆开计算节点动态调度到底在调什么、调度器是怎么做决策的以及我在真实集群上踩过的坑。正在做存算分离改造或者被资源浪费和数据倾斜折磨的平台工程师、架构师可以重点看看。1. 存算分离架构为什么需要“动态调度”这味药1.1 存算分离的本质是“分”动态调度是“合”存算分离的常规做法是把原来HDFS同时承担的存储和计算职责拆开存储下沉到对象存储、分布式文件系统或者独立存储集群计算层变成一套无状态的资源池节点上不再长期保存业务数据。这么做的好处非常直接——存储和计算可以各自独立扩容冷数据可以用更低成本的存储介质计算集群也能在跑批结束后迅速缩容而不是养着一批磁盘快满但CPU闲置的机器。但问题也随之而来。原来的Hadoop生态里数据本地性是最重要的调度原则之一任务调度器会优先把计算放到持有输入数据块的节点上尽量让计算去找数据而不是让数据找计算。存算分离之后数据在远端存储计算节点本地基本没有原始数据任务要读输入就必须走网络把数据拉过来。如果你不做任何调度约束所有任务一股脑全上最直接的结果就是存储网关被打爆、节点间网络拥塞任务大量时间花在等待数据上计算利用率反而比存算一体时更差。我常用快递分拣来类比这件事。存算一体相当于每辆货车自带一个小仓库货物和分拣员在同一辆车里效率天然高存算分离是建了一个统一大仓货车只负责拉货分拣员集中在一个大厅。这时候如果没有一个调度员决定“哪批货进哪个通道、哪个分拣员接哪条线”所有人就会全堵在大仓门口抢货。动态调度就是那个调度员它要做的核心事情是让每个计算任务在“去远端拿数据”的前提下仍然尽量少搬数据、少排队、不把某几条链路压死。从工程上讲动态调度解决的不只是“任务放哪台机器”这一个问题。它包括三个层面第一任务放置——每个作业的每个任务该分到哪个节点第二资源池伸缩——波峰时加节点波谷时减节点第三数据协同——哪些数据块被缓存在哪台节点调度时如何利用这些缓存。这三件事互相耦合缺一个都算不上真正的存算分离。1.2 静态调度为什么不够用早期很多平台做存算分离是带着“静态调度”的思路去做的给每个团队固定几个节点、固定一批资源槽位任务只能在指定队列里跑。这种模式在业务负载稳定、集群规模不大的时候是够用的但存算分离之后负载的波动会被放大。原因在于存算一体时任务和数据绑定在一起调度范围再乱数据读取成本也基本可控而存算分离后数据访问成本取决于任务离数据缓存和存储网关的远近如果资源是静态分配的你可能出现这样的情况A队列的节点上缓存了一大块热点数据但任务却因为队列隔离跑不到这些节点上B队列节点空闲但任务要读的数据全在远端每次跑批都要全量拉取。静态调度完全无法感知这些差异它天然把节点和任务做了物理绑定等于把“动态”两个字扔掉了。更麻烦的是静态分配无法应对节点故障和流量倾斜。一个节点宕机它的任务要么等重启要么人工迁移某张表突然变成热点只有一个静态队列承载其他队列再闲也帮不上忙。动态调度的价值就在于把集群当成一个整体池子任务按需流动、节点按需伸缩数据缓存位置和任务放置实时对齐。它不是为了让调度系统更“智能”而是存算分离这个架构本身要求调度器必须感知数据位置和实时负载否则分离出来的存储和计算就是两块互相拖累的孤岛。2. 动态调度器的核心实现原理2.1 先有全局视图节点状态如何被组织起来调度器要做决策第一个前提是知道集群里每台节点有多少资源、哪些任务正在跑、本地缓存了什么数据、到存储网关的网络带宽还剩多少。这一步看起来简单但真实集群里节点状态是高频变化的不能每次调度都现查监控系统。常见的做法是让每个节点上的Agent周期性上报心跳上报内容包括CPU核数和当前利用率、可用内存、本地磁盘缓存容量、当前运行的task列表、到对象存储网关的实测吞吐。心跳周期一般设3到5秒太短会制造大量无效消息太长又会让调度器对节点状态感知滞后。节点本身发生重大变化宕机、网络断连、磁盘故障时Agent要立刻发送事件消息而不是等下一个心跳周期这样调度器能快速把故障节点移出候选列表。调度器端需要把心跳信息聚合成一份全局资源视图。我建议用内存缓存维护三张表节点资源表、节点任务表、数据块位置表。节点资源表以nodeId为key保存实时资源快照和标签所属机架、可用区、机型、是否挂载SSD缓存盘节点任务表记录每个节点上正在运行的任务和已分配的槽位用于计算剩余资源数据块位置表是一个倒排索引记录“数据块的某个副本缓存在哪些节点上”。这三张表不需要实时到毫秒级能容忍数秒的陈旧性但必须在调度决策前异步刷新不能让调度主流程去等一个同步查询。2.2 调度决策的本质约束条件下的一次匹配每次调度本质上是在解一道匹配题任务提出需求CPU核数、内存大小、预估输入数据量、需要访问的数据块集合节点提供可用的资源池我们要找一个既满足硬约束、又让综合收益最高的节点。用一个生活化的例子患者任务挂号有的要看心内科数据在某个缓存节点有的需要重症监护需要超大内存号源节点有不同医生的剩余号。正常情况下不是哪个医生最闲就去哪个而是先过滤掉不能看的科室硬约束再在能看的医生里选排队最短、匹配度最高的软评分。落到实现上动态调度器的匹配过程分三步。第一步是硬过滤剔除资源不足的节点、剔除被标记为Quiescing或Reclaimable的节点、剔除不符合任务数据位置约束的节点——比如任务明确要求访问某个存储网关覆盖的数据就不要调度到另一个地域的节点。第二步是软打分对过滤后的节点逐个计算综合分数分数高的优先。第三步是更新预留选中的节点要在资源视图里临时扣减资源防止同一节点被并发调度器重复分配。注意“预留”和“实际分配”之间要有一个超时机制如果预留后任务迟迟没有在该节点启动要释放预留资源避免资源被占着不用。2.3 数据亲和性调度把任务调度到离数据最近的地方存算分离之后“数据本地性”并没有消失而是换了一种存在形式。数据原始副本虽然在远端存储但计算节点上可能挂了本地缓存盘缓存了一部分热数据块。调度器要做的是在决策时判断这个任务的输入数据哪个节点缓存命中率最高就把任务优先调度到那去。我建议给调度器维护一个数据块到节点列表的倒排索引数据来源可以监听存储侧的元数据事件或者周期性扫描任务访问日志来构建。任务提交后调度器先把任务的输入文件拆成block粒度然后逐个block查索引汇总每个候选节点的本地命中数据量。计算方式是一个简单的比例假设任务总共要读S字节的数据某个节点本地缓存命中了H字节那这个节点的数据命中率就是H/S。命中率越高说明任务启动后需要从远端拉取的数据越少执行速度自然更快。但这里有一个容易踩的坑不能只追求缓存命中率否则所有任务都会往已经缓存了热数据的几台节点上挤导致这些节点的CPU和网络被打满。数据亲和性必须和负载均衡叠加使用命中率只作为打分里的一个重要因子而不是唯一因子。如果几个节点的命中率都差不多再结合剩余CPU和内存去选只有命中率差异非常大的时候才允许牺牲一定的负载均衡程度。2.4 弹性伸缩动态调度器的另一只手任务放置解决的是“在已有的节点池里怎么放”而弹性伸缩解决的是“节点池本身该多大”。这层如果做不好动态调度就成了一个局部优化器——集群整体扩容慢任务照样排队集群不缩容钱照样烧。扩缩容的依据不能只看节点CPU利用率。典型的反例是凌晨跑批任务启动CPU瞬间冲高触发扩容脚本加节点结果任务跑完只要10分钟新节点刚ready就已经没活干了第二天早上又触发缩容。我推荐用调度器内部的两个信号做判断一是任务队列的等待深度也就是有多少任务在Pending状态、已经等了多久二是节点的连续空闲时间。当任务队列里新任务等待时间超过阈值比如30秒并且等待任务的数量持续增长就触发扩容。扩容时不是简单增加一台同规格节点而是根据等待任务的资源需求画像来选择节点规格——如果等待任务主要是大内存任务就扩内存型节点如果主要是IO密集任务就扩带大缓存盘的节点。缩容则只允许回收“已经连续N分钟没有被分配任何新任务”的节点而且必须先进入Quiescing状态等存量任务跑完再真正回收。弹性伸缩特别要防止“抖动”。我一般会加冷却时间一次扩容操作完成后的10分钟内不允许再次扩容一次缩容操作完成后的15分钟内不允许再次缩容。这样能避免调度器在两个状态之间反复横跳也避免扩出来的节点因为还没被充分利用就又被缩掉造成资源浪费和元数据反复迁移。3. 从原理到落地调度关键环节的实现示例3.1 调度器的模块划分一个能上生产的动态调度器不建议做成一个大循环里什么都干。我常用的模块划分是这样的资源采集模块负责接收心跳、维护全局资源视图任务队列模块负责管理Pending、Scheduled、Running、Finished、Failed几个状态调度引擎模块负责执行“过滤-打分-预留-派发”的主流程扩缩容模块根据队列深度和节点空闲度决定节点池大小元数据缓存模块维护数据块位置索引和任务历史画像。调度器自身的高可用可以通过主备模式实现主调度器负责实际派发任务备用调度器持续从同一个状态存储同步数据一旦主节点失联就接管。状态存储可以用分布式协调服务或者外部数据库关键是所有“任务状态变更”都要先持久化再返回成功否则宕机后会出现任务被派发到节点但状态丢失的情况导致任务重复执行或者悬挂。3.2 调度主循环的伪代码调度器的主循环适合用事件驱动而不是每秒钟全表扫描一遍所有节点和任务。事件来源包括任务提交、任务完成、节点心跳更新、节点故障、扩缩容指令。下面是精简版的伪代码可以直接作为一个实现骨架while True: event next_event() # 从事件队列拿一个事件 if event.type TASK_SUBMITTED: candidates filter_nodes(event.task) # 硬过滤 scores { node_id: compute_score(node, event.task) for node_id in candidates } target max(scores, keyscores.get) if target is not None: reserve(target, event.task) # 预留资源 dispatch(target, event.task) else: enqueue(event.task) # 进等待队列 maybe_scale_out() # 考虑扩容 elif event.type TASK_FINISHED: release_reserved(event.task) check_scale_in() # 看看有没有可缩容节点 elif event.type NODE_METRICS_UPDATE: update_node_view(event.node_id, event.metrics) elif event.type NODE_FAILURE: mark_unschedulable(event.node_id) reschedule_tasks(event.node_tasks)注意几个容易忽视的细节。reserve和dispatch必须放在一个事务里避免资源预留后派发失败留下一个永远不会被释放的“幽灵预留”。任务如果在节点上启动失败调度器要捕获启动反馈事件把任务重新放回Pending队列并短暂将该节点加入黑名单避免反复往同一坏节点派发。compute_score不应该在事件循环里调用重量级的机器学习模型否则单次调度耗时过长事件队列就会积压。3.3 给打分公式一个可以落地的版本打分公式是动态调度器最核心的决策逻辑。我给的版本是一个线性加权公式实际工程里可以按需扩展score(node, task) w1 * data_locality_hit_rate(node, task.inputs) w2 * resource_fit_score(node, task) w3 * bandwidth_avail_score(node) - w4 * migration_penalty(node, task)四个分量各有各的含义。data_locality_hit_rate是数据命中率表示任务输入集合中有多少比例的数据块已经缓存在该节点本地。这个值需要结合元数据缓存实时计算不应该在每个任务上重复查存储元数据而是由任务提交时批量获取输入文件的block列表然后和节点缓存索引做一次集合交集。resource_fit_score衡量节点剩余资源与任务需求之间的匹配程度。最理想的状态是任务需求刚好占剩余资源的60%到80%剩太多说明任务浪费了资源剩太少说明节点可能很快又不够用。计算时把CPU、内存、磁盘IO分别归一化再取加权平均。bandwidth_avail_score是节点到存储网关的可用带宽余量。这个值来自最近1分钟内的网络实测如果节点当前正在大量拉取数据它的带宽分数就应该很低。这个因子的权重在数据命中率不高时尤其重要。migration_penalty是一个惩罚项。如果任务之前已经在一个节点上跑过把中间结果放在本地那调度器应该尽量让它继续留在原节点避免频繁迁移带来的序列化、网络传输、缓存失效成本。初始权重可以参考w10.5w20.3w30.2w40.1。这个配置的前提是集群里数据访问成本是主要矛盾如果你们的作业大多是计算密集而不是数据密集可以调高w2、调低w1。更严谨的做法是跑一批历史作业样本用调度结果回放去搜索最优权重网格搜索在参数只有四五个时也够用。3.4 与Spark/Flink的联动方式调度器真正的落地场景是给计算框架派发执行单元。以Spark on K8s为例Spark自身有动态资源分配但它的动态分配只管executor数量不管executor在哪些节点上。如果调度器只告诉Spark“你扩容”executor还是可能被K8s默认调度器随机放到不合适的节点上那前期的数据亲和性计算就白做了。我建议在提交Spark作业时由外部调度器先做一次预调度生成一个“推荐节点列表”然后把这个列表通过Pod模板的nodeSelector注入到executor Pod上。比如spark.kubernetes.driver.podTemplateFile: /opt/template/driver-pod.yaml spark.kubernetes.executor.label.nodeType: compute-spot spark.kubernetes.executor.annotation.preferredNodes: node-0101,node-0102 spark.kubernetes.executor.volumes.hostPath.cache.path: /data/cacheFlink的实现思路类似调度器在提交作业前根据作业的并行度和每个TaskManager需要的Slot数预先规划好一组合适的节点把任务部署到这些节点上。因为Flink的TaskManager是长驻的调度器还需要考虑Checkpoint会写远端存储还是本地缓存这些路径的不同会影响打分公式中的带宽因子。4. 我踩过的坑和排查思路4.1 调度器变成新的元数据热点调度器上线后的第一个事故是存储元数据服务被调度器打爆。原因很简单我在最初版本里每次调度决策都同步去查输入文件的数据块列表和缓存位置。一个正常的离线作业可能有几万个输入文件每个文件十几个block每个block都要查一次缓存索引调度器单次决策就产生了几十万次元数据查询直接把底层的元数据服务QPS打到报警线。解决办法是把“查询闭环”改成“预取缓存闭环”。任务提交后调度器异步批量拉取输入文件的数据块列表在本地构建一个LRU缓存设置5到10分钟的过期时间。热点作业的热门输入文件缓存命中率能到95%以上元数据服务的压力瞬间掉下来。真正需要实时查的只有那些首次出现、且过期时间还特别短的冷门文件。4.2 用瞬时负载做判断集群疯狂扩缩容有一次我发现集群每天凌晨都在做无意义的扩缩容次数多到能让云厂商账单出现明显的“启动费用”。排查后发现是扩缩容模块用了“最近1分钟CPU平均利用率”作为唯一指标。跑批任务启动瞬间CPU冲高触发扩容任务跑到中间CPU趋于平稳又触发缩容新节点刚被扩出来还没有数据缓存调度器打分时会因为数据命中率低而一直不派任务给它于是它一直被判定为空闲马上又被缩掉。这次踩坑之后我把扩缩容信号改成了两个扩容只看任务等待队列深度和等待时长缩容只看节点连续无任务时间。同时加上冷却时间和最小存活时间——新节点扩容后至少存活10分钟缩容前必须连续空闲15分钟。这套组合在后续几个集群上表现稳定基本消除了扩缩容抖动。4.3 缩容把正在跑任务的节点回收了还有一次更严重的缩容流程把一台正在跑长任务的节点直接标记为可回收节点上几十个task全部被杀。看日志发现缩容判断只看“该节点当前没有新任务被分配”忽略了已经有旧任务在Running状态。长任务从提交到跑完可能要几个小时节点一直有存量任务当然不应该被缩。修复方式是给节点引入三种生命周期状态Active正常接收新任务、Quiescing不再接收新任务等待存量任务自然结束、Reclaimable存量任务为空且空闲超过阈值可以进行回收。缩容流程只允许从Quiescing进入Reclaimable而且Quiescing如果超过一定时间比如30分钟还有任务没结束需要先通知作业侧做checkpoint或迁移而不是直接杀进程。强制回收应该做成最后手段而不是默认路径。4.4 网络带宽其实是最大瓶颈数据亲和性打分跑了一段时间后发现有些节点缓存命中率明明很高但任务还是慢。追查网络监控发现那几台节点到存储网关的出网带宽已经接近打满。原因是我只算了“本地命中率”没算“未命中部分的拉取成本”。任务输入是100GB命中率90%意味着还有10GB要拉但如果命中率85%的节点带宽很空闲实际要拉的15GB可能比10GB更快完成。这提醒我把网络因子正式加进打分公式。具体做法是每台节点实时记录最近1分钟的实际出网吞吐和可用带宽上限两者相减得到带宽余量然后做归一化。调度时如果两个节点的数据命中率相差不超过5个百分点就优先选带宽余量大的节点。上线后整体任务平均耗时下降了12%效果非常明显。5. 一套可参考的动态调度参数模板5.1 关键参数速查表参数项建议初始值说明节点心跳周期3秒太短消息量大太长故障感知慢节点指标滑动窗口5分钟平滑负载波动避免瞬时抖动影响调度数据块位置缓存过期时间10分钟热点数据可缩短到1分钟扩容触发条件Pending任务等待超过30秒结合队列深度一起判断缩容触发条件节点连续空闲15分钟必须先进入Quiescing新节点最小存活时间10分钟避免扩了就被缩Quiescing超时时间30分钟超时后走优雅退出流程打分权重w10.5,w20.3,w30.2,w40.1按业务场景调整这些数值都是参考起点不要直接抄到生产环境。我见过电商大促场景需要把扩容阈值压到10秒也见过离线数仓场景可以放宽到60秒。最靠谱的方式是先把参数模板跑一个月汇总任务平均等待时间、节点平均利用率、数据命中率三个指标再针对瓶颈做微调。5.2 部署拓扑建议调度器本身建议独立部署不要嵌在某个作业里。主备两个实例共享一份任务状态存储协调服务用分布式协调组件保证同一时间只有一个主调度器在分配任务。每个集群部署一套调度器不要跨地域统一调度因为数据位置和网络延迟在跨地域场景下变化太大一个全局调度器很难兼顾。如果集群跨可用区调度器必须在节点标签上标明可用区打分时优先同可用区节点只有在同可用区资源不足时才允许跨可用区调度。对象存储本身有内网接入点调度器要确认节点到存储内网接入点的路由是经过专线还是公网这直接影响带宽余量的计算方式。我个人在实际操作中的体会是存算分离能不能出效果八成取决于计算节点动态调度这套系统做得是否细致。调度器不是一个可以上线后就不管的组件它需要持续根据作业特征和集群状态修正自己的打分逻辑。如果正在做类似改造我建议先用一个业务线的作业接进来盯着数据命中率和任务等待时间这两个核心指标跑两周再逐渐放大开关范围。最后分享一个小技巧每次调整调度策略前保留一份完整的调度回放日志改动上线后如果出现异常可以用回放工具对比新旧策略在相同输入下的决策差异定位问题会快非常多。