简介:《云计算之分布式计算》是一份面向云计算、大数据及相关专业学习者、教师和从业者的PPT课件,聚焦分布式计算在云计算体系中的核心作用与实现思路。内容从大数据时代的数据爆炸背景切入,引用加州大学和IDC研究数据说明海量信息增长趋势,系统讲解批量计算与实时计算的分类、技术演进方向,并梳理Google与Apache生态中的代表性框架:MapReduce、Pregel、Percolator、Dremel、Tenzing、Hive,以及GFS/HDFS、BigTable/HBase等文件系统与分布式数据库。课件以信息图式和架构对比方式呈现,帮助读者快速建立分布式计算整体认知,理解离线批处理、准实时更新、大规模查询等典型场景的技术选型,同时涉及数据中心、文件系统、数据库与查询引擎的协同关系,适合课堂教学、自学入门、期末复习或作为课程PPT选题参考。全包仅含1个pptx文件,压缩包大小840KB,便于下载与随堂使用;目前已有106人学习下载。
1. 云计算之分布式计算:先从「任务怎么拆」而不是「框架怎么选」开始
很多人翻到「云计算之分布式计算」这类标题,第一反应是「分布式计算不就是多台机器一起算吗」。真上手才发现,把单机上能跑的程序原样丢到多台机器上,几乎必翻车:有的节点闲着等结果,有的节点CPU飙满卡死,数据对不上,重试越重试越乱。这一篇就把分布式计算从PPT概念拆成能落地、能排错、能运维的真实方案。适合两类人:一类是刚接触分布式计算、想找一条可靠路径入门的开发者;另一类是已经在跑分布式任务、遇到问题只能靠重启恢复的云计算运维工程师——这篇文章至少能让你少走几次弯路。以下内容绝大多数都能在笔记本上复现,不需要一上来就租集群。
2. 分布式计算的三个底层决策:任务划分、数据分片与一致性取舍
2.1 先回答三个问题:任务、数据、状态各自归谁
分布式计算的核心不是「并行」而是「划分」。同样是统计一批日志里关键词出现次数,单机版是一个进程从头读到尾;分布式版要把数据切成块、把每个块派给不同节点、最后把各节点的局部结果合并成全局结果。这个过程中必须做三个决策:任务怎么拆、数据怎么分、共享状态怎么处理。
任务拆分有两条常见路线。第一条是数据并行:所有节点执行同一种处理逻辑,各算各的数据块,典型的场景是日志清洗、图片转码、词频统计。第二条是任务并行:把一个完整流程拆成多个不同的子任务,节点之间形成流水线,典型的场景是ETL里的抽取、清洗、加载三段式。单机没有这个概念,因为一个进程按顺序做完了;分布式环境下,任务拆分的粒度直接决定调度效率。我的建议是:第一次做分布式,别想任务并行,先把数据并行跑通。数据并行的任务边界清晰、失败恢复简单,是分布式计算里性价比最高的第一步。
数据划分决定了任务能不能平均执行。常见分片策略有哈希分片、范围分片和取模分片。哈希分片把数据的某个字段做哈希后再取模,适合数据分布未知的情况;范围分片按时间戳或ID区间切,适合按时间维度查询的业务数据,但容易产生热点分片;取模分片最简单,唯一要记住的是节点数变化时必须全局重分。分片不均匀时,任务会倾斜到一个节点上,其他节点进入空等状态,最后总耗时等于那个最慢节点的耗时——这就是「分布式计算比单机还慢」的第一大元凶。
共享状态是最容易被忽视的一环。如果多个节点要同时写同一个文件、同一张数据库表或同一个内存变量,就引出了分布式锁、事务、幂等写入一整串复杂度。做第一步时,我建议你硬性要求自己做无状态任务:每个节点只处理自己分到的数据,产出独立结果文件或独立消息,需要合并的部分最后由调度方统一来做。等无状态版本跑通了,再引入Redis或数据库这类外部组件来处理有限状态。
2.2 分片策略:Hash、Range 与倾斜的现实世界
分片策略不能拍脑袋,要看你手里的数据长什么样。拿用户ID举例:Range分片把ID 1~10000分给节点A,10001~20000分给节点B;如果ID是自增的,新用户永远落在后段节点上,后段节点很快就变成热点。Hash分片把用户ID做哈希后取模,数据相对均衡,但如果某个大客户贡献了全站一半的日志量,哈希分片照样挡不住倾斜。更现实的方案是按业务维度分片,比如按城市、按渠道、按用户ID段,每个维度单独评估数据量,宁可人工调,也别迷信哈希。
处理倾斜有一个常见做法叫「二次分发」。第一轮按业务维度分片,统计每个分片的数据量,把数据量异常大的分片再切成更细的子分片,按子任务重新调度。这多出来一轮调度开销,但换来的是执行时间的大幅缩减。在运维侧看,任务倾斜的直接表现是某个节点CPU居高不下、其他节点几乎空闲、任务总耗时接近热点节点的耗时。如果你在监控面板上看到这种情况,先不要怀疑代码有bug,先检查数据分片是否均匀。
数据并行和任务并行也可以混用。一个计算量集中在少数几条数据上的任务,可以先用任务并行把大任务拆开,再用数据并行处理其中耗时的分支。这套玩法的落地复杂度会明显上升,少说也要多维护一套任务依赖关系。作为第一步,我始终建议把分片方案单独抽成一个函数,输入是原始数据列表,输出是分片结果列表。这样后期想调整分片策略,只需要改一个函数,不用动分布式框架的配置。
2.3 一致性取舍:离线任务里以最终一致为主
分布式计算必须做一致性决策。在线服务要求强一致,用户下单后立刻能看到订单状态;离线计算任务则普遍采用「先算完、再归并、以最后一次结果为准」的最终一致。原因不复杂:离线任务处理的是只读快照,节点之间不需要实时同步状态,只要保证每个分片至少算一次、结果可以合并,最后结果就是正确的。同一批数据被两个节点重复计算,只要不把重复结果同时写入最终产出,就不会造成损失。
实际操作中,我会把一致性交给外部存储处理。一种做法是用消息队列收集各节点的计算结果,由消费者做幂等写入;另一种是每个节点写独立的结果文件,最后用一个归并任务做合并。我在项目里一定会为每个任务生成唯一的job_id,所有中间结果文件名都带job_id和分片号,比如output/{job_id}/part-{shard_id}.json。这样即使某个分片任务被重试,新写入的临时文件和旧文件也不会互相覆盖,归并阶段按job_id筛选即可。这个习惯在重试频繁的线上环境里价值极高,可以挡住一整类数据重复问题。
如果你要在分布式任务里做全局计数或全局去重,比如统计全量用户数,就要考虑引入外部存储加分布式锁。但分布式锁本身又是一个复杂度来源,比如锁超时、锁续期、脑裂。做得不好,锁会成为全系统最不稳定的点。我的血泪经验是:第一版系统尽量避免强一致的全局操作,把它拆成「先各自算局部结果,再在归并阶段串行处理」,大部分离线场景根本不缺那几秒的实时性。
3. 用 Ray 在本地跑通第一套分布式任务:最小命令与三个必调参数
3.1 装环境与最小代码:把词频统计拆成可并行算子
理论上讲清楚了,下面用Ray跑一套最小示例。Ray是目前上手成本最低的分布式计算框架之一,它把多机调度、对象存储和任务重试封装成了Python装饰器,适合先在本地模拟多节点执行。用到的工具链只有Python和pip,不需要买服务器。
pip install ray安装完成后,先做一个简单的词频统计任务。注意这里的写法刻意用了函数式风格,因为分布式环境里,任务被调度到另一台机器上时需要能被序列化,闭包和类方法在传递时会踩坑。
import ray from collections import Counter ray.init( num_cpus=4, object_store_memory=1024 * 1024 * 512, # 512MB ignore_reinit_error=True ) @ray.remote(max_retries=2) def count_words(text: str) -> dict: return dict(Counter(text.split())) if __name__ == "__main__": with open("sample.log", "r", encoding="utf-8") as f: raw_text = f.read() chunk_size = 8192 chunks = [raw_text[i:i + chunk_size] for i in range(0, len(raw_text), chunk_size)] futures = [count_words.remote(c) for c in chunks] partials = ray.get(futures) total = {} for part in partials: for word, cnt in part.items(): total[word] = total.get(word, 0) + cnt print(total)这段代码的逻辑是:先把原始文本切成固定大小的块,每块作为一个独立任务发给Ray调度器;Ray把任务分配到不同的CPU核心上并行执行,这里在本地跑,4个核心模拟了4个节点的效果;ray.get(futures)会阻塞等待全部任务返回,返回后逐一把分片的词频合并进全局字典。chunk_size是分块大小,它对执行效率有直接影响——太大时并发度低,太小则任务数和调度开销上升,8192到65536之间通常是安全区间。
3.2 必调参数:num_cpus、object_store_memory 与 max_retries
ray.init里最值得关注的是num_cpus。它不是服务器物理核数,而是Ray调度器认为集群拥有的并行槽位数。如果你的机器是4核8线程,设置num_cpus=4比设置8更合理,因为Ray的任务不仅消耗CPU,还有对象序列化和网络传输的开销;过高的并行度会让任务在CPU切换上损耗时间。我自己一般会按物理核数的一半起步,观察任务耗时曲线再往上加。
object_store_memory是Ray对象存储的内存上限。ray.get拉取的结果会先写入这块共享内存,如果预估单批结果总量超过这个值,Ray会开始把对象刷到磁盘,性能瞬间掉一个量级。512MB适合小规模实验,真实业务里我通常按「所有结果对象总大小的1.5倍」来设置,宁多勿少。
max_retries是任务的自动重试次数。要特别注意的是,Ray默认会对失败任务做重试,但重试只保证「任务重新执行」,不保证「任务结果幂等」。如果你的任务里带了写数据库、发消息这类副作用,必须把重试次数调成0,或者把副作用改写成幂等操作后再重试。这是我踩过的最深的一个坑:一个同步任务在节点执行到一半崩溃,Ray自动重试后,数据库里出现两条重复记录。
3.3 除了 Colab,还有哪里能免费跑分布式计算
很多人在Colab上写过深度学习脚本,但Colab免费版单实例规格有限,跑分布式任务很憋屈。除了它,常见的免费云计算路径有这么几条:各家云厂商的新用户免费试用额度,通常有1~3个月的时限,资源规格也够跑中型分布式实验;有条件的话自己攒一台多核二手工作站,装个Linux用Ray跑本地多进程,实践效果和真实集群非常接近;另外各类教育云平台也会提供短期的计算资源,具体时效以平台公告为准。我一般会建议把本地多进程方案当作起点,原因很朴素:本地不用考虑网络延迟、安全组、数据上传这些外部因素,可以专心验证任务逻辑本身。
从本地切到云端时要记住一个变化:ray.init()默认启动的是单机多进程模式,在云主机上也只是把这台机器的并发用满;真正的多机部署需要指定ray.init(address="auto")并先启动一个Ray集群头节点。这一步的复杂度在于网络配置和安全凭证,逻辑代码基本可以原样复用。上云前先确认一下数据缓存的位置,云端节点之间的数据传输走内网流量,而你本地上传数据到云端走的是公网流量,成本差异很大。
4. 分布式计算的避坑清单:节点假死、任务倾斜与无效重试
4.1 节点无响应被误杀:心跳超时的陷阱
现象:某节点在执行任务过程中没有任何报错,但整个任务的日志里出现了「节点心跳丢失,标记为死亡」的记录,节点上的任务被强制迁移到其他节点重跑,耗时翻倍。
原因:分布式框架用周期性的心跳判断节点是否存活。心跳超时设置得太短,比如默认3秒或5秒,一旦节点因为GC停顿、IO等待或网络瞬时抖动,心跳没能及时发出,框架就把该节点判定为死亡并触发任务迁移。迁移本身没问题,但被误杀会导致大量已完成的任务重算。
解决:在框架配置里调整心跳超时和心跳间隔。以Ray为例,ray.init时可以通过namespace和runtime_env调整调度器的健康检查策略,更直接的做法是在集群模式下修改头节点的环境变量,把心跳超时从默认值调大到10秒以上。同时要区分「进程假死」和「节点宕机」:进程假死时CPU占用通常正常或偏高,节点宕机时网络直接不通。监控里同时看CPU和网络指标,能快速判断是哪个方向。
4.2 任务倾斜:数据分片看起来均匀但执行时间差5倍
现象:任务总耗时远大于预期,点开任务列表发现一个节点跑了40分钟,其余节点5分钟就跑完了。
原因:之前提到过,Range分片遇到自增ID时,后半段数据量都会压到最后一个节点;或者业务数据里存在少量超大文本,它们恰好落在同一个分片。
解决:先给所有分片记一个预估耗时,统计每个分片的处理时间,找出异常耗时的分片。然后对热点分片做二次切割,把一个大分片拆成多个小分片,重新进入调度队列。如果在代码里发现分片函数是按文件大小切的,建议改成按行数或按事件数切,这样更接近真实计算耗时。倾斜是分布式系统里最典型的「黑匣子」问题,表面看数据量均匀,实际上每条数据的处理成本差异巨大,动手之前先做一个抽样估算。
4.3 重试不可重入任务导致重复数据
现象:任务第一次失败后自动重试,这次成功了,但最终结果里出现了重复的记录;把重试次数改为0后,重复问题消失。
原因:任务函数里有非幂等操作,比如「如果记录不存在则插入」「把本地文件追加写到一个共享文件」「自增一个外部计数器的值」。第一次任务可能已经做了一半,重试时再做一遍,副作用就重复了。
解决:把任务函数改写成幂等设计。常见做法是任务开头检查结果是否已存在,存在则直接返回;或者所有副作用写入都带job_id和分片号作为唯一键,写入前先检查唯一键。更严格的方案是开启框架的「恰好一次」语义,但这通常需要配套的消息队列和状态存储。对大多数入门项目,保证「每个分片由同一个任务执行、结果可安全覆盖」就够了。
4.4 小任务太多,调度开销比计算还大
现象:把数据切成很多小块后,任务总量从1000涨到100000,执行时间反而变长;节点数增加也没有改善。
原因:每个分布式任务都有调度、序列化、网络传输的固定开销。当单任务执行时间在几毫秒级别时,这些开销远大于计算本身,增加节点只会让调度器更忙。这类问题在日志清洗类任务里特别常见——每条日志单独处理,每条任务执行成本极低。
解决:把多个小数据块合并成一个批次任务,让每个任务处理一批数据而不是一条。批次大小要根据单条数据的平均耗时估算,目标是把单任务耗时做到1秒以上,这样调度开销占比会降到可接受区间。批量引入后还要注意内存占用,单个批次过大会把节点内存打爆,建议从每条任务处理1000条数据起步,观察内存曲线再调整。
5. 从运维视角看分布式计算:日志、指标与扩容水位线
5.1 日志里必备的四个字段:job_id、节点、分片号、耗时
分布式系统的日志最怕的是无法串联。一个任务崩了,你看到的可能只是某个节点上的某条异常堆栈,但不知道它属于哪个job、处理的是哪个分片、为什么重试了两次。我自己维护的分布式任务,日志统一输出成一行JSON,至少包含四个字段:job_id标识一次整体任务;node_id标识运行节点;shard_id标识数据分片;duration_ms记录任务执行耗时。有了这四个字段,排查问题时可以直接按job_id过滤,按node_id看分布,按shard_id对比耗时,第一步就把问题范围缩到最小。
日志级别也要想清楚。业务日志用INFO,数据分片开始和结束必须打点,异常堆栈用ERROR,分片数量、数据量大小这类信息放DEBUG。最怕的是把大量敏感业务数据打进INFO日志,既占磁盘又容易泄露。所有日志统一走标准输出到日志采集系统,不要各节点自己写本地文件,否则排查时你还得一台一台机器登录去看。
5.2 用指标判断系统健康:队列积压、任务耗时与节点负载
指标是分布式系统健康度的仪表盘。我日常重点看三个指标:队列积压量、任务耗时分位数、节点负载。
队列积压量反映调度是否顺畅。如果任务持续产生,消费速度跟不上,积压量会线性上涨,此时系统在累积「隐形超时」。任务耗时比平均值更重要的是p95和p99,因为平均值掩盖了慢任务的存在;p99环比上涨30%且没有对应数据量增长,通常意味着某个节点出了问题或分片倾斜。节点负载看CPU和内存两个维度,CPU中位数长期超过80%建议扩容,内存则要看有没有OOMKilled事件。
| 指标 | 正常区间 | 预警线 | 建议动作 |
|---|---|---|---|
| 队列积压量 | 稳态波动 | 持续上升超过容量三成 | 先查消费端瓶颈,再加并发 |
| 任务p95耗时 | 无明显增长 | 环比上升超过30% | 查慢任务和分片分布 |
| 节点CPU中位数 | 30%~70% | 单节点持续95%以上 | 优先排查倾斜,再考虑扩容 |
5.3 扩容水位线:什么时候加节点,加到多少
扩容是运维里最容易被情绪带跑的动作。任务慢了就加节点,加完发现没快多少,还多付了机器钱。扩容的决策依据应该是:在当前负载下,任务的p95耗时里有超过一半是在等资源,而不是在算数据。判断方法是压测——逐步增加并行度,观察p95耗时是否线性下降;如果增加并行度后耗时几乎不变,说明瓶颈不在计算资源,而在数据读取或外部服务调用。
合理的扩容水位线是:节点CPU中位数超过70%、任务队列积压持续增长、并且p95耗时已经达到业务容忍上限。三项同时满足才扩容。扩容幅度按当前并发度的50%起步,加到p95耗时不再显著下降为止。我见过不少团队把8个节点加到32个节点,结果耗时只缩短了20%,原因就是瓶颈在数据库连接数上。先把数据分片、批次大小、外部依赖这些变量排干净,再动节点数量,才是可靠的路径。
6. 从玩具到可运维的分布式计算:先问数据怎么分
前面五章走完,你已经能用一套框架在本地模拟分布式执行、处理分片倾斜、优化重试策略、看懂调度器的健康指标。但落到真实业务之前,还有一个验证动作不能跳过:在小数据集上,拿单机版的输出和分布式版的输出做逐字段对比。做法是取一份带标准答案的小样本,分别用单机脚本和分布式任务各跑一遍,用哈希对比结果文件是否完全一致。这个动作能提前暴露分片边界问题、重复计算问题和顺序问题,比你在大数据量上反复调试效率高得多。
下一步进阶方向是引入真正的资源调度器,把任务从「在代码里定义」升级到「在集群里排队执行」。常见的路径是给Ray套上Kubernetes,由Kubernetes负责节点生命周期,Ray负责任务调度;或者改用完全托管的云计算服务,把集群管理交给云厂商。这一步的收益是节点可以按负载自动伸缩,代价是你得理解容器镜像、健康探针、资源配额这一整套新概念。
我自己现在拿到任何分布式任务,第一步永远是问「数据怎么分」,而不是「用哪个框架」。分片不均匀,再强的调度器和再多的节点都在陪跑;分片设计好了,框架只是执行工具。带着job_id规范打日志,用p95而不是平均值做告警,重试前先想幂等,做到这三条,你的分布式系统就已经超过不少线上项目的可靠度了。希望帮到你。
本文还有配套的精品资源,点击获取