ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

基于异构计算边缘化的分布式任务调度:用劣质硬件搭建弹性集群实践

基于异构计算边缘化的分布式任务调度:用劣质硬件搭建弹性集群实践 1. 项目背景为什么要把“破硬件”重新用起来我先说个场景你可能也有同感。公司机房角落堆着一批淘汰下来的办公电脑配置五花八门有的是八年前的i3有的是内存只有8G的小主机硬盘还是机械盘。正常来说这些东西该走报废流程但真扔了又觉得可惜尤其是边缘侧的项目——摄像头识别、传感器数据处理、日志清洗、定时爬虫——这些活得没那么“重”但数量多、很碎专门买新服务器去跑反而亏。我这次做的事就是把这些“劣质硬件节点”组到一个分布式集群里让它们在边缘侧承接计算任务。标题里的“异构计算边缘化”拆开看就是两件事一是异构节点CPU架构不同、内存大小不同、磁盘速度不同甚至有的节点装的是ARM开发板二是边缘化任务不往中心机房送在离数据最近的地方就地处理。再加上“基于劣质硬件节点的分布式弹性调度”说白了就是用一套调度系统把一堆性能参差不齐、随时可能掉线的普通设备当成一个还算稳定的计算资源池来用。这个题目适合谁看我建议这几类朋友重点参考一是手头有闲置设备但不知道怎么利用起来的二是在做边缘计算或者IoT项目、不想给每个点位都配高性能服务器的三是已经在玩分布式任务调度、但只接触过K8s或XXL-Job这类成熟框架想自己理解底层调度逻辑的。有些人可能会问直接用K8s不就行了吗说实话K8s确实能做节点管理但K8s对节点稳定性、系统组件的要求都不低让一堆Windows老电脑和Linux开发板混进同一个K8s集群光折腾kubelet和容器运行时就能劝退一大半人。我们的场景更“野”一点所以最后决定自研一个轻量调度器。这篇文章会把整个思路、关键代码、踩坑过程全部拆开讲你可以直接当参考也可以拿来改造成自己的调度框架。2. 整体设计思路不是管理机器而是管理“信用分”2.1 先理清楚我们到底要调度什么设计之前我先给自己提了一个问题任务都有什么特征调研了一圈边缘侧的任务类型挺杂的我总结成三类定时批处理任务比如每天凌晨汇总昨天的监控数据、定期清理日志、定时拉取外部接口数据。事件触发任务比如摄像头拍到移动物体后触发一次图片识别传感器数值超阈值后触发一次告警计算。常驻轻服务比如一个HTTP接口服务负责接收传感器上报数据并写入消息队列。这三种任务对调度的要求差异很大。批处理任务需要的是“到点了一定有人跑”事件触发任务需要的是“尽量快的响应”常驻轻服务需要的是“尽量别重启”。所以我从一开始就不打算做一个大一统的任务系统而是聚焦在“批处理任务”和“事件触发任务”这两种上。常驻轻服务的调度本质上是服务发现和负载均衡的活儿跟“劣质硬件”这个话题的匹配度不够先放一边。2.2 异构节点的难点不只是性能差异“异构”这个词听起来高大上但实际接触后你会发现真正难的不是“性能差异”而是“不稳定”和“不可预测”。性能差异A机器跑一个识别任务需要200毫秒B机器可能要用800毫秒这还不算等你把20个任务塞给AA的处理时间可能会线性恶化到几秒。资源口径不统一有的机器内存8G但系统占了3G实际可用只有5G有的机器CPU是4核但其中两个核被别的进程占满了。你要是只按“机器标称配置”来调度肯定翻车。随时可能掉线办公电脑被同事重启、开发板供电不稳、老机器风扇积灰导致温度过高自动关机……这些在数据中心里不太常见在边缘侧就是家常便饭。网络环境差节点之间可能是Wi-Fi连接甚至走4G延迟高、抖动大。就算调度器把任务发过去了结果能不能传回来都是个问题。所以我的核心设计理念是不要试图精确管理每一台机器的资源而是给每台机器建一个“信用档案”用历史表现说话。机器说自己内存多大、CPU多快这些只能作为参考真正决定任务分给谁的是它过去一段时间里成功完成了多少次任务、平均耗时多少、掉线过几次。2.3 调度器自身架构中心化调度 节点自治现在主流分布式调度大致分两派中心化调度所有节点向一个调度中心汇报由中心做决策和去中心化调度比如一致性哈希、gossip协议。我选的方案是中心化调度为主节点保留自治能力。理由很简单边缘节点数量不会特别大百台以内是常态中心化调度完全撑得住而且实现简单、问题好排查。去中心化虽然听起来更“高级”但在节点质量参差不齐的环境下反而容易因为节点间通信不稳定导致脑裂。中心化调度器的组成也不复杂就是这几块节点管理模块负责接收节点心跳、维护节点状态、计算节点信用评分。任务队列模块接收任务请求按优先级和类型排队。调度引擎从任务队列里取出任务结合节点信用评分决定派给哪个节点。结果回收模块接收节点执行结果超时检测失败重试。但要是调度中心挂了整个集群就瘫了所以节点本身也要保留一定的自治能力。比如节点发现自己和调度中心失联超过一定时间就会把本地正在跑的任务继续执行完结果暂存在本地等恢复连接后再补报。3. 核心机制实现节点画像、信用评分与弹性伸缩3.1 节点画像先摸清每一台机器的底细3.1.1 静态信息注册节点第一次启动时会向调度中心发送一份“自我介绍”报文包含以下信息{ node_id: edge-device-014, os_type: linux, cpu_arch: x86_64, cpu_cores: 4, cpu_freq_mhz: 1800, mem_total_mb: 8192, disk_total_gb: 256, net_upload_kbps: 5120, net_download_kbps: 10240, labels: { location: factory-a, owner: security-team } }这里的labels很重要它可以帮你实现“标签调度”。比如摄像头识别任务只派发给装了GPU或者带AI算力NPU的节点日志清洗任务只派发给磁盘空间充足的节点。没有标签调度器就只能瞎猜应该在哪儿跑。但是注意这个静态信息只是“参考值”不是“承诺值”。因为机器可能被其他进程抢占资源所以更关键的是动态信息。3.1.2 动态健康度上报节点启动后每隔10秒向调度中心上报一次动态状态格式类似这样{ node_id: edge-device-014, timestamp: 1713600000000, cpu_usage_percent: 23.5, mem_available_mb: 4096, disk_free_mb: 51200, load_avg_1m: 1.2, heartbeat_seq: 1234 }调度中心收到这份数据后不只是简单地入库而是要算出几个派生指标CPU可用率100 - cpu_usage_percent内存可用率mem_available_mb / mem_total_mb综合负载因子正常情况下load_avg_1m应该小于cpu_cores如果大于说明机器已经过载调度优先级要降低。我实际用的时候发现光看单次上报不够还需要用滑动窗口来平滑。因为CPU使用率这玩意儿波动极大可能是瞬时尖峰也可能是持续过载。我维护了一个长度为5的滑动窗口每次计算平均值避免了因为一次尖峰就把节点拉黑的情况。3.2 信用评分体系让烂机器“将功补过”这是整个调度器最有价值的部分。我把节点信用评分定义为(0, 100]之间的一个值初始值设为60也就是说每台新加入的机器都先默认“能干活”但不够优秀得靠实际表现加分。影响信用分的主要因素有三个任务成功率节点执行任务成功的次数占比。每次成功加1分每次失败扣5分。为什么失败扣分更多因为失败意味着任务要重跑浪费的是整个集群的时间代价更高。任务超时率如果节点领了任务但迟迟不回结果超时一次扣3分。这比失败还讨厌因为调度中心要一直占用内存保存这个任务的状态直到超时判定。心跳稳定性如果节点连续三次心跳丢失直接扣10分并标记为“疑似离线”如果恢复了每次心跳加0.5分慢慢养回来。这个评分体系的关键在于它不歧视配置差的机器。老机器只要每次都兢兢业业完成分配给它的任务信用分会稳定上升新机器就算配置再好如果总超时、总失败信用分也会掉下去。这跟现实中用人是一样的态度和结果比出身重要。3.3 调度算法不是选最好的而是选最合适的当调度中心从任务队列里取出一个任务时要决定把它派给谁。我的策略分三步走第一步候选节点过滤节点状态必须是online信用分大于30。节点标签必须满足任务的标签要求比如任务要求locationfactory-a。节点当前已分配的任务数不能超过上限上限由节点配置决定比如CPU为2核的机器并发上限设为24核的设为4。第二步候选节点排序排序考虑了四个因素按权重计算总分因素权重说明信用分0.4历史表现越可靠越优先当前负载率0.3(mem_used cpu_used) / 2负载越低越优先历史平均耗时0.2同类任务耗时越短越优先网络延迟0.1与任务数据源所在节点的网络延迟越低越优先最终分数 信用分*0.4 (100 - 负载率)*0.3 耗时得分*0.2 延迟得分*0.1取分数最高的节点。第三步任务下发与超时控制任务数据通过HTTP或gRPC推送给节点节点收到后立刻返回ACK然后异步执行。调度中心同时启动一个超时定时器默认超时时间由任务类型决定批处理任务10分钟事件触发任务30秒到1分钟可指定自定义超时时间如果超时未收到结果调度中心会把这个任务重新放回队列尾部并把对应节点的超时次数加一。3.4 弹性调度怎么应对节点“说没就没”的突发情况“弹性”这个词被很多云厂商用滥了但在这里它是一个很朴素的诉求节点减少时任务不能断流节点增加时任务能自动分流过去。我做了两个机制来保证弹性一是任务级容错搬迁。假设一个节点在执行任务的过程中突然掉线了调度中心怎么知道情况A调度中心很久没收到心跳判断节点离线主动将该节点上所有未完成任务重新入队。情况B节点在处理任务前先把任务ID和状态写入本地磁盘重启后如果发现之前有未完成的任务重新上报给调度中心申请恢复执行。加上一个“任务幂等”设计同一个任务即使被分发两次节点通过唯一任务ID判断是否已经执行过避免重复计算。二是按节点压力自动扩缩容。我设了一个指标系统整体积压量即任务队列里未分配的任务总数。如果积压量连续3次超过阈值比如50调度器会在候选节点里选信用分最高的节点把它的最大并发数上调20%相当于“压榨”高性能节点如果积压量连续5分钟内低于10就把并发数回调。这个扩缩容的粒度是“并发数调整”不是“增减容器实例”因为这里是物理节点不能用K8s那套Pod伸缩逻辑。3.5 调度核心代码一个简化但能跑通的最小实现这里我写一个简化版的核心调度逻辑方便你理解整体流程用的语言是Python实际项目中如果要用建议换Go或Java性能和并发控制会好很多。import time import threading from collections import defaultdict from dataclasses import dataclass dataclass class NodeInfo: node_id: str cpu_cores: int mem_total_mb: int status: str online credit_score: float 60.0 current_load: float 0.0 running_tasks: int 0 max_concurrent: int 2 history_avg_ms: float 500.0 heartbeat_timeout_count: int 0 class Scheduler: def __init__(self): self.nodes {} self.task_queue [] self.task_result {} self.lock threading.Lock() def register_node(self, node_info: dict): with self.lock: node NodeInfo( node_idnode_info[node_id], cpu_coresnode_info[cpu_cores], mem_total_mbnode_info[mem_total_mb], ) # 根据配置不同设置不同的并发上限 if node.cpu_cores 8: node.max_concurrent 8 elif node.cpu_cores 4: node.max_concurrent 4 else: node.max_concurrent 2 self.nodes[node.node_id] node def update_heartbeat(self, node_id: str, cpu_usage: float, mem_available_mb: float): with self.lock: node self.nodes.get(node_id) if not node: return if cpu_usage 60 and mem_available_mb 1024: node.credit_score min(100, node.credit_score 0.1) else: node.credit_score max(0, node.credit_score - 0.2) node.current_load (cpu_usage (100 - mem_available_mb/node.mem_total_mb*100)) / 2 def add_task(self, task): with self.lock: self.task_queue.append(task) def schedule_once(self): with self.lock: if not self.task_queue: return task self.task_queue.pop(0) candidates [] for node in self.nodes.values(): if node.status ! online: continue if node.credit_score 30: continue if node.running_tasks node.max_concurrent: continue score ( node.credit_score * 0.4 (100 - node.current_load) * 0.3 max(0, 100 - node.history_avg_ms / 20) * 0.2 80 * 0.1 # 延迟得分简化处理 ) candidates.append((score, node)) if not candidates: self.task_queue.insert(0, task) return candidates.sort(reverseTrue, keylambda x: x[0]) _, target_node candidates[0] target_node.running_tasks 1 # 异步执行任务这里省略具体执行逻辑 print(f任务 {task[id]} 派发给节点 {target_node.node_id}) threading.Thread(targetself._exec_task, args(task, target_node)).start() def _exec_task(self, task, node): try: # 模拟执行任务耗时 time.sleep(1) with self.lock: node.running_tasks - 1 node.credit_score min(100, node.credit_score 1) node.history_avg_ms (node.history_avg_ms * 9 1000) / 10 except Exception: with self.lock: node.running_tasks - 1 node.credit_score max(0, node.credit_score - 5)这段代码已经覆盖了调度器的核心骨架节点注册、心跳跟新信用分、按综合评分派发任务、任务失败扣分。你如果要扩展可以在_exec_task里增加任务结果回传、超时重试、任务迁移等逻辑。实际项目中我建议把并发控制从threading.Lock换成分布式锁或者Go的channel避免单点瓶颈。数据存储也不要存在内存里用Redis或者SQLite持久化否则调度中心一重启节点信息和任务状态全丢了。4. 实操过程部署与参数调优从0搭一套能跑的集群4.1 部署拓扑一台指挥 九台干活我这次搭建的测试环境是这样的一台调度中心用的是一台闲置的i5迷你主机8G内存跑Kafka用来接收任务、MySQL用来存节点历史数据、调度服务。九台边缘节点配置差异很大3台老式办公电脑i3-21004G内存机械硬盘2台迷你主机J41258G内存固态硬盘2台树莓派4B4G内存microSD卡2台淘汰的笔记本i5-4200M6G内存机械硬盘这九台机器分布在办公区、仓库、门口三个位置中间通过局域网连接有一个位置只有Wi-Fi覆盖网络质量一般。4.2 关键参数怎么定压出来的不是想出来的这节我直接给出几个我实测过后比较稳的参数配置以及调整依据。1. 心跳间隔10秒超时阈值30秒太短了浪费带宽太长了发现不了节点掉线。10秒一次心跳连续3次未收到就判定离线这个节奏在LAN环境下基本够用。但如果你的节点走4G网络心跳间隔建议调到20秒超时阈值60秒免得因为网络抖动误判离线。2. 任务队列堆积阈值50扩缩容检测周期3次我一开始设的阈值是10结果发现稍微有个任务峰值就会触发扩容调度器忙得不可开交。后来调到50稳了很多。判断是否扩容不能只看单次要连续3次检测都超标才动作防止抖动。3. 每个节点的最大并发数默认等于CPU核数最高不超过核数的1.5倍为什么不能超过太多因为在劣质硬件上任务往往是IO密集型和CPU密集型混杂的并发太高会导致线程切换开销剧增反而拖慢整体速度。我压测过J4125这颗4核处理器并发设到6的时候吞吐量达到峰值再往上就掉头了。所以最终公式是max_concurrent min(cpu_cores * 1.5, 8)。4. 任务超时时间批处理默认10分钟事件触发默认30秒这里遵循“留足余量但不惯着”的原则。如果节点信用分高于80我会允许它申请“长任务模式”超时时间自动翻3倍。4.3 压测过程和观察到的现象我写了一个压测脚本模拟“图像识别”任务随机生成一张带噪声的图片要求节点运行一个简单的边缘检测算法返回处理耗时。第一轮压测一次投递180个任务每个任务数据量大概100KB。结果很有意思树莓派的CPU利用率瞬间飙到95%但任务耗时一直在800毫秒左右不算太差。老办公电脑反而翻车了——机械硬盘IO成了瓶颈任务还没开始算写日志就把时间吃掉了大半。J4125迷你机表现最稳定处理耗时稳定在180毫秒左右。这说明一个很重要的规律在异构环境里瓶颈往往不在CPU算力而在IO子系统和内存带宽。后来我针对机械硬盘节点做了个优化所有任务写入先写到内存队列由单独线程批量刷盘减少了频繁的IO中断老电脑的任务耗时从800毫秒降到了400毫秒。第二轮压测我模拟节点突然崩溃。压测进行到一半直接把一台笔记本的电源拔了。调度中心大概在30秒后判定节点离线然后把它正在执行的4个任务重新入队派给其他节点。整个过程任务没有丢只是整体耗时多了大约1分钟。这个表现说实话比我预想的好。5. 常见问题与排查技巧实录5.1 问题一节点反复“离线-上线”任务一直被踢皮球这个问题我一开始特别头大表现是某些节点每隔十几分钟就掉线一次然后又自己回来信用分被扣到谷底任务都被分给别的机器了。排查过程是这样的先看调度中心的日志确认节点确实有断连记录。然后登到节点上看系统日志发现是网络接口在重启进一步查看是这块机器装了某个自动更新的软件更新完会重置网络。最后定位到是有线网卡和Wi-Fi网卡同时启用系统在两个网卡之间频繁切换导致网络闪断。解决办法把不用的网卡禁用固定IP设置同时把自动更新改为手动。如果你的节点分布在弱网环境还有一个更稳妥的办法——给节点本地加一个“任务缓冲池”就算网络断了任务也在本地排队等网络恢复再上报结果。5.2 问题二任务失败后重试结果把集群“打爆”了这是很多分布式系统新手都会踩的坑。有段时间某个数据源接口不稳定拉取任务频繁超时调度器一看失败就重新入队结果任务越积越多最后把整个集群的任务队列塞满连正常任务都排不进去。这里我做了三个改动重试次数上限单个任务最多重试3次超过3次进入“死信队列”由人工处理。退避策略第一次重试等待5秒第二次30秒第三次5分钟。这个指数退避能有效防止“刚失败立刻重试又失败”的无效循环。熔断如果某个节点连续失败10次调度中心直接把它拉黑10分钟不让任何新任务派给它避免一个坏节点拖垮整个调度流程。5.3 问题三结果回传丢失任务到底算成功还是失败有一次遇到个诡异情况节点明明执行成功了但调度中心一直没有收到回传结果超时后把任务重新派发给了其他节点导致同一份数据被处理了两遍。查了之后发现是网络断开的时间点太尴尬——节点执行完成刚要回传的时候网络恰好断了。我最后用了“两阶段提交”的思路第一阶段节点执行完任务后先把执行结果写入节点本地缓存并标记为“待同步”。第二阶段节点与调度中心的连接恢复后主动把本地缓存中的结果同步过去调度中心收到结果后返回确认节点再删除本地缓存。这个机制很简单但能确保结果不丢、不重复。如果你用的是成熟消息队列比如Kafka或者RabbitMQ也可以用它们来做结果回传的缓冲道理是一样的。5.4 问题速查表现象可能原因排查方法解决办法节点信用分莫名下降心跳间隔太长或网络抖动检查节点端的ping延迟和丢包率适当放宽心跳超时阈值任务全部堆积在队列所有节点信用分低于30查看各节点信用分变化趋势手动重置信用分并排查原因某任务执行时间差异巨大不同节点的CPU架构、主频差异对比各节点同类任务的平均耗时按耗时设定任务分桶分别调度调度中心内存持续增长任务结果或状态没有及时清理用jstat或pprof查看内存占用定期清理已完成任务的状态记录节点收到任务却迟迟不执行节点端任务队列被其他任务堵住查看节点CPU和IO占用率限制节点单线程处理任务数6. 一点经验总结与后续演进思路折腾完这套系统我最大的感受是做异构计算最重要的事情不是算法多精巧而是心态要摆正——你必须接受“有的节点就是会把事情搞砸”这个现实然后设计一套机制让整体系统不因个别节点而崩溃。信用评分加弹性重试这套组合虽然不是新东西但在劣质硬件场景下特别好用因为它在不增加硬件成本的前提下尽量榨干了每一台机器的价值。如果你也想在自己环境里复现这个东西我的建议是别一上来就追求复杂。第一步先把节点注册、心跳、信用分三个基础功能做出来跑通每个任务都能被分配到节点并返回结果然后再加上失败重试和超时控制最后再考虑弹性扩缩容和标签调度。每一步都找一批真实任务来压测别用纯模拟任务因为验证不出IO瓶颈和网络抖动这些真实问题。后续我打算做两件事一是把调度器改成基于插件式的架构方便接不同的执行器比如Docker容器、进程、Serverless函数二是把节点的信用评分和任务执行结果做一次离线分析看看能不能训练出一个更智能的调度模型提前预判节点会不会掉线。不过这些都是后话了先把当前这套稳定跑起来再说。最后分享一个不太起眼但对整体体验影响很大的小技巧调度中心一定要可视化。我后来给调度器写了一个简单的Web面板能看到每台节点的实时状态、信用分曲线、任务排队情况。没有这个面板的时候出了问题基本靠猜有了面板故障定位时间从小时级降到了分钟级。做分布式系统可观测性永远是第一位的。
返回列表