简介:这份《云计算之分布式计算》PPT课件面向云计算、大数据方向的学习者与技术人员,系统梳理分布式计算的核心概念与主流技术体系,帮助读者建立从批量计算到实时计算的完整认知框架。课件围绕分布式计算的定义、批量计算与实时计算的区别、技术演进趋势展开,并对比Google与Apache两套技术栈:文件系统GFS与HDFS、分布式数据库BigTable与HBase、批量计算框架MapReduce、迭代计算框架Pregel与Hama、SQL查询引擎Tenzing与Hive等,同时结合MapReduce工作流程的实例演示,讲解并行化、容错、数据分布与负载均衡等关键问题。资源包内含1个pptx文件,大小约840KB,以图文并茂的幻灯片形式呈现,便于课堂讲解与自学梳理。目前已有106人学习浏览,适合作为云计算课程配套资料或技术入门参考,帮助读者快速掌握分布式计算的知识脉络与典型框架选型思路。
1. 一份 PPT 撑不起分布式计算:从「云计算之分布式计算.pptx」看落地到底缺什么
很多人第一次接触分布式计算,是从一份叫「云计算之分布式计算.pptx」的课件开始的。幻灯片里画着 Master 和 Worker 的框图,讲 MapReduce 怎么把任务切开再合并,讲云平台怎么按需分配算力。看完感觉懂了,一上手就懵——因为 PPT 只告诉你「有这回事」,没告诉你「怎么让它跑起来」。这份课件真正对应的技术方向,是把一个大计算任务拆成若干子任务、分发到多台机器上并行执行、再把结果汇总回来,而云计算提供的是弹性算力底座。适合谁看?适合已经会写单机 Python 脚本、但一遇到数据量翻倍就卡住的后端或数据工程师,也适合云计算运维方向想补齐分布式计算实操的人。接下来我不复述 PPT,而是把这条链路拆成能复现的步骤。
2. 分布式计算到底在算什么:拆任务、分机器、收结果
2.1 从单机瓶颈到分布式拆分
单机跑一个计算任务,瓶颈通常出现在三个地方:CPU 核数不够、内存装不下中间数据、磁盘 IO 成为瓶颈。分布式计算的核心思路不是「让一台机器变快」,而是「让多台机器同时干」。这里有个关键前提:任务必须可拆分。如果一个任务的后一步强依赖前一步的完整结果,那它天然不适合分布式。
常见的可拆分模式有三种。第一种是数据并行,把同一套计算逻辑应用到不同的数据分片上,比如对一亿条日志分别统计 PV。第二种是任务并行,不同子任务之间没有依赖,比如同时训练多个超参组合的模型。第三种是流水线并行,把计算过程分成若干阶段,每个阶段由不同节点处理,数据像流水线一样流过。MapReduce 本质上属于第一种,把 Map 阶段并行化,再用 Reduce 汇总。
理解这一点很重要,因为它决定了你选什么框架。数据并行选 MapReduce 或 Spark,任务并行选 Celery 或 Ray,流水线并行选更偏底层的消息队列加计算节点方案。选错了框架,后面全是坑。
2.2 云计算在这里扮演什么角色
云计算对分布式计算的价值,不是「提供了分布式计算能力」,而是「让分布式计算的门槛降到可以按小时计费」。以前要搭一个十节点的集群,你得先买十台服务器、装系统、配网络、部署调度器,一周过去了。现在在云平台上开十台按量付费的实例,半小时内就能跑起来,跑完就释放。
但这里有个容易被忽略的点:云上的分布式计算,网络开销和单机完全不同。单机里内存访问是纳秒级,云上跨节点通信是毫秒级,差了六个数量级。所以任何在单机上看不出来的通信开销,在分布式环境里都会被放大。这就是为什么很多在本地跑得好好的代码,一上云就慢得离谱。
2.3 最小可跑的分布式任务长什么样
先不急着上 Hadoop 或 Spark,用 Python 的 multiprocessing 在单机模拟分布式拆分逻辑,把「拆任务、并行算、合并结果」这条链路跑通。这是理解后续所有框架的基础。
import multiprocessing as mp import time def worker(chunk): # 模拟一个计算密集型任务:对数据块求和 result = sum(x * x for x in chunk) return result def split_data(data, n): # 把数据均匀切成 n 份 k, m = divmod(len(data), n) return [data[i*k + min(i, m):(i+1)*k + min(i+1, m)] for i in range(n)] if __name__ == "__main__": data = list(range(10_000_000)) n_workers = 4 chunks = split_data(data, n_workers) start = time.time() with mp.Pool(n_workers) as pool: results = pool.map(worker, chunks) total = sum(results) print(f"结果={total}, 耗时={time.time()-start:.2f}s")这段代码的逻辑是:先把一千万个整数切成四份,每份交给一个进程去算平方和,最后把四个结果加起来。split_data里的divmod是为了处理不能整除的情况,保证每份大小尽量均匀。mp.Pool创建进程池,pool.map把任务分发下去并收集结果。
参数方面,n_workers一般设成 CPU 核数,设太大反而因为进程切换变慢。chunk的大小要适中,太小则调度开销占比高,太大则负载不均。这个例子在单机上跑,但它和真正的分布式计算在逻辑结构上是一致的:切分、分发、计算、汇总。区别只在于进程换成了网络节点。
提示:如果你在云主机上跑这段代码,注意实例的 vCPU 数和
n_workers要匹配,否则会出现「开了四进程但只有两核」的假并行。
3. 用云主机搭一个能跑的分布式计算环境
3.1 选型:什么时候用 MapReduce,什么时候用 Spark
这是实操中第一个要做的决策。MapReduce 适合一次性的、磁盘 IO 密集的批处理任务,它的中间结果落盘,容错性好但速度慢。Spark 适合迭代计算和内存密集的任务,中间结果尽量留在内存,速度快但对内存要求高。
判断标准很简单:如果你的任务需要反复扫描同一份数据,比如机器学习训练中的多轮迭代,选 Spark。如果是一次性的 ETL 清洗,选 MapReduce 或直接用云平台的数据流水线服务。如果任务本身不复杂但需要调度大量异构任务,考虑 Ray 或 Celery。
我一般会先用单机 Spark 的 local 模式验证逻辑,确认没问题再上集群。这样能省掉大量「在集群上调试」的时间。
3.2 在云主机上部署 Spark 单机模式
先在一台云主机上把 Spark 跑通,这是最小验证单元。以下步骤在 Ubuntu 22.04 上验证过。
# 安装 Java 运行环境,Spark 依赖 JVM sudo apt update && sudo apt install -y openjdk-17-jre-headless # 下载 Spark(以 3.5.x 为例,具体版本按官方页面选择) wget https://dlcdn.apache.org/spark/spark-3.5.1/spark-3.5.1-bin-hadoop3.tgz tar -xzf spark-3.5.1-bin-hadoop3.tgz sudo mv spark-3.5.1-bin-hadoop3 /opt/spark # 配置环境变量 echo 'export SPARK_HOME=/opt/spark' >> ~/.bashrc echo 'export PATH=$SPARK_HOME/bin:$PATH' >> ~/.bashrc source ~/.bashrc # 验证安装 spark-submit --version安装完成后,用 PySpark 跑一个词频统计,这是分布式计算的「Hello World」。
from pyspark import SparkContext sc = SparkContext("local[4]", "WordCount") # local[4] 表示用 4 个本地线程模拟分布式 text = sc.parallelize([ "cloud computing distributed", "distributed computing scalable", "cloud scalable elastic" ]) counts = (text .flatMap(lambda line: line.split(" ")) .map(lambda word: (word, 1)) .reduceByKey(lambda a, b: a + b) ) for word, count in counts.collect(): print(f"{word}: {count}") sc.stop()flatMap把每行拆成单词,map给每个单词打上计数 1,reduceByKey把相同单词的计数加起来。这三个操作就是 MapReduce 的 Map 和 Reduce 在 Spark 里的表达。local[4]里的 4 是并行度,上集群后换成spark://master:7077就变成真正的分布式执行。
3.3 从单机到集群:加机器时要改什么
单机跑通后,上集群主要改三个地方。第一,SparkContext的 master 地址从local[N]改成集群 master 的地址。第二,数据源从parallelize本地集合改成从对象存储或 HDFS 读取。第三,资源配置从默认改成显式指定 executor 数量和内存。
from pyspark.sql import SparkSession spark = (SparkSession.builder .appName("CloudDistJob") .master("spark://10.0.0.1:7077") # 集群 master 地址 .config("spark.executor.instances", "4") # executor 数量 .config("spark.executor.memory", "4g") # 每个 executor 内存 .config("spark.executor.cores", "2") # 每个 executor 核数 .getOrCreate() ) df = spark.read.parquet("s3a://your-bucket/data/") # 从对象存储读 df.createOrReplaceTempView("logs") result = spark.sql("SELECT level, count(*) FROM logs GROUP BY level") result.show() spark.stop()executor.instances决定并行度上限,executor.memory决定每个节点能缓存多少数据,executor.cores决定每个 executor 能同时跑几个 task。这三个参数要匹配云主机的规格,比如你开的是 4 核 16G 的实例,那executor.cores设 2、executor.memory设 6g 左右比较合理,留一部分内存给系统和缓存。
注意:云上跨可用区通信有额外延迟,尽量把 master 和 worker 放在同一个可用区,否则网络开销会吃掉并行带来的收益。
4. 避坑:分布式计算在云上翻车的五个典型场景
4.1 数据倾斜导致个别节点跑满、其余空闲
现象:任务卡在最后几个 reduce 上,99% 的 task 都完成了,剩下 1% 跑了几个小时。原因:某个 key 对应的数据量远超其他 key,所有这个 key 的数据都被分到同一个 reducer 上。解决:先采样找出热点 key,对它加随机前缀打散,计算完再去掉前缀合并。在 Spark 里可以用repartition加盐处理。
4.2 executor 内存溢出但日志只报 OOM
现象:任务失败,日志里只有java.lang.OutOfMemoryError,看不出是哪一步出的问题。原因:可能是数据缓存太多,也可能是单个 partition 太大。解决:先看 Spark UI 里各 stage 的 shuffle spill 情况,如果 spill 量大就加内存或减少 partition 大小。另外把spark.memory.fraction从默认 0.6 适当调低,给用户代码留更多空间。
4.3 云主机按量计费跑完忘记释放
现象:月底账单比预期高出一大截。原因:集群跑完后没有释放实例,按量付费的机器一直在计费。解决:用脚本在任务结束后自动释放,或者用云平台的定时任务功能设置最大运行时长。我一般会在提交任务时加一个超时参数,超过时间自动 kill 并释放资源。
4.4 网络带宽成为瓶颈而不是 CPU
现象:CPU 利用率一直上不去,但任务就是慢。原因:跨节点 shuffle 的数据量太大,网络带宽被打满。解决:尽量在 map 端做预聚合,减少 shuffle 数据量。在 Spark 里可以用reduceByKey替代groupByKey,因为前者会先在 map 端合并。另外检查云主机的内网带宽规格,有些入门级实例内网带宽很低。
4.5 本地能跑、上云就报序列化错误
现象:本地local模式跑得好好的,一上集群就报NotSerializableException。原因:闭包里引用了不可序列化的对象,比如数据库连接、文件句柄。解决:把不可序列化的对象在 executor 内部创建,不要通过闭包传进去。用foreachPartition替代map可以在每个 partition 内创建一次连接,既避免序列化问题又减少连接数。
5. 把分布式计算用对:从能跑到跑得省的三个进阶习惯
第一个习惯是先用小数据集验证逻辑,再放大数据量。我见过太多人直接拿全量数据上集群,跑了一天发现逻辑写错了。正确做法是取千分之一的数据在本地跑通,确认结果正确后再上集群。这样调试成本从小时级降到秒级。
第二个习惯是给每个任务设资源上限和超时。云上最怕的不是任务失败,而是任务卡住一直占着资源。在 Spark 提交时加上--conf spark.executor.memoryOverhead=1g和超时参数,在 Celery 里设task_time_limit,都是花五分钟配置、省几小时账单的事。
第三个习惯是看 Spark UI 而不是猜。任务慢的时候,先打开 Spark UI 看哪个 stage 耗时长、哪个 task 数据量大、shuffle 读写多少。这些数字比任何猜测都准。我一般会重点关注 shuffle read/write 的大小和 task 的 duration 分布,如果某个 task 明显比其他长,那就是数据倾斜。
| 检查项 | 正常范围 | 异常信号 | 处理方向 |
|---|---|---|---|
| shuffle spill | 接近 0 | 持续大量 spill | 加内存或减小 partition |
| task duration 分布 | 相对均匀 | 个别 task 极长 | 检查数据倾斜 |
| CPU 利用率 | 60%-80% | 长期低于 30% | 检查网络或 IO 瓶颈 |
| executor 内存使用 | 稳定 | 持续增长后 OOM | 检查缓存或内存泄漏 |
最后说一个我自己的教训。早期做分布式计算时,我总觉得「机器越多越快」,有一次开了二十台实例跑一个并不大的任务,结果因为调度开销和网络通信,比单机还慢。后来才明白,分布式计算的收益来自「任务可并行且单机装不下」,不是来自「机器多」。先判断任务值不值得分布式,再决定开几台机器,这个顺序不能反。希望帮到你。
本文还有配套的精品资源,点击获取