搞大数据开发的这几年,被问得最多的一个问题就是:刚接触大数据,到底先学什么?我的答案一直没变过——Hadoop 打底,Python 上手,中间用 PySpark 把它们串联起来。这个组合既是入门的最佳路径,也是很多公司生产环境的真实标配:Hadoop 负责分布式存储和资源调度,Python 负责分析效率和开发体验,而 PySpark 让你的 Python 代码能跑在分布式集群上,处理几个 GB 甚至几十个 TB 的数据都没问题。这篇文章就把我这些年关于 Hadoop 与 Python 的实战经验做一次完整整理,围绕“如何用 PySpark 高效完成大数据处理”这条主线,把环境搭建、核心原理、代码实操、性能调优和常见坑位全部串起来。适合三类人看:一是刚开始学大数据、想跑通第一个 Hadoop 环境的学生;二是用 Python 做数据分析、但 Pandas 已经明显带不动海量数据、迫切需要切换到分布式方案的工程师;三是准备大数据岗位面试、想快速把 Hadoop 和 PySpark 核心知识点梳理成体系的人。我尽量不讲教科书套话,所有命令和代码都是跑过之后才敢贴出来的版本。
1. 项目整体思路与技术选型
1.1 单机处理为什么撑不住
先说个很现实的场景。以前很多数据分析师习惯用 Excel 处理百万行数据,再大一点就上 Pandas。Pandas 读一个 500MB 的 CSV,内存经常吃掉 3 到 4 个 GB,做个 groupby 聚合往往要等几十秒;如果还要跟另外十几份文件做 join,计算量一上去,风扇就开始狂转,机器卡到怀疑人生。这不是你代码写得不好,也不是机器不够好,而是单机方案的本质瓶颈摆在那里:内存大小和 CPU 核心数都有物理上限,你再怎么优化,也绕不开单机资源的天花板。
这时候有两条路可以走:一是升级机器配置,换更大的内存、更多的核心,但价格是近乎指数上涨的;二是横向扩展,用很多台普通机器组成一个集群,把数据分散存储、把计算分散执行。大数据领域主流方案明显选的是第二条路,Hadoop 就是从这条路上长出来的生态。
1.2 为什么选 Hadoop 而不是某类数据库
总有人问:现在 ClickHouse、Doris 这类分布式数据库不是也很火吗,为什么还要学 Hadoop?我把这个事说透。分布式数据库解决的是特定场景下的高速查询,比如 OLAP 报表分析,它的存算一体架构确实简单好用。但 Hadoop 的独特之处在于,HDFS 是一个通用的分布式文件系统,你可以把任意格式的数据都放进去,不需要预先定义表结构;YARN 则提供一个通用的资源调度层,上面可以跑 MapReduce、Spark、Flink 等多种计算引擎。
也就是说,Hadoop 解决的不仅仅是“数据怎么存”,还有“怎么让一堆计算框架共享同一批机器资源”的问题。对数据处理链路长、数据来源杂、格式乱的场景——比如一堆日志要先清洗、再关联、再算出指标、最后供下游报表使用——HDFS + YARN + 计算引擎的组合明显更合适。而且 Hadoop 生态的学习路径非常清晰,你把 HDFS 和 YARN 搞透了,后面再接触云上的大数据产品,核心概念几乎都能一一对应,学习曲线会平滑很多。
1.3 计算引擎为什么选 Spark 和 PySpark
Hadoop 解决了“数据怎么存、资源怎么分”的问题,Spark 则解决“数据怎么算得快”的问题。传统 MapReduce 每一轮计算都要把中间结果写到磁盘,遇到迭代计算场景,效率低到让人抓狂;Spark 基于内存计算,把中间结果尽量留在内存里,性能直接提升一个量级。而 PySpark 是 Spark 的 Python API,对这个时代的主力军 Python 开发者来说,学习成本低很多。
我在实际项目里做技术选型时有过一个很直观的体会:同一个数据清洗任务,用 Java 写 MapReduce 可能要写两百行代码,改用 PySpark DataFrame 接口十几行就能搞定,而且跑得还更快。这背后并不是说 PySpark 比 MapReduce 本身强多少,而是 Spark 引擎做了大量优化(比如 Catalyst 优化器和 Tungsten 执行引擎),加上 DataFrame 这种高级 API 帮你省去手动优化的工作量。所以结论很清楚:Hadoop 提供底层存储和调度,Spark 提供高性能计算,Python 提供开发效率,三者结合起来刚好覆盖大数据处理的核心链路。
2. 从零搭建 Hadoop + PySpark 环境
2.1 版本选型:先把坑踩在前面
搭建环境第一步不是下载安装包,而是确定版本组合。这里我强烈建议直接照着一个已验证过的组合来,省得自己瞎试。我这边稳定使用的一套是:操作系统 Ubuntu 20.04 LTS(64 位)、JDK 1.8(最新 Hadoop 3.3.x 也兼容 JDK 11,但 8 最稳)、Hadoop 3.3.6、Python 3.10(PySpark 官方目前完整支持 3.8 到 3.10)、PySpark 3.5.x。
很多人一上来就装最新版 Hadoop,结果跟 Java 版本不兼容,或者跟操作系统的 glibc 版本冲突,启动时各种莫名其妙的问题。别问我是怎么知道的,第一套环境就是被版本组合折磨到崩溃的。记住一个原则:大数据组件追求的是稳定组合,不是最新版本。Hadoop 官方文档里有一张 Java 兼容性表格,安装之前务必先核对一下。
2.2 Java 与 SSH 基础配置
Hadoop 的启动脚本依赖 SSH 来做免密登录,所以 Java 之后要先配好 SSH。这个过程不复杂,但顺序别搞错。先安装 JDK:
sudo apt update sudo apt install -y openjdk-8-jdk java -version确认输出里能看到1.8.0_xxx之类的结果就行。然后生成 SSH 密钥并配置本地免密:
ssh-keygen -t rsa -P '' -f ~/.ssh/id_rsa cat ~/.ssh/id_rsa.pub >> ~/.ssh/authorized_keys chmod 600 ~/.ssh/authorized_keys ssh localhost如果你执行ssh localhost之后能直接进入 shell 而不用输密码,说明免密配置成功。这一步在伪分布式模式里几乎是必须的,因为 Hadoop 的守护进程会通过 SSH 连接到本机启动不同角色。
2.3 Hadoop 伪分布式配置实战
伪分布式模式就是在单台机器上模拟完整的 Hadoop 集群,NameNode、DataNode、ResourceManager、NodeManager 都跑在这台机器上。这种方式做学习和开发验证再合适不过,也是从零开始理解 Hadoop 集群架构的捷径。
下载解压 Hadoop,并配置环境变量:
wget https://archive.apache.org/dist/hadoop/common/hadoop-3.3.6/hadoop-3.3.6.tar.gz sudo tar -xzf hadoop-3.3.6.tar.gz -C /usr/local/ sudo mv /usr/local/hadoop-3.3.6 /usr/local/hadoop sudo chown -R $(whoami) /usr/local/hadoop然后编辑/etc/profile,加入 Hadoop 环境变量:
export HADOOP_HOME=/usr/local/hadoop export PATH=$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin接着进入 Hadoop 配置目录,修改四个核心配置文件。第一个是core-site.xml,指定 NameNode 的地址和临时目录。这里要提前建一个目录,比如/home/yourname/hadoop_tmp,并且权限别给错:
<configuration> <property> <name>fs.defaultFS</name> <value>hdfs://localhost:9000</value> </property> <property> <name>hadoop.tmp.dir</name> <value>/home/yourname/hadoop_tmp</value> </property> </configuration>第二个是hdfs-site.xml。伪分布式模式建议把副本数设成 1,否则三副本机制在单机上会报磁盘空间不足:
<configuration> <property> <name>dfs.replication</name> <value>1</value> </property> <property> <name>dfs.namenode.name.dir</name> <value>/home/yourname/hadoop_tmp/name</value> </property> <property> <name>dfs.datanode.data.dir</name> <value>/home/yourname/hadoop_tmp/data</value> </property> </configuration>第三个是mapred-site.xml,这个文件在模板目录里叫mapred-site.xml.template,需要先复制一份出来,然后指定用 YARN 做资源调度框架:
<configuration> <property> <name>mapreduce.framework.name</name> <value>yarn</value> </property> </configuration>第四个是yarn-site.xml,关键是把 auxiliary 服务配置成 mapreduce_shuffle,不然 MapReduce 任务跑不起来:
<configuration> <property> <name>yarn.nodemanager.aux-services</name> <value>mapreduce_shuffle</value> </property> <property> <name>yarn.nodemanager.aux-services.mapreduce_shuffle.class</name> <value>org.apache.hadoop.mapred.ShuffleHandler</value> </property> </configuration>配置完成后,第一次启动前必须格式化 NameNode,这个操作相当于给 HDFS 文件系统做一次初始化:
hdfs namenode -format start-dfs.sh start-yarn.sh启动后执行jps,如果能看到NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager这几个进程,说明环境基本就绪。再用浏览器打开http://localhost:9870看 HDFS 的 Web UI,打开http://localhost:8088看 YARN 的资源页面,两个页面都能正常显示,Hadoop 伪分布式就算搭完了。
2.4 安装 Python 和 PySpark
Hadoop 起来之后,接下来配 Python 侧。我建议直接用虚拟环境管理依赖,避免污染系统 Python。先装好 Python 3.10 和 pip,然后创建虚拟环境并安装 PySpark:
python3 -m venv ~/pyspark_env source ~/pyspark_env/bin/activate pip install pyspark==3.5.1注意 PySpark 安装包挺大的,大概一两百 MB,包含完整的 Spark 二进制,所以不用额外下载 Spark,直接 pip 就够。装完之后写个最简单的验证脚本:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .master("local[*]") \ .appName("EnvTest") \ .getOrCreate() df = spark.range(0, 10) df.show() spark.stop()如果你看到输出 0 到 9 这十行数据,说明 PySpark 能正常读取本地 Spark 运行环境。这里local[*]是让 Spark 用本机所有 CPU 核心跑在本地模式,学习阶段完全够用。如果后续想连接 Hadoop 集群,把master改成yarn,并把HADOOP_HOME环境变量指好就行。
3. Hadoop 核心机制:光会启动远远不够
3.1 HDFS 的存储设计逻辑
很多教程让你把 Hadoop 启动就算了,但我建议至少要理解 HDFS 为什么这么设计,不然写 PySpark 的时候,遇到数据读写的性能问题会一头雾水。HDFS 的核心思想是把大文件切分成固定大小的块(默认 128MB),然后分散存储到集群的不同 DataNode 上。每个块默认保存三个副本,副本放在不同机器上,这样任何一台机器宕机都不会丢数据。
NameNode 是 HDFS 的“大脑”,负责维护文件系统的目录树和每个块的元数据;DataNode 是“肌肉”,真正存数据。你在 shell 里执行hdfs dfs -put,文件会被切成块并分发到 DataNode;执行hdfs dfs -cat,客户端会先问 NameNode 要元数据,然后直接从对应的 DataNode 读数据,不经过 NameNode 转发。这个设计保证了数据读写不会被单点瓶颈卡死,但也带来一个注意事项:NameNode 是整个集群的“单点”,所以在生产环境里,NameNode 的高可用配置是重中之重。
3.2 YARN 的资源调度机制
YARN 解决的问题是:一个集群里有多种计算任务(MapReduce、Spark、Flink),它们怎样才能公平、高效地共享同一批机器的 CPU 和内存。YARN 里有两个核心角色:ResourceManager(RM)负责全局资源分配,NodeManager(NM)负责管理单台机器上的资源。
当你要跑一个 PySpark 任务时,客户端会向 ResourceManager 提交 Application;RM 找到合适的 NodeManager,启动一个 ApplicationMaster 负责协调这个任务;ApplicationMaster 再向 RM 申请容器(Container),然后在容器里启动 Executor 进程,真正执行计算。这套机制用生活类比来解释就是:ResourceManager 是酒店前台,NodeManager 是楼层服务员,ApplicationMaster 是会议的会务组,会务组找前台要会议室,前台协调楼层服务员来布置,会议才能顺利开起来。
3.3 为什么实际写代码很少直接碰 MapReduce
MapReduce 是 Hadoop 最早的计算引擎,思想非常经典:Map 阶段把数据拆分成键值对,Shuffle 阶段按 key 分组,Reduce 阶段做聚合。但它的硬伤也很明显——每个阶段的中间结果都要落盘,复杂任务可能有几十个 MapReduce 串起来,每次落盘都是巨大的 I/O 开销。Spark 之所以快,核心就在于把中间结果尽量保留在内存里,加上 DAG 调度引擎,能自动合并多个计算步骤,避免频繁落盘。
所以实际开发中,如果做离线批处理,大家更倾向于直接用 Spark 而不是裸写 MapReduce。PySpark 就是在这一层给 Python 开发者开的一扇窗户。你不需要知道 MapReduce 的每个细节,但你要明白 PySpark 底层的 Shuffle 机制和 MapReduce 的 Shuffle 是同源的,理解了这个,后面调优时候你就知道哪些操作会触发 Shuffle、为什么会慢。
4. PySpark 核心实操:从 RDD 到 DataFrame 再到完整任务
4.1 RDD 与 DataFrame:底层和界面
PySpark 里有两个层次的东西,RDD 是底层的弹性分布式数据集,DataFrame 是上层的结构化 API。RDD 的优势是灵活,什么都能干,但写起来啰嗦,而且没有自动优化机制;DataFrame 类似 Pandas 里的 DataFrame,但又跑在分布式引擎上,多了一整套 Catalyst 查询优化器。我的经验是:日常业务开发优先用 DataFrame;只有在需要做 RDD 底层操作(比如自定义分区器)的时候,才把数据.rdd转下去处理。
举个例子,给 DataFrame 增加一列,用 RDD 方式要写 map 函数、处理 Row 对象,还得关注序列化问题;用 DataFrame 的withColumn一行就完了,而且优化器会自动帮你做向量化执行。这就是为什么我对初学者只有一句忠告:别从 RDD 开始学,直接学 DataFrame,把 DataFrame 用熟练之后,再回头理解 RDD 会发现它其实很简单。
4.2 DataFrame 高频操作速览
下面这段代码是读取一个模拟的用户访问日志(CSV 格式),做排序、过滤、分组统计,用的都是日常开发频率最高的操作:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, count, avg, sum, desc from pyspark.sql.types import StructType, StructField, StringType, IntegerType spark = SparkSession.builder \ .appName("UserLogAnalysis") \ .config("spark.sql.shuffle.partitions", "4") \ .getOrCreate() schema = StructType([ StructField("user_id", StringType(), True), StructField("service_name", StringType(), True), StructField("request_time", IntegerType(), True), StructField("status", StringType(), True) ]) df = spark.read.option("header", False)\ .schema(schema)\ .csv("file:///home/yourname/logs/*.csv") # 过滤出请求时间大于 100 毫秒的慢请求 slow_df = df.filter(col("request_time") > 100) # 按服务统计慢请求数量和平均耗时 result = slow_df.groupBy("service_name") \ .agg( count("*").alias("slow_count"), avg("request_time").alias("avg_time") ) \ .filter(col("slow_count") >= 10) \ .orderBy(desc("slow_count")) result.show()这里有几个容易忽略的关键点。第一,读取 CSV 时一定要手动指定 schema,不要依赖自动类型推断。自动推断虽然方便,但会额外扫一遍数据,大数据量下开销很大。第二,groupBy后面跟的聚合是宽依赖操作,会触发 Shuffle,所以在本地测试时可以通过config("spark.sql.shuffle.partitions", "4")控制输出分区数量,这个参数在生产环境里尤其重要,后面我还会细讲。第三,filter在聚合前后的含义完全不同。先filter再聚合,过滤的是原始数据;先聚合再filter,过滤的是聚合结果,两个结果可能截然不同,写代码前先想清楚你要的是哪种。
4.3 UDF 自定义函数
DataFrame 自带的内置函数解决大部分场景,但总有一些业务逻辑需要你自己写。这时候就需要 UDF(User Defined Function)。比如要根据请求时间判断性能等级:
from pyspark.sql.functions import udf from pyspark.sql.types import StringType def judge_level(time_ms: int) -> str: if time_ms < 50: return "fast" elif time_ms < 200: return "normal" else: return "slow" judge_level_udf = udf(judge_level, StringType()) df.withColumn("level", judge_level_udf(col("request_time"))) \ .select("user_id", "service_name", "request_time", "level") \ .show()但我要提醒一个特别大的坑:普通的 Python UDF 在 PySpark 里是一次一行调用的,序列化和 Python 解释器开销很大,数据量一大,性能会退化得非常明显。如果你只是在做原型验证,用 UDF 没问题;但生产环境里,能用内置函数或者 Spark SQL 的表达式,就尽量别用 UDF。实在绕不开 UDF,考虑用 Pandas UDF(也叫 Vectorized UDF),它一次批量处理一批数据,性能能提升好几倍:
from pyspark.sql.functions import pandas_udf from pyspark.sql.types import StringType import pandas as pd @pandas_udf(StringType()) def judge_level_pd(time_ms: pd.Series) -> pd.Series: return time_ms.apply(lambda x: "fast" if x < 50 else ("normal" if x < 200 else "slow")) df.withColumn("level", judge_level_pd(col("request_time"))).show()4.4 完整案例:日志清洗与指标统计
把前面内容串起来,做一个更接近生产场景的任务。假设我们现在有几份日志文件,字段包括用户 ID、服务名、请求耗时、状态码,需要完成三件事:清洗掉缺失字段的行,算出每个服务每天的平均耗时和请求量,最后把结果写回 HDFS。
from pyspark.sql import SparkSession from pyspark.sql.functions import col, count, avg, to_date spark = SparkSession.builder \ .appName("LogETL") \ .config("spark.sql.shuffle.partitions", "8") \ .getOrCreate() df = spark.read.csv( "hdfs://localhost:9000/input/logs/", header=True, inferSchema=True ) # 清洗:去掉 user_id 为空或者请求耗时小于 0 的异常数据 clean_df = df.filter( col("user_id").isNotNull() & col("request_time").isNotNull() & (col("request_time") >= 0) ) # 加一列事件日期 clean_df = clean_df.withColumn("event_date", to_date(col("timestamp"))) # 指标计算:按服务名和日期聚合 result = clean_df.groupBy("service_name", "event_date") \ .agg( count("*").alias("req_count"), avg("request_time").alias("avg_time_ms") ) \ .orderBy("event_date", "service_name") # 写回 HDFS,覆写模式 result.write.mode("overwrite").parquet("hdfs://localhost:9000/output/metrics/") spark.stop()这段代码做完之后,你可以执行:
hdfs dfs -ls /output/metrics/看到 parquet 文件和_SUCCESS标记文件,就说明任务确实写到了 HDFS 上。写代码时有两个容易踩的坑值得注意:一是路径前缀,HDFS 路径一定要写hdfs://localhost:9000开头,别跟本地文件系统搞混;二是to_date函数依赖时间戳字段能被正确解析,如果原始日志的时间格式不规范,建议先清洗成标准格式再做转换,否则很容易产生大量空值。
5. 性能调优与问题排查实录
5.1 分区数:不是越多越好
PySpark 性能调优里面,分区数是直接影响并行度的参数。数据分发到各个分区,每个分区由一个 task 处理,所以分区数决定了任务并行度。理论上分区数越多,并行度越高,但分区太多也会带来额外的调度开销和 Shuffle 网络开销,反而变慢。
我通常的经验是:每个分区处理 128MB 到 256MB 数据比较合适。比如 1GB 的数据,分 4 到 8 个分区就够了。如果你发现某个任务 Executor 数量不少,但大多数 Executor 都是闲着的,很可能就是分区数太少;反过来,如果你看到大量小任务在秒级启动、秒级结束,大概率是分区数太多,增加了不必要的调度成本。调节手段无非是repartition()、coalesce()和spark.sql.shuffle.partitions这几个工具,coalesce()只能减少分区,而且不会触发 Shuffle,用于处理后的结果回收很适合;repartition()既可以增加也可以减少,但会触发一次全量 Shuffle,要用在合适的位置。
5.2 数据倾斜:最头疼的老大难问题
我在实战中最常遇到的大坑就是数据倾斜。表现形式很典型:一个任务跑了几十分钟,其他任务早就结束了,就卡在最后一个 task 上慢慢磨。根本原因是数据里某个 key 的值特别多,比如日志里的某个用户 ID 是爬虫攻击,请求量占了全量的 90%,所有数据都倾斜到同一个分区去了。
处理思路主要有四个。第一个思路是过滤异常 key,如果倾斜的 key 本来就不参与核心逻辑,直接在过滤条件里去掉;第二个思路是加盐(salting),把倾斜 key 变成若干个加了随机前缀的新 key,打散到多个分区后再聚合,最后再把前缀去掉聚合一次;第三个思路是改用广播 join,如果一个表很小,就把它广播到每个 Executor 内存里,避免 Shuffle 阶段的倾斜;第四个思路是两阶段聚合,先局部聚合再加全局聚合。
我举个加盐的简略思路:
from pyspark.sql.functions import concat, lit, rand, substring # 给倾斜 key 加随机后缀打散 salted_df = df.withColumn( "salted_key", concat(col("key"), lit("_"), (rand() * 10).cast("int")) ) # 第一阶段:按加盐 key 预聚合 partial = salted_df.groupBy("salted_key").agg(count("*").alias("cnt")) # 第二阶段:还原出原始 key,再做全量聚合 final = partial.withColumn( "original_key", substring(col("salted_key"), 1, 4) ).groupBy("original_key").agg(sum("cnt").alias("total"))这只是最简单的演示,真实场景要根据 key 的长度和格式调整恢复逻辑。但记住核心思想:把热点数据打散,分两步聚合,是处理数据倾斜的通用套路。
5.3 缓存与血缘机制
Spark 的任务天然有“血缘关系”,每一步操作都会记录依赖链条,这样某一步出错了可以从源头重新计算。但这也带来一个副作用:如果某个中间结果要被多个下游任务反复使用,每次都重算一遍,代价极大。
解决办法就是缓存。df.cache()把数据缓存在内存里,df.persist(StorageLevel.MEMORY_AND_DISK)还可以配置内存不够时落盘。我一般会把那种进行过多轮 join 和过滤的中间表做缓存,后续有好几个 DataFrame 都要从这个中间表继续派生。但注意,缓存不是银弹,用一次缓存就要占一份内存,缓存太多反而会撑爆 Executor 内存。用完之后记得调用df.unpersist()释放。
5.4 高频报错与排查速查表
把我在实际调试中遇到的高频报错整理成一张表,方便大家直接对号入座。
| 报错信息 | 常见原因 | 解决办法 |
|---|---|---|
java.net.ConnectException: Connection refused | Hadoop 服务没启动,或端口配置不对 | 检查jps确认 NameNode/DataNode 进程,检查core-site.xml端口 |
Container killed by the ResourceManager | Executor 内存超限 | 调大spark.executor.memory或者减少单个 Executor 的核心数 |
OutOfMemoryError: Java heap space | 数据量太大,Executor 堆内存不够 | 增加分区数、减少缓存数据、调大堆内存 |
FileNotFoundError: input path does not exist | HDFS 路径写错,或文件还没上传 | 用hdfs dfs -ls确认路径,注意hdfs://前缀 |
Cannot connect to the cluster | YARN 集群连接失败 | 检查yarn-site.xml和core-site.xml,确认 ResourceManager 地址 |
IllegalArgumentException: Wrong FS | 把 HDFS 路径和本地路径混用了 | 统一使用明确前缀,别省略 scheme |
除了这六类,我还想说一个排查通法:遇到问题先看 Web UI。YARN 的http://localhost:8088页面能看到任务有没有失败、失败在哪一步,点进去能看到具体日志。Spark 也有一个自己的 Web UI,直接http://localhost:4040,里面能看到每个 stage 的任务执行时间、Shuffle 读写量、内存使用情况。很多性能问题在这两个页面上都是一眼能看出来的。
6. 生产实践与经验补充
6.1 开发环境与生产环境的差异
很多人在本地用local[*]模式跑通了代码,觉得万事大吉,结果一提交到生产集群就挂。原因很简单:开发模式在单机执行,没有网络传输、没有资源竞争、没有权限管控。生产环境里提交 PySpark 任务,通常要改用spark-submit,并指定集群资源:
spark-submit \ --master yarn \ --deploy-mode cluster \ --num-executors 4 \ --executor-cores 2 \ --executor-memory 4g \ your_job.py这里面有几个参数是你必须理解的。--num-executors是启动多少个 Executor 进程,--executor-cores是每个 Executor 用几个 CPU 核心,--executor-memory是每个 Executor 分多少内存。它们直接决定了你的任务能拿到多少集群资源。我曾经做过一次测试,把 Executor 数量从 2 加到 6,一个两小时的任务压缩到四十分钟,这就是并行度的威力。
6.2 数据本地性的讲究
生产环境里跑 PySpark,还有一个经常被忽略的性能因素叫数据本地性。简单说,Spark 在执行任务时,会优先把计算尽量调度到数据所在的那台机器上,这样不用把数据从网络里拖来拖去。如果资源紧张,调度器会让任务跑到离数据比较远的节点上,通过网络拉数据,性能就会差不少。
这个点你平时可能感觉不到,但一旦集群繁忙、资源碎片化严重,数据本地性等级会从NODE_LOCAL降到RACK_LOCAL,甚至ANY,任务耗时一下子就上去了。所以生产环境里,集群资源预留和任务排队策略往往比代码本身更影响性能。如果你发现自己代码怎么写都慢,不妨先看看是不是资源调度层面出了问题。
6.3 面试高频考点速记
最后顺便帮准备面试的读者串一下高频考点。我这些年面试别人和被面试,发现问来问去就是那几个核心点:HDFS 读写流程、YARN 调度流程、Spark 的 RDD 与 DataFrame 区别、宽依赖与窄依赖、Spark Shuffle 机制、数据倾斜处理方案、Spark 的容错与血缘机制。这些概念平时写代码未必都会直接用到,但理解它们能让你在排查问题时更快定位方向。
我的建议很简单:不要死背答案,而是对着自己跑过的任务去理解。比如你刚才跑了一个日志统计任务,中间哪一步触发了 Shuffle?哪个阶段产生了宽依赖?数据倾斜如果发生,你会怎么处理?能把这些问题结合自己的代码回答出来,面试官基本不会为难你。
说起来,写这套东西的时候我又想起了当年第一次把 Hadoop 伪分布式集群从零搭起来、第一次用 PySpark 跑通一个亿级数据聚合任务时的感觉。大数据处理这条路,最大的门槛其实是“第一次”。第一次搭环境、第一次分布式跑通、第一次定位到数据倾斜问题,只要经历了这些“第一次”,后面的路就会顺畅很多。如果你正卡在环境搭建或者概念理解的阶段,别灰心,照着文章里的步骤慢慢试,一定能把这套 Hadoop 加 PySpark 的组合跑起来。