ARTICLE DETAIL

资讯详情

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

Hadoop与Python实战:用PySpark构建大数据处理全流程

Hadoop与Python实战:用PySpark构建大数据处理全流程 搞大数据开发的这几年被问得最多的一个问题就是刚接触大数据到底先学什么我的答案一直没变过——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 和 PySparkHadoop 解决了“数据怎么存、资源怎么分”的问题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 LTS64 位、JDK 1.8最新 Hadoop 3.3.x 也兼容 JDK 11但 8 最稳、Hadoop 3.3.6、Python 3.10PySpark 官方目前完整支持 3.8 到 3.10、PySpark 3.5.x。很多人一上来就装最新版 Hadoop结果跟 Java 版本不兼容或者跟操作系统的 glibc 版本冲突启动时各种莫名其妙的问题。别问我是怎么知道的第一套环境就是被版本组合折磨到崩溃的。记住一个原则大数据组件追求的是稳定组合不是最新版本。Hadoop 官方文档里有一张 Java 兼容性表格安装之前务必先核对一下。2.2 Java 与 SSH 基础配置Hadoop 的启动脚本依赖 SSH 来做免密登录所以 Java 之后要先配好 SSH。这个过程不复杂但顺序别搞错。先安装 JDKsudo 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 namefs.defaultFS/name valuehdfs://localhost:9000/value /property property namehadoop.tmp.dir/name value/home/yourname/hadoop_tmp/value /property /configuration第二个是hdfs-site.xml。伪分布式模式建议把副本数设成 1否则三副本机制在单机上会报磁盘空间不足configuration property namedfs.replication/name value1/value /property property namedfs.namenode.name.dir/name value/home/yourname/hadoop_tmp/name/value /property property namedfs.datanode.data.dir/name value/home/yourname/hadoop_tmp/data/value /property /configuration第三个是mapred-site.xml这个文件在模板目录里叫mapred-site.xml.template需要先复制一份出来然后指定用 YARN 做资源调度框架configuration property namemapreduce.framework.name/name valueyarn/value /property /configuration第四个是yarn-site.xml关键是把 auxiliary 服务配置成 mapreduce_shuffle不然 MapReduce 任务跑不起来configuration property nameyarn.nodemanager.aux-services/name valuemapreduce_shuffle/value /property property nameyarn.nodemanager.aux-services.mapreduce_shuffle.class/name valueorg.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 和 PySparkHadoop 起来之后接下来配 Python 侧。我建议直接用虚拟环境管理依赖避免污染系统 Python。先装好 Python 3.10 和 pip然后创建虚拟环境并安装 PySparkpython3 -m venv ~/pyspark_env source ~/pyspark_env/bin/activate pip install pyspark3.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 里有两个核心角色ResourceManagerRM负责全局资源分配NodeManagerNM负责管理单台机器上的资源。当你要跑一个 PySpark 任务时客户端会向 ResourceManager 提交 ApplicationRM 找到合适的 NodeManager启动一个 ApplicationMaster 负责协调这个任务ApplicationMaster 再向 RM 申请容器Container然后在容器里启动 Executor 进程真正执行计算。这套机制用生活类比来解释就是ResourceManager 是酒店前台NodeManager 是楼层服务员ApplicationMaster 是会议的会务组会务组找前台要会议室前台协调楼层服务员来布置会议才能顺利开起来。3.3 为什么实际写代码很少直接碰 MapReduceMapReduce 是 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 自带的内置函数解决大部分场景但总有一些业务逻辑需要你自己写。这时候就需要 UDFUser 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/, headerTrue, inferSchemaTrue ) # 清洗去掉 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 refusedHadoop 服务没启动或端口配置不对检查jps确认 NameNode/DataNode 进程检查core-site.xml端口Container killed by the ResourceManagerExecutor 内存超限调大spark.executor.memory或者减少单个 Executor 的核心数OutOfMemoryError: Java heap space数据量太大Executor 堆内存不够增加分区数、减少缓存数据、调大堆内存FileNotFoundError: input path does not existHDFS 路径写错或文件还没上传用hdfs dfs -ls确认路径注意hdfs://前缀Cannot connect to the clusterYARN 集群连接失败检查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 的组合跑起来。
返回列表