ARTICLE DETAIL

资讯详情

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

Hadoop+Spark大数据实战:从集群搭建到性能调优全指南

Hadoop+Spark大数据实战:从集群搭建到性能调优全指南 简介面向大数据开发与学习人群这份源代码包围绕Hadoop和Spark两大框架汇集了MapReduce、Spark SQL、Streaming及MLlib等场景的算法实现适合希望结合代码理解分布式计算原理的读者。包内共876个文件以Java和Scala源码为主辅以jar依赖、Shell运行脚本、Markdown说明文档及数据集样例可支撑本地编译与集群调试整个zip压缩包约204MB。目前已有682人学习下载属于上手即用的实战型资料。通过学习源码与配套数据可掌握词频统计、数据清洗、分类回归等典型任务同时了解HDFS读写、RDD算子调优、Spark任务提交等关键环节对系统提升Hadoop/Spark工程能力有直接帮助。 做大数据开发这几年我经手过的任务从几百MB的日志清洗到上TB的离线聚合都碰过踩过的坑一个比一个离谱。Hadoop和Spark这套组合到今天依然是离线处理的事实标准但网上资料要么是纯理论要么是讲基础API真正遇到“为什么我这集群跑得这么慢”“为什么格式化又失败”这类问题时翻半天也找不到靠谱答案。这篇文章我就基于自己的实战经验从环境搭建到源代码调优把能用得上的处理技巧、核心代码样例和排查思路完整梳理一遍。不管你是刚学大数据准备找工作的学生还是已经在集群上摸爬滚打想提升效率的工程师这篇都能给你些实在的参考。1. 内容整体设计与思路拆解1.1 Hadoop和Spark到底什么关系为什么要一起学很多新手一上来就搞混以为Spark是Hadoop的替代品其实两者是完全不同的定位。Hadoop的核心价值在于分布式文件系统HDFS它解决的是“数据存哪里、怎么存得下”的问题而Spark解决的是“数据怎么算得快”的问题。实际生产环境最常见的组合是HDFS存原始数据YARN做资源调度Spark作为计算引擎跑在这些数据上面。你完全可以不用Spark用Hadoop自带的MapReduce算但MR每次任务都要落盘延迟高得让人崩溃Spark把中间结果尽量放内存速度能快一个数量级。那为什么学Spark之前最好先碰一下Hadoop因为Spark的很多设计比如分区、shuffle、容错都是从MR那里进化过来的。你没见过MR的痛就很难理解Spark为什么有那些“奇怪”的配置参数比如分区的设定、shuffle时溢写文件的机制全是针对MR的短板做的改良。所以一个合理的路径是先用伪分布式把HDFS和YARN跑通再在它上面搭Spark最后才去啃源代码。1.2 为什么一定要看源代码纯调API写业务逻辑天花板很低。你只会在DataFrame上filter、groupBy遇到数据倾斜、OOM、Executor挂掉根本无从下手。但如果你读过Spark的RDD源码理解DAGScheduler怎么切分Stage、ShuffleManager怎么管中间文件很多性能问题自己就能推理出来。读源码不一定是为了二次开发更实际的价值是建立起“执行模型”的心智印象这在排查问题时比什么都有用。2. 核心细节解析与实操要点2.1 Hadoop伪分布式与全分布式搭建的关键差异很多教程上来就让你搭集群我建议新手先老老实实跑通伪分布式。你只有一台机器把NameNode、DataNode、ResourceManager都启动在本地一样可以感受完整的读写流程、YARN调度和任务提交逻辑排查问题也简单日志都在本地文件里。伪分布式需要改四个配置文件core-site.xml、hdfs-site.xml、mapred-site.xml、yarn-site.xml要特别注意fs.defaultFS要写hdfs://localhost:9000而不是默认的file:///否则HDFS根本不会生效。replication默认是3单机有一个DataNode这里要改成1不然块复制不了那么多副本会一直空等。全分布式和伪分布式的差别主要在两点一是规划主节点和从节点NameNode和ResourceManager放主节点DataNode和NodeManager放在从节点生产上为了高可用还会单独部署Zookeeper节点和JournalNode二是免密钥登录主节点要能ssh到所有从节点否则启动时没法远程拉起进程。这两个点每个都有人踩坑尤其是忘了配免密钥启动脚本一直卡在那里报Permission denied又慢又烦。2.2 格式化NameNode的一个大坑重复格式化失败“hadoop启动格式化失败”这个问题在社区里被问得极多。格式化NameNode的命令很简单就是hdfs namenode -format但它不是可以随便执行的。格式化相当于给文件系统清空重新做标记如果你之前已经format过再次执行会生成一个新的clusterId旧的DataNode上保存的还是以前的clusterId两边对不上启动的时候DataNode就会疯狂报错一直尝试连接NameNode然后失败。解决思路有两条。如果你确认数据没用最干脆的是把dfs.name.dir和dfs.data.dir指向的目录全部删掉还有tmp文件夹也清掉重新格式化一次如果集群还在运行、不想丢数据那就别用format用hdfs namenode -recover去恢复元数据但恢复过程比较考验日志分析能力。反正我的习惯是只有在刚部署、确认没有业务数据时才会执行format跑起业务后绝不在线格式化。2.3 Hadoop和Zookeeper整合实战HA高可用单NameNode的集群有一个致命问题NameNode挂了整个HDFS就不可用了HDFS作为存储底座一旦停摆上层Spark再快也没有意义。HA方案是部署两个NameNode一台Active一台Standby共享或镜像EditLogZookeeper负责故障时自动切换。生产环境还需要JournalNode集群来同步EditLog这个过程有点繁琐要配置zoo.cfg、hadoop-hdfs-ha.xml、core-site.xml里的nameservice等等还要手动初始化journalnode和zkfc。我第一次搭HA的时候栽在ZKFC上启动后一直在报连接Zookeeper超时后来发现是防火墙没放22002端口网络策略问题比配置问题还要隐蔽排查时间翻倍。2.4 Spark集群搭建与内存模型Spark本身不负责存储它需要跑在一个资源管理器上。本地学习可以只用local模式跑单机但碰真实数据就得搭集群。常见两种方式一种是Standalone模式Spark自己管资源另一种是Spark on YARN接在Hadoop的YARN上。生产多用后者好处是资源可以统一调度Hadoop和Spark任务混跑不会出现YARN集群闲着、Spark集群却挤爆的情况。Spark配置文件里最重要的是spark-env.sh要指定JAVA_HOME和HADOOP_CONF_DIR。启动history-server还要配置spark.history.fs.logDirectory否则作业跑完了UI一刷新全是空的查不到历史日志。内存模型这块新版Spark把堆内内存分成了Reserved、Execution和Storage三个部分。Reserved占300MB留给系统内部用Execution是给shuffle、join、aggregation这些操作用的Storage是给缓存数据和广播变量用的。两者之间有动态借用机制你缓存的数据多了执行区的内存可以抢占但反过来执行区急着要内存时缓存数据会被淘汰。很多人问“spark.executor.memory设了4G怎么堆内可用只有3G多”就是因为Reserved和用户代码还有一部分内存开销不是Bug是设计。实际调参数时executor内存不宜超过YARN容器上限核数也不宜分配太高避免IO密集任务之间抢带宽。3. 实操过程与核心环节实现3.1 源代码对比Hadoop版WordCount与Spark版WordCount先看一段最经典的入门代码用MapReduce实现词频统计。代码本身逻辑不复杂但你看完就知道为什么MR跑迭代任务那么痛苦。public class WordCount { 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); } } } 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); } } }再看Spark版本的实现同样的逻辑写起来简洁得多。关键是理解Spark的变换操作是惰性的只有遇到action算子才会真正提交任务map、flatMap、filterByKey这些变换只是构建了一个DAG执行图。val textFile sc.textFile(hdfs://namenode:9000/input) val counts textFile.flatMap(line line.split( )) .map(word (word, 1)) .reduceByKey(_ _) counts.saveAsTextFile(hdfs://namenode:9000/output)绕开复杂度看本质MR版本每个map和reduce之间都要把中间结果写到磁盘而Spark的reduceByKey在同一个Exector内会优先内存聚合只在必须shuffle时才落盘这个差异也正是两者性能差距的主要来源。3.2 一个真实的Spark SQL数据分析案例我接过的数据分析需求大部分都不需要写RDD算子直接用Spark SQL最顺手。比如有一份用户行为日志字段包括user_id、action、item_id、timestamp要统计每天各action类型的PV、UV以及人均操作次数用DataFrame API加SQL代码量可以压得非常小。读取JSON文件后用df.createOrReplaceTempView注册成临时视图然后直接写SQLgroup by分区字段最后写回Parquet格式结果集。Parquet是列式存储后续按字段查询时能大幅减少IO这是实际项目里很实用的小技巧。val df spark.read.json(hdfs:///data/user_logs) df.createOrReplaceTempView(logs) val result spark.sql( SELECT date, action, COUNT(*) AS pv, COUNT(DISTINCT user_id) AS uv FROM logs GROUP BY date, action ) result.write.mode(overwrite).parquet(hdfs:///result)需要注意的一点是COUNT(DISTINCT ...)在数据量极大的情况下很容易成为性能瓶颈因为精准去重需要全量shuffle。如果业务上允许近似值可以用approx_count_distinct代替它在底层做HyperLogLog估计能省大量时间。3.3 数据倾斜优化源代码思维指导下的实战方案数据倾斜是Spark跑得慢的头号元凶症状就是某个executor长期卡住其他executor早完成任务在等它。原因通常是shuffle时key分布不均匀比如热点用户、热点商品几十亿条数据里一个key占了30%reduce端那一个partition收到的数据量比其他partition大几个数量级。解决思路先想到的是加盐salting。对于热点key在map阶段先给key加一个随机前缀比如0到n之间的随机数把一个大key拆成n个小key让它分散到不同的reduce分区去处理计算完后再去掉前缀做一次聚合。但这方法不能乱用适用于聚合类算子对join要小心的关联逻辑。更简单直接的方法是给join的小表广播出去用broadcast join避免shuffle只有当两边都是大表时再考虑加盐拆散热点key。这些优化手段如果只背API是不好理解的真的要回到shuffle机制本身去推敲。4. 常见问题与排查技巧实录4.1 Spark on YARNCPU核心数怎么只有一个群里经常有人问“spark on yarnexecutor在yarn上跑的时候每个container只分配一个vcore资源配置明明改过为什么不管用”这个问题十有八九出在YARN的资源发现上。虚拟机默认情况下YARN的nodemanager不认识物理核和逻辑核的区别会自动把cpu-vcores识别成1。解决办法是在yarn-site.xml里设置yarn.nodemanager.resource.cpu-vcores为你期望的值比如8同时设置yarn.nodemanager.resource.detect-hardware-capabilities为true之后重启NodeManager才能生效。如果明明设置了还是没有效果那去看提交任务时有没有用--executor-cores参数强制覆盖。注意Spark任务提交参数的优先级是最高的配置文件里的默认值会被它盖掉别在这上面浪费时间。4.2 Spark内存与OOM排查路线Execrtor OOM的报错很多种但排查路径比较固定。先看是执行内存不足还是存储内存不足打开Spark UI的Executors页面看Shuffle Spill条和Storage Memory使用量如果磁盘溢写量特别大说明执行内存不够用了需要增加executor内存或减小并行度。如果是缓存数据把存储区占满了考虑缓存级别是否该换成MEMORY_ONLY_SER序列化后能压缩空间但会增加CPU开销。如果是Driver端OOM十有八九是collect()把所有结果拉到Driver内存里改成分批或直接写到HDFS就行。4.3 配置格式与日志的小坑Spark在启动时会输出一行日志“using sparks default log4j profile: org/apache/spark/log4j-defaults.propert”很多人看到这个以为出错了其实这只是告诉你没找到自定义log4j配置文件在用默认的。如果不希望INFO日志刷屏去$SPARK_HOME/conf复制一份log4j.properties.template重命名为log4j.properties设置成WARN级别即可。类似的还有Hadoop常见的一个报错“localhost:9000: Connection refused”基本都是NameNode没起来或core-site.xml配置没写对。5. 面试高频考题与学习路径建议5.1 Hadoop面试题里一定要会讲的几个点面试官爱问的点常年不变shuffle过程到底发生了什么、小文件问题怎么处理、NameNode压力太大怎么缓解、数据倾斜怎么解决。这些没有标准答案但在限定条件下有最优解。比如小文件问题底层原因是HDFS的特性每个文件都有对应的元数据占NameNode内存大量小文件会把NameNode堆内存撑爆。解决思路是输入端做合并把多个小文件合并成大文件或者用CombineFileInputFormat等自定义输入格式Spark任务写回数据时也尽量用coalesce()控制分区数别生成一堆碎文件。5.2 Spark面试题的逻辑层次Spark的面试题其实很能区分水平。初级会问RDD怎么创建、常用算子有哪些中级会问宽依赖和窄依赖的区别、Stage是怎么划分的、Spark为什么比MapReduce快高级会问什么时候用cache什么时候用checkpoint两者有什么区别shuffle调优参数怎么配数据倾斜有哪些更好的处理方式。我的建议是面试前自己画一张DAG切分图亲手写一个多次shuffle的任务然后对着Spark UI看Stage划分和shuffle read/write的数据量比背十遍八股文的效率高得多。5.3 推荐的进阶路线从环境搭建开始先用一份几百MB的日志练习HDFS操作文件、YARN跑MR、Spark跑统计接着读RDD源码里的getDependencies和getPartitions方法理解依赖和分区的概念。之后可以系统看一部讲Spark原理的书籍跟着代码走一遍DAGScheduler的事件循环最后做一个小型项目比如用户行为分析加上异常检测把Hadoop、Spark、Spark SQL、调优、源码排查全部串起来。根据我个人经验大数据这个方向动手实操的价值远大于看教程。同一个问题今天踩一遍坑比看别人写十遍避坑指南都记住得牢。先跑通再说跑完后多问几个为什么慢慢就能变成别人眼里“对Spark底层很懂”的那个人。本文还有配套的精品资源点击获取
返回列表