ARTICLE DETAIL

资讯详情

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

Hadoop与Spark实战:从环境搭建到数据算法调优全攻略

Hadoop与Spark实战:从环境搭建到数据算法调优全攻略 简介围绕Hadoop与Spark框架的数据算法实战源码包专为大数据开发者、数据工程师以及希望系统学习分布式计算技术的学习者准备。包内共876个文件包含360个Java源文件、34个Scala源码、242个Jar依赖包、63个Markdown讲解文档、31个Shell运维脚本以及csv、tsv、parquet等格式的测试数据压缩包整体约204MB目录结构清晰可按章节检索与复用。源码覆盖HDFS文件读写、MapReduce词频统计与数据清洗、RDD算子转换、Spark SQL结构化查询、MLlib分类回归等经典算法场景并附带可直接运行的输入数据能直观观察实际处理效果MapReduce的combiner与partition机制、Spark的shuffle调优等进阶细节也提供了对应示例非常适合动手实践有助于深入理解分布式计算的设计思想与调优策略。已有682人浏览学习是系统提升大数据处理技能、积累真实项目经验的重要参考资料。 我收到过不少读者留言问题出奇一致标题里这些词——Hadoop、Spark、数据算法、源代码——单独看都认识合在一起就不知道从哪里下手。去年有段时间我专门负责带新人做大数据入门训练发现90%的人卡在同一个地方不是原理听不懂而是跟着教程搭完环境之后跑一个真实的数据算法时集群就像故意跟你作对一样各种报错。这篇文章我不打算讲虚的就沿着一条完整的实战链路把Hadoop和Spark从关系、搭建、算法代码到性能排查的完整经验拆开每个环节都附上可以直接抄走改用的内容希望帮你少走那90%的人走过的弯路。1. Hadoop和Spark先分清“存储”和“计算”再动手1.1 它们不是替代关系而是搭伙干活很多新人进大数据这个领域第一个概念混淆就是Hadoop和Spark到底学哪个是不是有了Spark就不需要Hadoop了答案是它们根本不是同一层面的东西。Hadoop是一个生态体系核心由三块组成HDFS负责分布式存储YARN负责资源调度MapReduce负责分布式计算。你只要记住一个判断标准——凡是涉及“数据放在哪”的问题基本都是HDFS的事凡是涉及“谁来分配算力”的问题基本是YARN的事凡是涉及“数据怎么算”的问题才是计算引擎的事。而Spark本身只是一位“计算选手”它不负责存储也不负责资源管理。它可以把数据从HDFS上读出来跑完计算再写回HDFS。换句话说Spark是站在Hadoop肩膀上的计算框架你的集群可以没有MapReduce但通常不能没有HDFS。所以正确的组合方式不是“Hadoop还是Spark”而是“HDFS YARN Spark”。这也是绝大多数生产集群的标准形态。1.2 数据算法在分布式环境里的真实含义标题里提到“数据算法”我理解很多人的第一反应是是不是要学一堆高深的机器学习公式实际不是。在大数据入门和中级应用场景中所谓的“数据算法”更多是指那些在单机时代很简单、但扔到分布式环境下就需要重新设计的经典算法——词频统计、去重、排序、分组聚合、TopN、Join关联等。为什么“重新设计”这么重要给你一个例子求一个文件里的Top10热词单机环境下你读文件、建HashMap计数、排序取前10三分钟完事。但放在分布式环境下同样的逻辑要考虑数据分散在多台机器上每台机器只能算自己那一份算完局部TopN之后还需要一个合并动作把所有机器的结果汇总合并过程如果设计不当shuffle的数据量可能比你原始数据还大。这就是“设计分布式算法”和“写单机算法”最大的差异。所以这篇博文后面所有代码都会围绕这个思路展开先写单机逻辑再拆成分布式步骤最后落地成Hadoop或Spark的代码。这也是我认为“数据算法源代码”最有价值的学习路径。2. 从伪分布式到真实集群环境搭建最容易翻车的几个环节2.1 伪分布式搭建与格式化失败的真正原因如果你只是想先跑通代码、验证算法不需要一上来就搞三台机器。Hadoop的伪分布式模式Pseudo-Distributed是学习阶段效率最高的方式一台机器模拟所有角色。但很多人在这个阶段就会遇到一个极其经典的报错格式化Namenode时一切正常启动之后DataNode起不来日志里报Incompatible clusterIDs。回想我第一次遇到这个问题的排查过程也是花了半天才搞明白根因。原因是你第一次执行hdfs namenode -format之后NameNode和DataNode各自生成了一套clusterID它们互相匹配才能正常通信。但如果之后你因为某些操作比如修改配置再次执行了格式化NameNode会生成一套全新的clusterID而DataNode的data目录里保存的还是旧的clusterID两边对不上DataNode自然拒绝注册。这个问题的解决其实很简单但网上很多教程没说明白导致新手反复格式化反复失败# 1. 先停掉集群相关进程 stop-dfs.sh # 2. 找到配置文件里指定的数据目录 # 默认在 /tmp/hadoop-${user.name}或查看 hdfs-site.xml 中的 dfs.namenode.name.dir / dfs.datanode.data.dir # 3. 手动清掉 namenode 和 datanode 的目录 rm -rf /tmp/hadoop-${user.name} # 4. 重新格式化并启动 hdfs namenode -format start-dfs.sh核心思路是格式化之前把两边可能残留的旧元数据全部清干净让NameNode和DataNode在同一次启动中重新生成并协商新的clusterID。2.2 集群部署策略学习、面试、生产分别怎么配再往下一步你总会面临一个问题要不要搭集群搭几台我在带新人的时候给过一个经验值分三种场景场景机器数量配置建议说明学习/入门1台伪分布式4核8G起跑通原理和代码足够面试/深度练习3台每台4核8G1台Master 2台Worker能完整演示分布式效果生产/真实项目至少5台起每台16核64G以上Master节点独立NameNode和ResourceManager分开很多新手容易犯的错是一上来就租一堆云主机结果发现集群搭好之后95%的时间都在空转。学习阶段伪分布式完全够用等到你真正需要调优、测并发、验证数据倾斜时再上3台集群这个节奏比较合理。2.3 高可用必须引入Zookeeper这里有一个理解门槛热搜词里有“hadoop和zookeeper整合实战”这也是一个容易卡住人的主题。当你有3台以上机器时NameNode会成为整个HDFS的“单点”——一旦它挂了整个集群就无法读取文件元数据等于全部瘫痪。解决方式就是高可用HA方案跑两个NameNode一个Active一个Standby。Zookeeper在这里的角色是什么简单说它是一个“分布式协调器”。当Active NameNode出现故障时Zookeeper需要快速感知到这个故障然后触发自动切换把Standby NameNode提升为Active。整个过程不需要人工干预。配置HA时的核心文件是hdfs-site.xml你需要在里面声明nameservices、两个NameNode的地址、以及Zookeeper的仲裁地址。这里强调一个最容易踩的坑dfs.namenode.shared.edits.dir必须指向一个JournalNode集群很多人在这一步图省事直接复用DataNode节点导致故障切换时元数据无法同步HA形同虚设。3. 数据算法落地从MapReduce到Spark的源代码拆解3.1 MapReduce版WordCount理解分布式计算的“最小单元”WordCount对大数据的意义相当于HelloWorld对编程的意义。不要因为它简单就跳过它实际上是理解分布式计算模型的最好载体。一个完整的MapReduce任务包含三个核心组件Mapper、Reducer、以及驱动类。Mapper负责把数据切成KV对Reducer负责对相同Key的Value做合并。public static class TokenizerMapper extends MapperObject, Text, Text, IntWritable { private final static IntWritable one new IntWritable(1); private Text word new Text(); public void map(Object key, Text value, Context context ) throws IOException, InterruptedException { StringTokenizer itr new StringTokenizer(value.toString()); while (itr.hasMoreTokens()) { word.set(itr.nextToken()); context.write(word, one); } } }这段Mapper代码的逻辑很简单读进来一行文本按空格切分成单词每遇到一个单词就输出一个(word, 1)的键值对。public static class IntSumReducer extends ReducerText, IntWritable, Text, IntWritable { private IntWritable result new IntWritable(); public void reduce(Text key, IterableIntWritable values, Context context ) throws IOException, InterruptedException { int sum 0; for (IntWritable val : values) { sum val.get(); } result.set(sum); context.write(key, result); } }Reducer做的事情更简单把所有相同单词的计数加起来得到一个最终计数值。这里的关键在于Map阶段产生的中间数据会经过一个叫Shuffle的过程由框架自动把相同Key的数据路由到同一个Reducer上你在代码里完全不用管这个逻辑——这也是MapReduce被称为“隐藏了分布式复杂性的模型”的原因。驱动类不再完整贴出核心就一行job.setMapperClass(TokenizerMapper.class); job.setReducerClass(IntSumReducer.class);3.2 用Spark RDD重写同一算法对比着看才看得出差距等你理解了MapReduce的套路再用Spark重写同一个算法你会立刻感受到“代码量”和“思维复杂度”的双重下降。同样一个WordCountSpark写法如下val lines sc.textFile(hdfs:///input/data.txt) val wordCounts lines .flatMap(_.split( )) .map(word (word, 1)) .reduceByKey(_ _) wordCounts.saveAsTextFile(hdfs:///output/wc_result)如果你已经看懂了MapReduce的版本这段代码你甚至不需要注释就能读懂flatMap负责把每行拆成单词map负责把单词变成KV对reduceByKey负责按Key聚合计数。但你注意到没有这中间少了一个非常重要的东西Shuffle。MapReduce里每个单词和它的计数从Mapper流向Reducer的过程是需要落盘写到磁盘的而在Spark里reduceByKey会先在每个分区内部做一个局部聚合再把合并后的结果通过网络传输到下游节点大幅减少了需要搬运的数据量。这也是Spark“比MapReduce快”的核心原因之一不只是因为它在内存里算更因为它减少了很多不必要的磁盘和网络开销。3.3 字数统计之外TopN和分组排序的两种写法如果你真正想应付面试和实际需求WordCount之后建议立刻啃掉TopN问题。它在分布式环境里比WordCount多了一层考验如何避免把所有数据都shuffle到一个节点上。先说最简单但不推荐的做法// 全局排序取TopN数据量大时很容易造成OOM val topN wordCounts.sortBy(_._2, ascending false).take(10)这种方法写起来很爽但问题是sortBy通常是全局排序需要把所有数据集中处理数据量一旦上来性能急剧下降。更推荐的做法是“两步走”思路先在每个分区内部取局部TopN再对局部TopN结果合并取全局TopN。// 第一步每个分区内部统计Top10 val localTop wordCounts.mapPartitions { iter iter.toList.sortBy(-_._2).take(10).iterator } // 第二步对局部结果做全局Top10 val globalTop localTop.sortBy(_._2, ascending false).take(10)这种写法的精髓在于每个分区只需要输出10条数据最终需要shuffle的数据量被压到了可以忽略不计的量级整个作业性能会有一个质的提升。3.4 结构化数据直接用Spark SQL别自己造轮子如果你要处理的数据本身有结构JSON、Parquet、CSV带表头我建议直接上Spark SQL而不是自己写RDD逻辑。代码可读性和维护性是RDD写法完全没法比的。val df spark.read.json(hdfs:///data/events.json) df.createOrReplaceTempView(events) val result spark.sql( |SELECT user_id, COUNT(*) AS cnt |FROM events |WHERE event_type click |GROUP BY user_id |ORDER BY cnt DESC |LIMIT 10 |.stripMargin) result.show()这段SQL如果你有一点关系型数据库的基础几乎零学习成本。Spark SQL底层的Catalyst优化器会自动做谓词下推、列裁剪等优化很多时候比你手工写RDD性能还好。这就是“选择比努力重要”的典型场景——同样的结果你用RDD硬写可能还需要调优半天换Spark SQL一行优化都不用做。4. 性能疑案Spark on YARN明明配了4核为什么只用了1个4.1 从一次真实排查看CPU配置的完整链路热搜词里有一条让我印象很深“spark on yarn cpu只能用1个是为什么”。这个问题我印象深是因为当初带项目时真被它坑过而且网上答案非常分散。当时的情况是这样用户用spark-submit提交作业时明明通过--executor-cores 4指定了每个Executor要4个CPU核但在YARN的资源管理界面上看到的实际使用核数始终只有1个作业跑得奇慢无比。排查过程我按这个链路走第一步确认spark-submit的参数是否生效spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-cores 4 \ --executor-memory 8g \ --num-executors 3 \ your-job.jar参数看起来完全没毛病。第二步去看YARN的调度配置。如果你的YARN被配置成了yarn.scheduler.capacity.resource-calculator或者DominantResourceCalculator在分配资源时是按“CPU 内存”整体来算的。如果每个容器的yarn.nodemanager.resource.cpu-vcores没有正确设置YARN根本不知道每台机器有几个CPU可以用所有容器都会被默认限制为最小值1个vCore。第三步去检查yarn-site.xml里这两个参数property nameyarn.nodemanager.resource.cpu-vcores/name value8/value /property property nameyarn.scheduler.maximum-allocation-vcores/name value8/value /property问题往往就出在这里cpu-vcores没有显式设置成物理机的真实核数YARN默认只认1个核然后无论你--executor-cores写几都会被限制成1。把这两个参数按实际核数配置后重启YARN集群再提交同样的任务CPU核数就正常了。4.2 Spark内存模型为什么executor-heap设了10G还是OOM提到CPU就绕不开内存。很多新人设置--executor-memory 10g以为Executor最大就能用10G结果跑着跑着频繁报OOM非常困惑。这里的关键在于Spark对Executor内存做过精细划分。在Executor的堆内内存里主要分为三块Execution内存用于shuffle、join、sort等计算过程Storage内存用于缓存RDD数据Reserved内存留给系统自身使用默认300MBspark.executor.memory设置的是Executor进程的总堆内存但这部分内存并不是“纯计算可用内存”。实际执行计算时能够真正被一个task使用的内存需要按spark.memory.fraction默认0.6再乘以spark.memory.storageFraction默认0.5来计算。这是很多人OOM的真正原因你设置了10G堆内存但减去Storage保留的30%左右真正的Execution内存可能只有4G左右。如果作业里Join或Aggregate操作的数据量远超这个值很容易出现内存溢出而它并不是因为你“总内存不够”。调优的方向有两类一是增加Executor内存二是降低spark.memory.storageFraction把更多内存让给Execution。4.3 一套可以直接参考的配置基线根据不同作业类型我自己常用的配置从这样一个表起步参数建议值适用场景--executor-memory8g-16g通用默认--executor-cores4-8建议不超过物理机核数的1/4--num-executors依据队列资源单作业不建议超过50spark.memory.fraction0.6-0.75常规即可spark.sql.shuffle.partitions200集群规模较大时可调大这个表不是标准答案但作为初始值基本不会出大问题。后续再结合监控面板里Execution内存的波峰波谷做精细调整。5. 大数据面试高频考点宽窄依赖、数据倾斜、Spark快在哪5.1 Spark为什么快别再回答“因为有内存”面试官特别喜欢问“Spark为什么比MapReduce快”但绝大多数人只会回答“因为Spark基于内存计算”。这个答案没有错但是远远不够。专业一点的回答应该包含三个层次。第一层内存计算确实减少了大量磁盘I/O尤其是迭代式算法第二层Spark引入的DAG调度引擎能够记住整个计算流程的依赖关系属于“懒执行”模式前面有多个操作时能自动合并优化第三层也就是我在前面提到过的Spark在很多操作上会做分区内聚合减少Shuffle数据量。你如果能从这三层来回答面试官会立刻知道你是真正写过代码做过调优的而不是背了八股文。5.2 RDD的宽窄依赖理解数据倾斜的理论基础另一个高频考点是“宽依赖和窄依赖的区别”。我习惯用生活化方式解释给新人听。窄依赖就像每个家长只负责自己的孩子父RDD的每个分区最多被子RDD的一个分区使用不需要跨节点传输代表操作有map、filter。宽依赖就像开家长会所有家长父分区都要把孩子数据送到同一个班级子分区代表操作有groupByKey、reduceByKey、join。这个区别如此重要是因为宽依赖会触发Shuffle而Shuffle是整个作业性能瓶颈的第一来源。数据倾斜某个Key的数据量远远超过其他Key之所以那么难搞本质也是因为Shuffle阶段把大量数据集中到了同一个节点上导致部分节点超时、部分节点空闲。针对数据倾斜一个最实用的处理思路就是“把倾斜的Key单独拆出来处理”。先把大数据量的Key过滤出来单独跑一个作业其余的Key正常跑最后合并两段结果。这个思路简单粗暴但解决90%的倾斜问题都够用了。5.3 学习路线建议先把一条链路走通再横向铺开最后聊一点路线问题。我自己比较推荐的学习节奏是先花一周时间把MapReduce WordCount完整走通理解Mapper、Reducer和Shuffle的基本流程再用一周时间把同样的算法用Spark RDD重新实现一遍对比两者的差异接着用Spark SQL处理一个真实数据集把数据清洗、聚合、TopN这些基础操作跑熟最后再回头去补理论比如宽窄依赖、内存模型、调度原理。这个路线的特点是每个阶段都能产出看得见的结果。每天都能有一小段可运行的代码“跑通”这种正反馈对学习来说是至关重要的。很多人的问题不是不努力是花大量时间在“看教程”而不是“写代码”导致原理讲得头头是道一打开终端就卡在环境变量上。数据算法这个东西动手永远比看教程重要十倍。像“源代码”这块我的建议是不要满足于把示例代码复制粘贴跑出结果。你可以做这样一个实验把WordCount的Reducer逻辑从“求和”改成“求平均数”或者把TopN从“取最大”改成“取最小”每次改一个点看看输出变化。改动和输出之间的因果链路建立起来之后你对分布式算法的理解深度会远超那些只刷过十遍教程的人。本文还有配套的精品资源点击获取
返回列表