简介:这是一份面向大数据架构师、解决方案人员及技术学习者的Hadoop与Spark项目案例分析文档。文档从实际工作常见场景出发,梳理了数据整合、专业分析、Hadoop作为一种服务、流分析、复杂事件处理、ETL流、更换或增加SAS七类典型大数据项目,逐一说明项目特点、技术栈选型与适用边界,例如数据湖常用HDFS+Hive/Impala,流分析可选择Spark Streaming/Storm,复杂事件处理则需关注毫秒级延迟的架构差异。资源为1个docx文档,压缩包约105KB,内容结构清晰,适合作为大数据项目方案评审、技术选型或内部培训的参考材料。目前已有479人学习/下载,对于希望快速建立大数据项目全局认知、规避重复建设与资源闲置等常见问题的读者,具有直接参考价值。
1. 一份Hadoop和Spark大数据项目案例分析,怎么把它变成能跑的集群
很多订阅了"Hadoop和Spark大数据项目案例分析"这份文档的人,第一反应是打开Word读一遍架构图和代码段,然后存进网盘吃灰。但真正的从业者都知道,案例分析这类资料的价值不在那几页纸,而在于你能不能照着它把自己的集群跑起来,把案例里的MapReduce和Spark作业重新实现一遍。这篇文章就围绕这个标题,从伪分布式环境搭建、Spark集群配置、典型案例拆解,到数据倾斜和版本兼容这些坑,给出一条可以照做的落地路径。
适合谁看?正在做Hadoop课程设计或毕业设计的学生,准备大数据面试的选手,以及想在企业里快速搭一套实验环境验证想法的工程师。你会看到怎么在虚拟机上装Hadoop伪分布式、怎么把Zookeeper整合进来、怎么用Spark读JSON跑分析,还会看到那些面试必问但网上说法不一的概念——InputSplit、分区数、内存参数——在这次实操里到底起什么作用。我不会复述某份别人的案例分析,只会按这个方向最常见的方案,把能抄的配置和命令给你,说清参数为什么这样设,以及失败时看什么。
2. 从虚拟机到Hadoop伪分布式:把案例的数据底座搭起来
2.1 为什么先用伪分布式而不是直接上HA集群
大型案例分析里常出现Hadoop HA架构图,NameNode和ResourceManager各两台,看着很正式。但本地复现时,我强烈建议先用伪分布式(单节点同时跑NameNode、DataNode、ResourceManager、NodeManager)跑通流程,再去拆成HA。原因很简单:伪分布式的配置文件最少,JVM进程全部挤在一个节点上,出问题排查链路短。企业里的HA集群涉及JournalNode、ZooKeeper、双NameNode切换,任何一个服务启动顺序错乱,新手都会卡半天,而很多案例的瓶颈并不在HA本身。
2.2 安装与配置:core-site.xml和hdfs-site.xml的最小可用写法
在虚拟机上安装Hadoop,我一般用Ubuntu 20.04或CentOS 7,JDK版本对应Hadoop 3.x选JDK 8或11。解压hadoop安装包到/opt/hadoop后,先改etc/hadoop/hadoop-env.sh,把JAVA_HOME写死,然后就是两个核心文件。
core-site.xml里最关键的是fs.defaultFS,它决定了HDFS的入口。以下是我常用的配置:
<configuration> <property> <name>fs.defaultFS</name> <value>hdfs://localhost:9000</value> </property> <property> <name>hadoop.tmp.dir</name> <value>/opt/hadoop/tmp</value> </property> </configuration>hdfs-site.xml里重点设置副本数和NameNode的检查点目录:
<configuration> <property> <name>dfs.replication</name> <value>1</value> </property> <property> <name>dfs.namenode.name.dir</name> <value>/opt/hadoop/namenode</value> </property> <property> <name>dfs.datanode.data.dir</name> <value>/opt/hadoop/datanode</value> </property> </configuration>副本数设为1,是因为伪分布式只有一个DataNode,设置为默认3会导致日志里一直报复制管道未完成的警告,虽然不影响跑,但会干扰排查。name.dir和data.dir一定要手动指定,否则默认走hadoop.tmp.dir,而tmp目录在系统清理时可能被删,删了之后NameNode格式化信息丢失,启动会报NameNode is not formatted。
格式化NameNode是新手必踩的第一步:执行hdfs namenode -format,看到successfully formatted才能继续。格式化只需要一次,之后不要因为在启动失败就反复格式化,那会导致集群ID变更,DataNode注册不上。
2.3 启动、验证和把案例数据传进HDFS
启动HDFS需要执行sbin/start-dfs.sh,然后执行sbin/start-yarn.sh。很多教程建议start-all.sh,但分开启动更容易定位是HDFS挂还是YARN挂。启动后用jps命令查看进程,伪分布式环境里至少要有NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager这五个进程。少一个都说明对应的服务没起来,直接看logs/目录下对应的log文件,不要瞎猜。
案例里常用到的操作是把本地文件上传到HDFS。比如一份访问日志access.log,用命令行上传:
hdfs dfs -mkdir -p /case/input hdfs dfs -put access.log /case/input/ hdfs dfs -ls /case/input/hdfs dfs和hadoop fs在这条命令上等价,但新版Hadoop推荐用前者。-put之后用-ls验证文件大小和副本数,确认Replication=1。如果上传的文件是几百MB以上的大文件,建议先hdfs dfs -du -h /case/input看实际占用量,避免直接把临时文件塞进根目录。
2.4 Hadoop和Zookeeper整合实战:给HA铺路
案例分析里一旦涉及HA或HBase,就一定需要Zookeeper。即使当前只是伪分布式,我也建议把Zookeeper配起来,因为后续Spark写入Hive或者Kafka对接都要用到分布式协调。
我一般下载zookeeper-3.4.x或3.6.x版本,解压后把conf/zoo_sample.cfg复制成zoo.cfg,至少改三个参数:
dataDir=/opt/zookeeper/data clientPort=2181 server.1=localhost:2888:3888单机模式下server.1可以注释掉,但如果要模拟真实HA,就在同一个节点上配三个Zookeeper实例,端口分别改成2181/2182/2183,dataDir也要分开。然后对每个实例的myid文件写入不同的id号。启动时逐个执行zkServer.sh start,用zkServer.sh status查看哪个是Leader。
Hadoop整合Zookeeper最直接的应用是自动故障转移。在hdfs-site.xml里加两个属性:
<property> <name>dfs.ha.automatic-failover.enabled</name> <value>true</value> </property> <property> <name>ha.zookeeper.quorum</name> <value>localhost:2181</value> </property>伪分布式不做HA的话,这两个属性不生效,但提前写进去能让你的环境后续平滑升级到双NameNode。注意启动顺序必须是Zookeeper先于Hadoop,否则NameNode会一直报连不上/hadoop-ha这个ZNode。
3. Spark集群搭建:从独立模式到YARN模式
3.1 Spark选型和下载:版本号对齐是玄学
Spark集群搭建是面试热词,也是让人头大的一个环节。最常见的问题就是Spark和Hadoop版本兼容,其实只要把Hadoop客户端的hadoop-client-api版本对齐就行。我用Spark 3.3.x配Hadoop 3.x,下载时选spark-3.3.4-bin-hadoop3.tgz这种预编译版本,不需要自己编译。如果你的Hadoop是2.7,就要下载bin-hadoop2.7的包。
下载后配置spark-env.sh,注意以下参数,这是Spark内存规划的起点:
export JAVA_HOME=/usr/lib/jvm/java-8-openjdk-amd64 export SPARK_HOME=/opt/spark export HADOOP_CONF_DIR=/opt/hadoop/etc/hadoop export SPARK_MASTER_HOST=localhost export SPARK_LOCAL_IP=localhost export SPARK_WORKER_CORES=2 export SPARK_WORKER_MEMORY=2gHADOOP_CONF_DIR这个环境变量很关键。Spark在YARN模式下需要读取Hadoop的core-site.xml和hdfs-site.xml来连接HDFS和ResourceManager。忘了配它,Spark作业提交时会直接报java.net.UnknownHostException,因为它不知道hdfs://localhost:9000在哪。
3.2 YARN模式搭建:让Spark跑在Hadoop资源调度里
企业里最常见的Spark集群搭建方式是Spark on YARN,好处是同一个Hadoop集群既可以跑MapReduce,也可以跑Spark,资源由YARN统一调度。配置方式很简单,在spark-defaults.conf里写入:
spark.master=yarn spark.yarn.jars=hdfs://localhost:9000/spark-jars/* spark.eventLog.enabled=true spark.eventLog.dir=hdfs://localhost:9000/spark-logsspark.yarn.jars指向HDFS上一个目录,需要先把Spark安装目录下的jars上传到HDFS:
hdfs dfs -mkdir -p /spark-jars hdfs dfs -put /opt/spark/jars/* /spark-jars/不设置这个,每次提交作业Spark都会把jars打成大包传到HDFS,导致提交过程卡很久。上传一次后,提交时间能从几十秒降到几秒。
3.3 在Spark中读取JSON:指定schema比让spark猜更好
热点里说"spark中读取json",这是Spark数据分析案例里的常见操作。很多人直接spark.read.json("hdfs:///case/data.json"),然后靠Spark自动推断schema,数据少时没问题,数据一多或者字段有缺失值,推断出的类型可能不稳定。更好的做法是先用一条数据探测字段,然后显式声明schema。
下面这个例子,我处理一个用户行为日志JSON文件,包含userId、action、timestamp三个字段:
import org.apache.spark.sql.types._ val schema = StructType(Array( StructField("userId", StringType, true), StructField("action", StringType, true), StructField("timestamp", LongType, true) )) val df = spark.read.schema(schema).json("hdfs://localhost:9000/case/input/user_log.json") df.select("userId", "action").where($"action" === "login").count()显式schema的好处是,即使某个JSON对象缺失action字段,读入后该字段也是null而不会变成其他类型。时间戳如果是从字符串用to_timestamp函数转换的,建议先存为Long原始值,计算性能更好。Spark读取JSON时,底层走的是FileSourceScanExec,你指定的schema会直接映射到Parquet或者JSON解析器,避免额外的一次schema合并。
3.4 Spark内存参数:为什么你的作业总是OOM
Spark内存是面试必问,也是实际调参的重灾区。默认情况下Executor内存是1g,很多案例数据一多就挂。我通常这样设置提交参数:
spark-submit \ --master yarn \ --deploy-mode cluster \ --driver-memory 2g \ --executor-memory 4g \ --executor-cores 2 \ --num-executors 2 \ --conf spark.memory.fraction=0.6 \ --class com.example.Analysis \ analysis.jarspark.memory.fraction是Executor内用于执行和存储的内存占比,剩下的留给Shuffle和溢写。Java堆外内存要单独看spark.memory.offHeap.size。遇到OOM时,首先看报错是Java heap space还是Direct buffer memory,前者调executor-memory,后者调spark.memory.offHeap.size或者spark.yarn.executor.memoryOverhead。另外,executor-cores不要设太高,单个Executor超过5个core时,GC反而变慢。这里的关键经验是,先满足内存,再看核数。
4. 项目案例分析核心:从InputSplit到数据倾斜的完整链路
4.1 InputSplit是什么:影响Map数量的根本机制
面试题里常问"在一个运行的Hadoop任务中,什么是InputSplit?"。这是案例分析绕不开的概念。InputSplit是Hadoop对输入数据的一个逻辑划分,它并不物理切割文件,而是记录了这个分片对应的文件路径、起始偏移和长度。Map任务的数量等于InputSplit的数量,所以理解它,你就知道为什么数据量相同,任务数却不同。
在FileInputFormat里,默认按文件大小切片,切片大小取minSize、blockSize和maxSize的中间值。HDFS块默认128MB,所以一个200MB的文件会被切成两个分片,产生两个Map任务。可以在mapred-site.xml里调整:
<property> <name>mapreduce.input.fileinputformat.split.maxsize</name> <value>67108864</value> </property>把最大切片设为64MB后,同一个文件变成3个分片。切片变小,Map数变多,每个Map处理的数据变少,适合CPU密集型的任务;但Map数太多会造成启动开销,要平衡。这是从Hadoop面试题里剥出来的实操点。
4.2 Spark中的分区机制:分区数决定并行度
Spark没有InputSplit这个概念,但它有分区(Partition)。理解两者对应关系,能帮你写出高效的分析代码。在读取HDFS文件时,Spark的默认分区数由spark.default.parallelism决定,读取文件时也会受到文件块数量影响。Spark里改变分区数有两个算子:repartition(num)会做全量shuffle,coalesce(num)尽量不做shuffle。前者能让分区数变大,后者用于减少分区数。
下面这段代码是案例分析里常见的操作——对日志数据做聚合,然后重新分区写入HDFS:
val logRDD = sc.textFile("hdfs://localhost:9000/case/input/*.log") .map(line => { val parts = line.split(",") (parts(0), parts(1).toInt) }) val result = logRDD.reduceByKey(_ + _) result.coalesce(1).saveAsTextFile("hdfs://localhost:9000/case/output/agg")reduceByKey会按key做shuffle,默认用HashPartitioner,分区数由spark.default.parallelism决定。如果分区数太少,同一个key的大量数据挤到同一个分区,堆内存涨得快。如果分区数太多,大量小文件写入HDFS,后续读取时Map任务数爆炸。我一般用rdd.partitions.size打印分区数,根据数据量估算,一个分区100MB到300MB比较合适。
4.3 经典案例复现:电商用户行为分析的MapReduce和Spark双实现
案例分析里最常见的是电商日志分析,比如统计每天的用户访问量。用MapReduce实现时,需要写一个Mapper和一个Reducer,整个类多达百行。用Spark实现,代码量少一个数量级。下面是Spark实现:
val df = spark.read.schema(schema).json("hdfs://.../user_behavior.json") val dailyUv = df.groupBy("date") .agg(countDistinct("userId").as("uv")) dailyUv.show()这里groupBy加agg就是典型的宽依赖Shuffle操作。countDistinct在数据量大时会产生精确去重开销,如果允许近似值,可以用approx_count_distinct,它在底层用HyperLogLog算法,误差控制在1%以内,内存占用低很多。案例分析里讲"UV精确统计"只是为了教学,实际生产中我直接用approx_count_distinct加阈值参数。
4.4 数据倾斜:让案例分析从"能跑"到"跑得快"的分水岭
数据倾斜是案例里很少直说,但生产中必踩的坑。现象是某个Reduce任务或Executor执行时间远大于其他,甚至卡死。原因通常是key分布极不均匀,比如按省份统计时,广东的数据量是其他省的几十倍。
解决数据倾斜的常用手段有四种。第一,过滤异常key,比如把空值或"unknown"过滤掉;第二,两阶段聚合,先给key加随机前缀,做部分聚合,再去前缀做全局聚合;第三,将reduceByKey改成mapValues配合groupByKey,但要注意groupByKey会把所有value放内存,风险更大;第四,调整分区策略,用自定义Partitioner让大key单独占一个分区。我一般先做第一种,因为最简单,而且筛选掉null后的数据分布往往就正常了。
以下是一个两阶段聚合的Scala示例,处理"同一userId点击次数特别多"的场景:
import org.apache.spark.util.Random val rdd = sc.textFile("hdfs://.../clicks.csv") .map(line => (line.split(",")(0), 1)) val prefixRdd = rdd.map { case (key, value) => val prefix = (Random.nextInt(10) + 1).toString() (prefix + "_" + key, value) } val partialAgg = prefixRdd.reduceByKey(_ + _) val finalAgg = partialAgg.map { case (key, sum) => val realKey = key.substring(key.indexOf("_") + 1) (realKey, sum) }.reduceByKey(_ + _) finalAgg.collect()随机前缀把原本集中在一个key上的数据打散到最多10个分区,每个分区的压力大幅下降。前缀数量的选择有讲究:设5时倾斜缓解不明显,设20时会增加Shuffle数据量。我通常设为总分区数的一半。这个技巧就是面试里常说的"加盐",但实际写代码时很多人拼字符串拼错,注意substring的索引边界。
5. 避坑/常见问题排查:三个让我半夜修bug的实战记录
5.1 NameNode启动失败,logs里一堆java.io.IOException: NameNode is not formatted
现象:执行start-dfs.sh后,NameNode进程一闪而过,查看logs/hadoop-hadoop-namenode-*.log,最后一行是NameNode is not formatted。
原因:dfs.namenode.name.dir配置指向的目录是空的,或者之前格式化过但目录被清掉了。很多教程让你"重装时再格式化一次",但格式化后的clusterID会变,DataNode的current/VERSION里还是旧的clusterID,导致DataNode无法注册。
解决:第一次配置时,先确认name.dir目录存在且为空,然后执行hdfs namenode -format。一旦启动失败想重来,不只格式化NameNode,还要把DataNode的data.dir也删掉,再重新格式化。用下面这组命令收干净:
rm -rf /opt/hadoop/namenode /opt/hadoop/datanode hdfs namenode -format start-dfs.sh5.2 Spark作业提交后一直卡在INFO YarnClientSchedulerBackend,但任务不执行
现象:spark-submit发出去了,控制台不停打印INFO YarnClientSchedulerBackend,但就是看不到Task执行。
原因:Spark在YARN模式需要把Spark的jars和依赖文件放到HDFS上,并把YARN的ResourceManager地址配好。我遇到过两种具体原因:一是spark.yarn.jars指向的HDFS路径不存在,导致ApplicationMaster启动时反复拉取jar包失败;二是Hadoop的yarn-site.xml里yarn.resourcemanager.hostname配置成了0.0.0.0,Spark客户端无法判断RM地址。
解决:先确认hdfs dfs -ls /spark-jars能看到jar包,然后在spark-defaults.conf里写spark.yarn.jars=hdfs://localhost:9000/spark-jars/*。再检查yarn-site.xml的RM地址,伪分布式里应改为localhost。还有一种低频原因是YARN没启动,执行jps看有没有ResourceManager进程,没有就先start-yarn.sh。
5.3 虚拟机上Spark读不了本地文件,一直报FileNotFound: /user/xxx/input/data.json
现象:用spark.read.textFile("data.json")读本地文件,报错说找不到/user/xxx/input/data.json。
原因:Spark默认的路径是HDFS路径,data.json会被解析成hdfs://localhost:9000/user/xxx/data.json。你以为它在本地,其实Spark找的是HDFS上相对路径。这是新手最容易搞混的坑。
解决:本地文件要写成file:///home/ubuntu/data.json,HDFS文件写hdfs://localhost:9000/case/input/data.json。更推荐直接先把文件传到HDFS,再用HDFS路径,因为集群模式下Executor不在同一台机器,本地路径只在Driver所在节点有效,Worker执行时会失效。在spark-shell里,也可以这么确认:
spark-shell --master local[2] scala> import org.apache.hadoop.fs.FileSystem scala> val fs = FileSystem.get(sc.hadoopConfiguration) scala> fs.exists(new org.apache.hadoop.fs.Path("/case/input/data.json"))用这个方法快速检查路径是否存在,比盯着日志猜快得多。
5.4 内存参数改了但是不生效,Executor还是1g
现象:spark-submit里加了--executor-memory 4g,但在Spark UI的Executors页面看到还是1g。
原因:Spark启动时读取配置的优先级是:spark-submit命令行的参数最高,其次spark-defaults.conf,最后是代码里的SparkConf。如果你在集群模式里设置的这些参数都没生效,多半是提交时--master写错了,或者spark-env.sh里写死了SPARK_WORKER_MEMORY,而Spark on YARN模式下,Executor内存由spark.executor.memory指定,SPARK_WORKER_MEMORY是Standalone模式的参数。
解决:提交时加一行--conf spark.executor.memory=4g,然后去Spark UI的Environment标签页看Spark Properties里是否出现这条。如果出现但Executor还是1g,检查yarn-site.xml里yarn.nodemanager.resource.memory-mb,它的值必须大于Executor内存总和,否则YARN会拒绝分配。伪分布式里物理机内存总共8g的话,Nodemanager最多给6g,--executor-memory 4g依然可行,但再加Driver 2g就可能会崩。
6. 把案例分析写成自己的:用Spark SQL重写一遍,再验证结果正确性
案例分析读十遍,不如自己重写一遍。我的习惯是拿到一份Hadoop和Spark项目案例分析后,先不改任何代码,把它的逻辑拆成三步:输入数据从哪里来,中间有哪些转换,输出结果长什么样。然后不直接照着代码敲,而是用Spark SQL写一套自己的版本交接。这样做的价值在于,Spark SQL能把很多手动map和reduce的逻辑用申明式写法激化,同时自动做谓词下推和列剪枝,比RDD版本更快。
下面是用Spark SQL重写UV统计的例子,它替换了前面第4章RDD里那一大段map和reduceByKey:
CREATE TEMPORARY VIEW user_logs AS SELECT * FROM json.`/case/input/user_log.json`; SELECT date, COUNT(DISTINCT userId) AS uv FROM user_logs WHERE action IS NOT NULL GROUP BY date ORDER BY uv DESC;这段SQL里,json.前缀告诉Spark直接以JSON格式解析指定路径,不需要提前定义schema。Spark SQL的COUNT(DISTINCT ...)底层和approx_count_distinct不同,它走全量精确去重。数据量大时,把COUNT(DISTINCT userId)改成approx_count_distinct(userId, 0.01),结果误差在1%内,但Shuffle数据量会明显减少。这是我从实际跑日志分析得到的经验,面试时讲出来,也会显得你真的碰过数据。
写完SQL后,我习惯做三件事验证结果没有被改坏。第一,运行MapReduce版和Spark版,各自输出结果到两个目录,然后对结果做逐行diff。第二,用Hive的hive-exec里自带的mapreduce.job.counters对比输入记录数和输出记录数,如果Reduce输出记录数大于输入,说明有异常展开。第三,在Spark作业里打印驱动端的执行计划,用df.explain("extended")查看Scan节点下的numFiles和filesize,确认实际读入的数据量是否和HDFS上的文件大小吻合。这三个方法就算不算精妙,但至少能把"看起来跑完了"和"结果真的对了"区别开。
最后提一个我用得最多的习惯:每次调参前,把原参数和改动后的参数都记录在文件的同一行注释里,而不是只记最终值。这会在你优化两三轮之后找不到后悔药时救你一命。很多案例里的参数推荐值来自特定环境下,你的虚拟机内存、磁盘IO、并发任务数都不同,照抄参数大概率翻车。先在大数据集上跑一遍,记录耗时和内存,再逐步调整,这样这套案例分析方案才能成为你自己的解法。希望这些经验能帮到你。
本文还有配套的精品资源,点击获取