1. 先拆掉那堵墙:Spark 和 Hadoop 到底是什么关系
入行大数据这些年,我被问得最多的一句话就是:“Spark 是不是比 Hadoop 快,所以我们直接用 Spark 就行?”这个问题背后的混乱,几乎每个团队都经历过。严格来说,这种问法本身就暴露出一个关键误解:Hadoop 根本不是“一个东西”,而是一个生态家族;Spark 也不是 Hadoop 的替代品,它只是家族里一个更能干的计算引擎。
用个生活化的比方。Hadoop 好比一套完整的厨房系统——有储藏室(HDFS)、有操作台(YARN)、有一本祖传菜谱(MapReduce),还有一堆配套工具(Hive、HBase、Zookeeper 等)。而 Spark 是一台多功能料理机,它确实能把菜做得更快,但前提是你得有食材(数据)、有地方放食材(存储层),并且有人告诉你放多少水、转几分钟(调度与资源管理)。所以真实的生产环境里,Spark 和 Hadoop 从来不是二选一,而是深度绑定。
很多新人在第一次接触这两个词时,会误以为它们是同一层面的竞品。这个误会的根源在于:Hadoop 的文档里 MapReduce 占了大量篇幅,而 Spark 的教程又动不动就“集群模式”“HDFS 读写”,导致大家天然觉得它们在抢同一个生态位。但从架构角度看,Spark 和 Hadoop 的关系更多是“寄生”与“协同”——Spark 读写 HDFS 作为存储底座,Spark 跑在 YARN 上作为计算调度,Spark SQL 又能无缝对接 Hive 的元数据。这不是竞争,这是标准的黄金搭档。
这篇文章我会从底层计算模型入手,讲清楚两者性能差异的根源,再结合我在生产环境中做过的日志分析、用户画像、离线报表等项目,给出选型判断标准和一个完整的 Spark on YARN 落地案例。如果你正在纠结“团队到底该上 Hadoop 还是 Spark”“两套都要的话怎么分工”,这篇应该能给你一个相对清晰的答案。
2. MapReduce 与 Spark 计算模型:性能差距到底差在哪
2.1 中间结果落盘 vs 内存管道:最根本的分水岭
我们得先回到计算模型本身。MapReduce 的核心流程是 Map 阶段 → Shuffle 阶段 → Reduce 阶段,每个阶段结束之后,中间结果几乎都要写入磁盘。为什么这么做?因为这套架构出生的年代,机器内存以 GB 为单位,磁盘廉价且可靠,为了保证任务挂掉之后能从断点恢复,落盘是最稳妥的方案。
但代价也肉眼可见:一个简单的 WordCount 任务,map 输出要落一次盘,shuffle 期间的排序合并要落一次盘,reduce 拿到数据可能还要落一次。三次磁盘 I/O 算下来,大量时间都耗在序列化、反序列化和磁盘读写上。如果数据量上了 TB 级,这种设计就是灾难。
Spark 走的是完全不同的路。它采用DAG(有向无环图)计算模型,把一个任务拆成若干 Stage,每个 Stage 内部尽可能把多个算子串联成一条流水线。数据在内存里以 RDD 分区为基本单位,算子之间不需要落盘,同一 Stage 内的 map、filter、flatMap 直接在一个内存管道里完成。只有当跨 Stage 发生 shuffle 时,才需要把中间结果写到本地磁盘。
我用一个非常直观的数据来说明差距。同样的 TPC-DS 基准测试,在 100GB 数据规模下,Spark SQL 的查询耗时通常是 Hive on MapReduce 的10 到 50 倍。这还是在 Spark 没有做任何调优的情况下。如果你把 Kryo 序列化、数据压缩、内存优化全部打开,差距还会进一步拉大。当然这个数字不是绝对的——MapReduce 在极端简单的任务上也能跑出不错的成绩,但绝大多数真实业务场景里,Spark 的优势是压倒性的。
2.2 DAG 调度与懒执行:为什么 Spark 更加“聪明”
MapReduce 的每次作业(Job)之间是完全独立的,上一个 Job 的产出必须落盘,下一个 Job 再重新读入。打个比方,你让一个实习生整理一批文件,他每完成一步就要把文件放回档案柜,下一步再重新拿出来——这不仅慢,而且每一步之间没有任何全局优化空间。
Spark 的执行计划则是由 Driver 端统一构建的 DAG。它会在真正执行之前做两件 MapReduce 做不到的事:
第一,逻辑优化。Spark SQL 执行前会经历 Catalyst 优化器,它会自动做谓词下推、列剪裁、常量折叠等优化。举个例子,你对一张 100 个字段的表做SELECT col_a FROM tbl WHERE date = '2024-01-01',Catalyst 会在扫描阶段就直接把 99 个没用的字段剪掉,只读需要的列。而 MapReduce 的 Hive 在早期版本里,是会傻乎乎地把整行数据都读出来再过滤的。
第二,物化策略优化。DAG 调度器会根据算子之间的依赖关系,决定哪些 Stage 可以流水线执行,哪些必须 shuffle,哪些 RDD 需要 cache 到内存供后续复用。你可以显式调用.cache()或.persist(),把一个被多次复用的数据集固定在内存里。这在迭代计算、交互式查询场景下效果极其明显——因为不需要反复从磁盘读取同一份数据。
就拿我们之前做过的 ALS 协同过滤推荐来说:MapReduce 实现每轮迭代都要完整地跑一个 Job,从 HDFS 重新读数据、计算、写回,一个 5 轮的迭代任务要花 2 小时;换成 Spark MLlib 之后,同样的数据规模压缩到 15 分钟内,而且代码量减少了一半。这就是 DAG 调度带来的范式差异。
2.3 Shuffle 机制与容错策略:各有取舍
谈到 shuffle,很多人以为 Spark 一定比 Hadoop 好,其实未必。MapReduce 的 shuffle 虽然慢,但它足够稳定——它的排序是全局的、确定性的,处理数据倾斜的方式也很简单粗暴:增加 reduce 数量,或者靠 Combiner 做预聚合。Spark 的 shuffle 默认基于 Hash 分区,虽然快,但一旦数据分布不均,就很容易出现某个 Executor 撑爆内存的 OOM 问题。
容错策略上两者走的也是不同路线。MapReduce 的容错粒度是“任务级别”——某个 Map 任务挂了,单独重跑这个任务即可,因为中间结果都落盘了。Spark 的容错则依赖RDD 的血缘 (Lineage)——一个分区数据丢失了,就根据血缘关系重新从父 RDD 算一遍。这条线路的优势是省掉了持久化的开销,但在血缘链特别长、或者某个 Stage 特别昂贵的情况下,重算代价也不小。所以生产里我们经常结合 checkpoint 机制,在某些关键 Stage 后手动把 RDD 写到 HDFS 上做快照,牺牲一点速度换取恢复效率。
这里想强调一个很多初学者没意识到的事实:Spark 的“快”建立在足够的内存之上,而这恰恰是它在云环境里成本更高的原因。当你评估“要不要从 Hadoop 迁移到 Spark”时,CPU 不是主要瓶颈,内存和网络带宽才是。内存不够,Spark 会退化为反复 GC 甚至 OOM,性能可能还不如 MapReduce 稳定。
3. 选型不是站队:什么场景该用 Hadoop,什么场景该用 Spark
3.1 十类典型业务场景的适用性对照
我在不少企业的技术评审会上看到过类似争论:架构师拍板“我们全面转向 Spark”,然后运维在台下苦笑——因为批处理里的 ETL 清洗任务,Hive on MapReduce 跑了三年都没出过问题,有什么必要为了“用新技术”而换引擎?反过来,也有人非要在 Spark 上跑 5MB 的小数据集,结果光启动 Executor 的时间就比整个计算时间长,纯粹是浪费资源。
真实世界的选型逻辑,从来不是“哪个先进用哪个”,而是“哪个更匹配你的负载特征”。我根据自己的实践经验,整理了一张对照表,可以在方案预选时直接参考:
| 场景 | Hadoop/MapReduce 适用性 | Spark 适用性 | 推荐选择 |
|---|---|---|---|
| TB 级以上的批量 ETL | 稳定,但耗时可观 | 内存充足时效率极高 | 优先 Spark |
| 百 GB 内的临时分析 | 略慢,可接受 | 快且交互性好 | 优先 Spark |
| 迭代式机器学习算法 | 每轮都落盘,几乎不可用 | 内存复用,天然契合 | 必须 Spark |
| 流式数据处理 | 原生不支持 | Structured Streaming 可用 | 必须 Spark |
| 超大表 JOIN | shuffle 稳定但慢 | 可能 OOM,需调优 | 视内存和倾斜情况定 |
| 冷数据归档/存储 | HDFS 是核心底座 | 不适用 | 保留 HDFS |
| 实时查询/即席分析 | 延迟太高 | Spark SQL 快很多 | 优先 Spark |
| 数据仓库分层建设 | Hive 成熟稳定 | Hive on Spark 效率更高 | 建议 Spark |
| 少量文件的小任务 | 可用 | 启动开销大,反而不划算 | 用 Hadoop 即可 |
| 数据湖/存算分离架构 | HDFS 继续做存储 | 计算层选 Spark | 二者结合 |
3.2 非要用 Spark 的硬性指标
如果你正在做一个架构选型,最需要搞清楚的一件事是:什么情况下 Spark 是“必须”而不是“可选”?根据我的实战经验,下面这几个硬指标只要踩中一个,你大概率就绕不开 Spark:
指标一:算法需要多轮迭代。机器学习、图计算、推荐系统这类任务,天然就是迭代式。比如 K-Means 聚类要跑几十轮收敛,逻辑回归要反复计算梯度。Spark 可以在一轮迭代里把需要复用的权重向量、特征矩阵 cache 在内存里;Hadoop 每一轮都得全部从 HDFS 重新读入,时间成本直接拉满。
指标二:交互式即席查询。业务方说“我改个过滤条件,重新跑一下看看结果”,MapReduce 的典型响应时间是几分钟到几十分钟(受 Job 启动开销和调度延迟影响)。Spark SQL 使用Thrift Server提供常驻服务,查询复用 Session 和缓存,秒级响应是常态。我们在给运营团队做的自助分析平台里,后端就是 Spark SQL Thrift Server,配合预热的 Hive 表,能把 90% 的查询控制在 10 秒内。
指标三:需要复杂的数据管道。一段数据管道里有不同的分析场景,比如读取 HDFS 文件后做一次过滤,再按用户维度聚合,再分别输出到多个目标源。Spark 的 DAG 调度会根据血缘关系做最优执行,MapReduce 则会拆成多个独立的 Job,每个 Job 都要排队等待调度。
3.3 保留 Hadoop 生态的不可替代部分
选型时最容易犯的错误就是“一刀切”——把 Hadoop 整个生态都排除了。事实上,Hadoop 里有几个组件是 Spark 完全替代不了的:
- HDFS:Spark 计算一百次,数据还是得存 HDFS 上。它提供的高容错存储、多副本机制、跨节点数据分布,是 Spark 运行的数据底座。
- YARN:虽然 Spark 也能用 Standalone 模式跑,但要让 Spark 和 Hive、MapReduce、Flink 等共享一个集群的资源,YARN 是更成熟也更省心的资源调度方案。
- Zookeeper:Hadoop 集群的 NameNode 高可用依赖它,HBase 的 RegionServer 协调也依赖它。Spark 自身虽然不强制要求,但一旦跟 HBase、Kafka 等生态组件集成,Zookeeper 基本绕不开。
我们团队有一条不成文的原则:存储用 HDFS,调度用 YARN,查询和计算首选 Spark,只对极少数轻量级任务保留 Hive on MapReduce 执行引擎。这个组合既发挥了 Spark 的性能优势,又稳住了 Hadoop 生态的可靠性和易运维性。
4. 生产级 Spark on YARN 部署:从资源规划到三步走落地
4.1 为什么生产环境优先选择 Spark on YARN
技术选型落到部署层面,第一个决策就是“Spark 跑在哪种模式”。Standalone 模式部署简单,适合学习和测试;Mesos 模式现在用的人越来越少了;Kubernetes 模式在云原生环境里前景很好,但运维门槛高,而且与 HDFS 的数据本地性调度还不够成熟。
我推荐生产环境老老实实用Spark on YARN(yarn-client 或 yarn-cluster),原因有三:
第一是资源统一管理。一套 YARN 集群上可以同时跑 Spark、Flink、MapReduce 任务,资源按队列隔离,不会出现“Spark 把内存吃光导致 Hive 没法跑”的尴尬局面。第二是高可用和安全性。YARN 的 ResourceManager 支持主备切换,且已深度对接 Kerberos 认证体系——这在政企和金融机构几乎是硬性要求。第三是简化运维。Spark 的 Driver 和 Executor 生命周期由 YARN 统一管理,不像 Standalone 模式需要额外部署 Master/Worker 进程;集群扩容时,只需要在 YARN 节点上准备好 Spark 客户端和依赖包即可。
4.2 资源规划与内存参数计算:不看这篇你迟早踩 OOM 的坑
Spark 性能调优里,内存配置是重中之重。规划时最核心的公式是:
单个 Executor 可用的堆内存 = spark.executor.memory - spark.memory.offHeap.size - 系统预留(约 300MB~600MB)
而spark.executor.memory里面,又分为Execution 内存(用于 shuffle、join、aggregation)和Storage 内存(用于 cache RDD),两者通过spark.memory.fraction(默认 0.6)共享。很多人一上来就把spark.executor.memory调得很大,比如 32GB,但忽略了 JVM 的 G1GC/ParNew 在超大堆下的停顿问题。我在实践中踩过这个坑:一个 64GB RAM 的 Worker 节点,我把 Executor 内存设为 48GB,结果 Full GC 频繁到任务直接卡死。
比较稳妥的做法是:
- 单节点 Worker 内存 64GB:配置 2 个 Executor,每个
spark.executor.memory=24g,spark.executor.cores=8,预留 16GB 给操作系统页缓存和 YARN NodeManager。 - 单节点 Worker 内存 128GB:配置 4 个 Executor,每个
spark.executor.memory=28g,spark.executor.cores=7。 - Executor 内存不宜超过 32GB:超过后 JVM 对象指针压缩失效,内存占用会显著上升,GC 停顿也不可控。宁可多开几个小 Executor,不要开一个巨无霸。
再解释一个经常被忽略的参数:spark.memory.fraction。默认 0.6 意味着 Executor 堆内只有 60% 的内存用于 Spark 自身管理,剩下的 40% 留给 RDD 对象、用户代码和 JVM 结构。如果你的任务以缓存为主(比如反复迭代读取同一份大 RDD),可以把这个比例调到 0.7;如果任务以 shuffle 为主(大量 join、groupBy),反而要调低一些,留更多堆内存给 shuffle 缓冲区。
4.3 部署步骤:从 Hadoop 集群到 Spark 跑通任务
假设你已经有一套可用的 Hadoop 集群(HDFS + YARN),下面是从零把 Spark 接到 YARN 上的标准流程:
第一步:下载并解压 Spark 二进制包。
选择与你的 Hadoop 版本兼容的 Spark 版本。以 Spark 3.5.x + Hadoop 3.3.x 为例,下载spark-3.5.0-bin-hadoop3包,解压到/usr/local/spark。注意不要贪新,建议看 Spark 官方文档的 Compatibility Matrix,避免出现不兼容的坑。
第二步:配置核心文件。
主要改两个文件。spark-env.sh里设置 Java 路径和 Hadoop 相关环境变量:
export JAVA_HOME=/usr/local/java export HADOOP_HOME=/usr/local/hadoop export HADOOP_CONF_DIR=$HADOOP_HOME/etc/hadoop export SPARK_HOME=/usr/local/spark export SPARK_DIST_CLASSPATH=$(hadoop classpath)spark-defaults.conf里配置资源上限和核心参数:
spark.master=yarn spark.yarn.am.memory=2g spark.executor.instances=10 spark.executor.memory=24g spark.executor.cores=8 spark.driver.memory=8g spark.serializer=org.apache.spark.serializer.KryoSerializer spark.sql.shuffle.partitions=200 spark.dynamicAllocation.enabled=true spark.dynamicAllocation.initialExecutors=5 spark.dynamicAllocation.minExecutors=5 spark.dynamicAllocation.maxExecutors=50第三步:提交测试任务。
/usr/local/spark/bin/spark-submit \ --class org.apache.spark.examples.SparkPi \ --master yarn \ --deploy-mode cluster \ --driver-memory 4g \ --executor-memory 8g \ --executor-cores 4 \ /usr/local/spark/examples/jars/spark-examples_2.12-3.5.0.jar \ 10如果能看到 YARN 的资源申请记录和最终输出的 Pi 值,说明 Spark on YARN 已经打通了。接下来就能把你的业务代码包(jar 或 Python 脚本)提交上去了。
注意:部署模式
cluster和client最大的区别是 Driver 进程跑在哪里。cluster模式下 Driver 由 YARN 的 AppMaster 托管,适合生产环境定时调度;client模式下 Driver 在你的提交机上,更适合调试阶段实时看日志。我建议调试用client,生产全部走cluster。
5. 一套可复用的日志分析管道:Hadoop 存数据,Spark 算数据
5.1 整体链路设计与数据模型
聊完理论就得上点真东西。我拿一个之前做过的用户行为日志分析项目做例子,这套管道从数据采集到最终报表,完整体现了“Hadoop 存储 + Spark 计算”的经典组合。
数据链路是:接收集群 NGINX 的访问日志(日均约 5 亿条,原始数据 600GB 左右)→ Flume 写入 HDFS 的原始日志目录 → Hive 建外部表 → Spark 定时任务做清洗、解析、聚合 → 结果写回 Hive 分区表 → 上层报表工具查询。
这个设计里最关键的决策是分层存储。HDFS 上分三个层:
/data/raw/nginx_log/:原始日志,按天分区,保留 30 天。这里除了数据本身,不建议做任何处理,保证数据不缺不重,随时可以追溯原始信息。/data/dwd/user_behavior/:清洗后的明细数据,Spark 做了解析、过滤、规范化。按天和小时分区。/data/ads/user_metrics_daily/:按用户维度聚合的日报表,直接被 BI 工具读取。
对应的 Hive 外部表分别为ods_nginx_log、dwd_user_behavior、ads_user_metrics_daily。这套模型的好处是:ODS 层存的是“事实”,DWD 层是“可用的事实”,ADS 层是“聚合后的结论”,每一层的错误都不会污染上一层。
5.2 Spark 清洗与聚合的代码骨架
清洗任务用 Spark SQL 写起来非常简洁。下面是一个简化版的每日清洗任务核心逻辑:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, regexp_extract, to_date, hour spark = SparkSession.builder \ .appName("nginx_log_etl") \ .enableHiveSupport() \ .getOrCreate() # 读取前一天 HDFS 上的原始日志 input_path = "/data/raw/nginx_log/2025-01-20" raw_df = spark.read.text(input_path) # 用正则解析 Nginx 日志 parsed_df = raw_df.select( regexp_extract("value", r'^(\S+)', 1).alias("remote_addr"), regexp_extract("value", r'\[([^\]]+)\]', 1).alias("request_time"), regexp_extract("value", r'"(\S+)\s(\S+)\s([^"]+)"', 2).alias("request_path"), regexp_extract("value", r'"(\S+)\s(\S+)\s([^"]+)"', 1).alias("http_method"), regexp_extract("value", r'\s(\d{3})\s', 1).alias("status_code"), regexp_extract("value", r'"([^"]*)"\s*$', 1).alias("user_agent") ) # 清洗过滤异常数据 cleaned_df = parsed_df.filter( (col("status_code") != "") & (col("request_path").rlike("^/api/")) & (col("request_time") != "") ) # 写回 DWD 层 Hive 分区表 cleaned_df.write \ .mode("overwrite") \ .partitionBy("dt", "hh") \ .format("parquet") \ .saveAsTable("dwd_user_behavior")需要注意的是partitionBy的字段要在写入前取出来,比如解析请求时间里的日期和小时。我这里为了代码整洁省略了这一步,实际生产里必须先withColumn("dt", to_date(col("request_time")))再分区写入。
聚合任务也很直观。假设运营要按天统计每个用户的访问次数、独立 IP 数和高频接口 TOP10:
daily_metrics = spark.sql(""" SELECT user_id, COUNT(*) AS pv, COUNT(DISTINCT remote_addr) AS uv, COUNT(DISTINCT request_path) AS api_cnt, dt FROM dwd_user_behavior WHERE dt = '2025-01-20' GROUP BY user_id, dt """) daily_metrics.write \ .mode("overwrite") \ .partitionBy("dt") \ .format("parquet") \ .saveAsTable("ads_user_metrics_daily")这段 SQL 在 Spark 里跑完,600GB 原始日志的清洗加聚合,我们当时用了 20 个 Executor(每个 24GB 内存)大约 25 分钟。如果换成 MapReduce,同样的数据量,经验上至少要 4 小时以上。这个差距在企业日常报表场景里非常致命——当天数据如果跑不出来,管理层第二天早上就拿不到前一天的经营分析。
5.3 调度与依赖管理:生产环境的工程化细节
上面的代码只是单机跑通。真正的生产环境还需要解决“任务什么时候跑、跑挂了怎么办、结果对不对”这三个问题。
我推荐用Apache DolphinScheduler(或 Azkaban)来做工作流调度。典型的工作流长这样:
- ODS 层检查任务:判断 HDFS 上
/data/raw/nginx_log/2025-01-20目录是否存在,且文件数大于阈值。 - Spark 清洗任务:依赖第 1 步成功,执行上面的 ETL 脚本。
- 数据质量校验任务:SQL 统计 DWD 表里的记录数、空值比例,如果异常则发告警并阻断下游。
- Spark 聚合任务:依赖第 3 步成功,生成 ADS 层数据。
- 报表数据导出任务:把 ADS 表的数据同步到 MySQL/ClickHouse,供 BI 查询。
这套流程加上失败重跑、告警通知机制,才是一个能交给运维同学安心睡觉的管道。很多人把 Spark 写得出神入化,但整个数据管道一到凌晨就挂,原因不在 Spark 本身,而在调度和监控没跟上。生产环境里,数据管道工程化的稳定性,远比某个算子的执行速度重要。
6. 真实踩坑记录:OOM、小文件和 Spark+Hive 协作的暗坑
6.1 Executor OOM 的排查链路,我之前是怎么一步步定位的
先说一个我们线上最典型的 OOM 事故。某天凌晨三点的离线任务突然大面积失败,错误信息清一色是java.lang.OutOfMemoryError: Java heap space。如果是新手,这时候第一反应就是调大spark.executor.memory——但这往往治标不治本,甚至会让 OOM 来得更晚一点但更猛烈。
我当时的排查链路是这样的:
第一步,先看 Spark UI 上的 Stage 详情。发现 OOM 集中发生在某个 join 操作的 Shuffle Read 阶段,而不是数据读取或计算阶段。这就把怀疑范围缩小到了“shuffle 数据倾斜”或“join 键分布不均”。
第二步,用spark.sql.adaptive.coalescePartitions.enabled和动态资源分配日志确认 Executor 数量。发现某个 Executor 上的 Shuffle Read 数据量是其他 Executor 的 30 倍——典型的热键问题。
第三步,检查 join 键的分布。因为我们的场景是用设备 ID 关联用户行为,有一批老设备 ID 占了所有数据的 80%——这些 ID 对应的日志量巨大,所以按哈希分区时全部挤到了同一个分区。
最终修复方案是给 join 加了一个Salt 前缀打散技巧:把设备 ID 后加一个 0~9 的随机后缀,join 时先关联打散后的临时键,最后再按真实设备 ID 做二次聚合。这个方法在绝大多数热点倾斜场景里都有效。另外还有一个常用配置是spark.sql.shuffle.partitions,当数据量不大时把它调小(比如 50),可以减少分区数、降低每个分区内的哈希碰撞概率。这里想提醒大家:OOM 不是简单的“内存不够”,它大概率是“某一块内存被不均匀地塞满了”。只看总量不看分布,永远修不到根上。
6.2 HDFS 小文件问题:Spark 写数据时容易忽视的定时炸弹
Spark 是个“制造小文件”的能手。尤其是你用.repartition(500)然后再写入 Hive 分区表时,每个分区可能对应 500 个文件,时间一长,HDFS 上堆满了 10KB、20KB 的小文件。NameNode 要管理所有文件元数据,几百万个小文件直接造成内存压力,甚至让整个集群进入保护模式。
这个问题我在很多团队都遇到过,根本原因是 Spark 写文件时分区数由最终 Stage 的并行度决定,而不是由你的预期文件数决定。解决办法有以下几种,按优先级排序:
第一,写 Hive 表前先repartition(分区个数)或coalesce(更少的分区)。比如你的 T+1 日报表最终只想生成 20 个文件,那就在写之前把分区数压到 20,底层的文件数量会直接减下来。
第二,如果已经产生大量小文件,用Hive 的INSERT OVERWRITE ... SELECT重新合并一次,或者用 Spark 读一遍再写一遍的方式做文件重整。我比较习惯的做法是在 DWD 层任务里加一个“文件瘦身”子任务,定期把目录下的小文件合并成大文件。
第三,尽量避免dynamicPartitionOverwrite误删整个分区。开启spark.sql.sources.partitionOverwriteMode=static可以避免你在覆盖某个分区时把其他分区的文件也清掉。这个坑我们栽过一次,第二天 BI 组的人来问“为什么昨天的数据没了”,排查半天发现就是分区覆盖模式设置不对。
6.3 Spark 读 Hive 表和直读 HDFS 的差异
最后聊一个容易踩的隐性差异:Spark SQL 读 Hive 表走的是Hive Metastore 的元数据,读 HDFS 目录则是直接列式扫描 Parquet 文件。前者能拿到表的 schema、分区信息、统计信息,所以查询规划更智能;后者则是你给的路径里有什么文件就扫什么文件。
所以我的建议是:如果上游数据最终要提供给报表工具使用,尽量建好 Hive 外部表,让所有计算统一走spark.sql。先建表:
CREATE EXTERNAL TABLE IF NOT EXISTS dwd_user_behavior ( remote_addr STRING, request_time STRING, request_path STRING, user_id STRING ) PARTITIONED BY (dt STRING, hh STRING) STORED AS PARQUET LOCATION '/data/dwd/user_behavior';然后做MSCK REPAIR TABLE或ALTER TABLE ... ADD PARTITION把 HDFS 上的分区注册到元数据里。这样 Spark 就能利用 Hive 的元数据信息进行列剪裁和分区剪裁,效率跟在 Hive 里跑是同一水平。
但要注意范式不一致的问题:同一个 HDFS 目录如果被外部程序(比如 Flume 直接写文件)改了文件格式,Hive 元数据不会自动感知。所以生产上要约定:文件写入必须通过 Spark SQL 或 Hive 任务完成,任何绕过元数据的写入都要先重建分区。
6.4 参数调优里的几个容易忽略的“小陷阱”
第一个陷阱是序列化方式。默认 Java 序列化虽然兼容性最好,但性能实在拉胯。我的建议是强制开启 Kryo,并注册需要用到的类:
spark.conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer") spark.conf.set("spark.kryo.registrationRequired", "false")开启后,shuffle 和 cache 数据的序列化大小通常能缩小到 Java 序列化的 1/5 到 1/10。在 shuffle 量很大的场景里,这个配置带来的提速比加内存还明显。
第二个陷阱是 GC 策略。默认的 Parallel GC 在超大数据集上容易频繁 Full GC。建议 JVM 参数里加上-XX:+UseG1GC,并在 driver 和 executor 中分别设置。G1GC 对几十 GB 的堆更友好,可以显著减少大堆下的停顿。
第三个陷阱是动态资源分配在团队共享集群上的副作用。这个功能很好,但如果你和 Hive 任务共享同一个 YARN 队列,Spark 的动态扩容会把队列资源吃满,导致 Hive 任务饿死。解决方案是给 Spark 和 Hive 分配不同的 YARN 队列,或者设置spark.dynamicAllocation.maxExecutors到合理上限,让资源分配有节流机制。
7. 写在最后:一套稳定的技术栈,不是看谁跑得快,而是看谁不掉链子
聊到这里,你应该已经形成了自己的判断。Spark 和 Hadoop 之间不是新与旧的代际更替,而是不同层级的组件各司其职。HDFS 负责可靠地存,YARN 负责公平地分,Spark 负责高效地算——三者组合起来才是真正的大数据分析底座。
我个人的体会是,技术选型到最后拼的其实是“稳定压倒一切”。Spark 再快,如果你们团队没有足够的内存预算,或者没人能 hold 住 Executor 参数调优,那它带给你的只有半夜三更的告警电话。Hadoop 再慢,它的成熟生态和庞大的社区经验积累,让它在“不出错”这件事上拥有极强的优势。我的建议是:小步快跑,先在团队里拿一两个非核心任务从 MapReduce 往 Spark 迁移,跑顺了再铺开。而不是一股脑把生产管道全部切过去,等出了问题再悔之晚矣。
最后分享一个一直沿用到今天的小习惯:每次 Spark 任务上线前,先在测试环境用近一周的真实数据量压一遍,观察它的 Shuffle Read 大小、GC 耗时和 Executor 内存水位。这三项数据比任何基准测试都更能说明问题。等你在生产环境踩过几次 OOM、小文件、倾斜的坑再回头看,你就能做到心里有数——知道什么场景该让 Hadoop 慢慢跑,什么场景必须把 Spark 推上去。