多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

为什么Hadoop天生是分布式系统?单机瓶颈与架构设计解析

为什么Hadoop天生是分布式系统?单机瓶颈与架构设计解析 有个学大数据的朋友问我你当年第一次把Hadoop环境跑起来的时候有没有想过一个问题——数据放一台机器上不就完事了吗为什么非要把系统拆成一大堆节点还每个节点各管一摊说实话这个疑问我在刚接触Hadoop那会儿也有。那时候我在一台单机上刚把伪分布式环境搭通兴奋之余心里始终打鼓一台机器也能装能跑那分布式系统到底图什么后来我啃了HDFS的实现思路翻了MapReduce的执行流程又用三台破虚拟机把集群反复折腾了小半个月才真正想明白一个道理Hadoop从第一行代码起就是冲着“分布式”去的这不是设计偏好而是它要解决的问题决定了它只能这么长。这篇文章不打算科普Hadoop怎么安装而是想把“为什么分布式”这件事讲透。我会把分布式系统和传统集中式系统放在一张台面上对比聊聊Hadoop的核心设计是怎么把分布式的优势落地的还会分享一些我从伪分布式到真集群迁移过程中踩过的坑和排查思路。适合正在学Hadoop的初学者、准备大数据方向面试的人以及做技术选型时犹豫要不要上集群的团队。1. 为什么Hadoop天生就是分布式的先看看它要解决什么问题1.1 量变引发质变单机系统撞上的“四堵墙”很多人以为分布式是一个“高级设计”好像是大厂为了炫技才搞出来的东西。真不是。Hadoop诞生于2005年左右当时搜索引擎要处理的数据量已经到了单台机器彻底扛不住的地步。这里的“扛不住”不是一个模糊的感觉而是四条非常具体、非常物理的墙。第一堵是容量墙。假设一块普通磁盘是4TB一台高配服务器能插8块盘也就是32TB原始容量。听着不小对吧但如果要存储的数据有50TB甚至500TB呢尤其像网页索引、日志、用户行为这类数据增长是指数级的。单机容量的扩张速度永远追不上数据增长速度这是第一堵墙。第二堵是吞吐墙。就算数据压缩后只有50TB单机顺序读盘速度按200MB/s算把全部数据完整读一遍要72个小时以上。算法再快数据读不出来也白搭。更别说数据增长之后扫描一遍全量数据需要一周甚至更久这在搜索引擎场景里是不可接受的。第三堵是算力墙。单颗CPU的核心数和主频是有物理极限的。即使你写了一段完美并行的代码单机并行度也就几十个线程封顶。而海量数据处理任务比如全量排序、全量统计、建倒排索引本质上都是可以拆成成千上万个子任务并行执行的。这一类“数据并行”的工作负载天生就需要成百上千个计算单元。第四堵是成本墙。一台可以稳定扛住几十TB数据、带冗余电源、冗余磁盘阵列的高端服务器价格是普通商用服务器的几十倍。而且买来之后它的算力仍然不够还得继续买更贵的。反观分布式方案买一堆普通PC服务器或者虚拟机靠软件层面把它们组织成一个整体成本曲线完全不同。所以不要被“分布式”这个词唬住。Hadoop采用分布式不是因为它“高大上”而是数据规模和计算规模到了一个量级之后单机在物理上就走到了尽头。1.2 从“零件损坏”到“系统可靠”分布式容错的核心哲学说到分布式系统我特别想强调一个关键的思维转变集中式系统假设硬件是可靠的而分布式系统假设硬件随时会坏。一台精心维护的服务器磁盘可能会坏内存可能会报错电源可能会烧掉。在集中式架构里一台机器挂了整个系统就挂了所以你需要上双机热备、磁盘阵列、UPS本质上是靠“更高的硬件可靠性”来换系统可用性。分布式系统的思路完全不同它不再追求单台设备永不故障而是默认故障是常态然后通过软件层面的冗余和自动恢复让整个系统在部分组件坏掉的情况下继续对外服务。HDFS的3副本机制就是这个哲学的典型案例。同一份数据块会在三个不同节点各存一份任何一个节点宕机数据仍然可以从另外两个副本恢复。你不需要买“永不出错”的硬盘只需要准备足够多的普通硬盘然后让系统自己去处理损坏和重建。这种设计让大规模低成本存储成为可能。生活里有个很贴切的类比集中式系统像一个独当一面的大厨手艺好但你很难找到替代者他一旦生病餐厅就得关门分布式系统像一条流水线每个工位做一道简单工序任何一个工人请假临时调个人补上就能继续运转。效率不一定比大厨单干高但胜在稳定、可替换、可扩张。1.3 Hadoop不是第一个分布式系统却让分布式走进了普通团队分布式系统在Hadoop之前早就存在比如各类高性能计算集群、分布式数据库。但那些方案要么是闭源商业产品要么对硬件和运维的要求极高一般团队根本玩不起。Hadoop的历史地位在于它把Google发表的GFS和MapReduce两篇论文变成了开源实现让一个刚毕业的工程师也能用一堆普通机器搭建自己的分布式存储和计算平台。这里有个观点我想多说一句Hadoop的成功不仅仅是技术上的成功更是“降门槛”的成功。如果Hadoop需要价值百万的专用硬件需要博士级别的分布式理论背景才能部署它不可能成为今天大数据生态的事实标准。它选择了普通网络、普通服务器、Java语言刻意把分布式系统的高门槛拉低到大部分工程师可以够得着的程度。这种“平民化”路线恰恰也是它和传统集中式系统拉开差距的另一个维度。2. 分布式系统与集中式系统两种架构逻辑的正面对比2.1 先看一张总表核心差异对照为了把对比讲清楚我整理了一张表把两套架构在不同维度的差异直接摆出来。对比维度集中式系统Hadoop分布式系统节点角色一台主机负责全部计算和存储多台节点分工协作各自承担存储或计算任务扩展方式垂直扩展升级CPU、加内存、换更强硬件水平扩展直接增加节点数量扩展成本高端硬件价格昂贵存在明显上限普通硬件即可成本随节点数线性增长容错策略依赖硬件冗余和高可用方案软件层默认故障常态通过副本、心跳、自动恢复容错数据位置数据集中在本地磁盘或专用存储数据分块打散到所有节点计算尽量靠近数据一致性通常提供强一致性多数场景为强一致或最终一致性需按场景权衡并行能力受限于单机CPU和内存天然支持大规模并行任务可拆成海量小任务运维复杂度相对简单一个系统一个团队就能管需要管理集群、监控节点、处理网络和协调问题适用场景小数据量、强事务、低延迟海量数据存储与分析、批处理、日志处理等这张表可以帮你快速建立一个整体认知。但要注意这张表不是“分布式永远更好”的证明而是说明两种架构在海量数据场景下面临截然不同的取舍。下面分开细说。2.2 集中式系统的“舒适区”与天花板集中式系统并没有过时它至今仍然是绝大多数业务系统的首选比如传统的企业管理系统、金融交易系统、在线事务处理系统。这些场景有几个共同特征数据量在可控范围内单表几百万到几千万行业务逻辑强依赖事务和强一致性需要毫秒级响应数据虽然重要但体量远没有到“单机装不下”的程度。在这些场景里集中式系统的优势非常明显架构简单一个数据库实例加几台应用服务器就能跑起来事务处理能力强ACID特性可以放心依赖运维成本低出了问题只需要排查几条链路。老实说用一个分布式集群去跑一个几GB的订单系统纯粹是给自己找麻烦。集中式系统的天花板也在于它的“集中”。当数据量增长到几十TB甚至更大当计算任务复杂到单机CPU长期跑满你只能被迫走上“垂直扩展”这条路换更强的CPU、加更多内存、上更高端的存储阵列。问题是高端硬件的性能提升总有尽头而且价格极其不友好。你可能会发现把钱花在买一台“超级服务器”上的费用已经够买几十台普通服务器了。这就是集中式系统在大数据时代绕不过去的死结。2.3 分布式系统的三个杀手锏线性扩展、故障容错、数据本地性集中式系统靠单点能力打天下分布式系统靠组织能力立身。Hadoop这类分布式系统最核心的三个优势我认为是下面这三条。第一是横向扩展能力。数据量从10TB涨到100TB集中式系统的方案可能是换一台贵十倍的机器而分布式系统的方案只是多加几台节点。虽然集群规模变大之后管理的复杂度也会增加但存储和计算能力确实能随着节点数量增加而增长。这种“不够就加机器”的底气是集中式系统给不了的。第二是故障容错能力。集群有几十上百台节点某台机器宕机是常态但整体服务不中断。HDFS里的数据块有多个副本某个DataNode挂了NameNode会自动把它的副本标记为待复制并在其他节点重新生成副本。任务调度也一样某个节点上的Map任务失败了ApplicationMaster会把任务重新调度到其他节点。对用户来说集群偶尔有节点出问题其实“无感”。第三是数据本地性。这一点很多人容易忽略但在海量数据场景里至关重要。Hadoop框架的核心思想是“移动计算而不是移动数据”move computation to data。比如一个查询要扫描1TB的数据如果数据分散在10个节点上框架会把计算任务下发给每个节点让它们只处理本地的数据分片最后汇总结果。这样一来网络传输的数据量只有最终汇总的结果而不是把1TB原始数据全部搬到一台机器上算。在集中式系统里数据在哪儿计算就必须在哪儿或者你得忍受把海量数据搬到计算服务器上的巨额网络开销。2.4 代价也要讲清楚CAP、协调与运维复杂度分布式系统不是万能药它也有需要付出的代价。面试里如果你能主动说出来反而是明显的加分项。第一个代价是网络分区和一致性权衡这就要提到CAP理论。分布式系统在网络分区发生时必须在可用性和一致性之间做取舍。HDFS在部分场景下选择了最终一致性某些操作可能不能像单机数据库那样提供严格的事务保证。MapReduce本身是批处理模型也天然不适合需要毫秒级延迟的在线事务。第二个代价是节点间协调的复杂度。集群越大节点之间需要同步的状态就越多谁还活着、谁宕机了、谁该承担哪个任务、数据副本是否足够。为了把这些事情理清楚Hadoop生态里需要ZooKeeper这样的协调组件来管理元数据、选主、做分布式锁。你在实际搭建Hadoop集群时可能还要配置Hadoop和ZooKeeper的整合目的就是让NameNode的高可用、ResourceManager的主备切换有可靠的协调底座。这也就意味着你多了一个要学习和运维的组件。第三个代价是运维和排查的复杂度。单机系统出了问题登录机器看日志基本能定位分布式集群出了问题可能是网络抖动、时钟偏差、某节点磁盘写满、某个进程内存溢出你得学会从几十台机器里快速定位问题节点。这些成本在你决定是否用Hadoop时一定得考虑进去。3. HDFS怎么把“分布式”落到实处三个关键设计拆解3.1 128MB数据块为什么切得这么大理解Hadoop分布式设计最好的入口是HDFS因为它把分布式的思想体现在最底层的数据组织方式上。传统文件系统按“文件”为单位管理数据文件多大就占多大空间HDFS则把文件切成一个个block每个block默认128MB然后分散存储到不同的DataNode上。这里有个值得停下来想一想的点为什么默认是128MB而不是64KB或者1GB实际上块大小的选择是一系列权衡的结果。先说为什么不能太小。如果块太小比如4KB那么一个1TB的文件会被切成2.6亿个块NameNode作为管理元数据的“大脑”要为每一块维护一条元数据记录每一条记录又包含文件名、块ID、副本位置等若干字段。块数量越多NameNode的内存压力就越大最后的瓶颈会出现在“元数据管理”而不是“数据存储”上。另外启动一个Map任务需要时间如果块太小任务数会多到调度开销超过计算本身。再说为什么不能太大。如果块太大比如1GB以上那么一个文件只有几个块并行处理的粒度太粗很难充分利用集群的算力而且当某个节点宕机需要重新复制数据块时太大的块会让复制时间变得不可控。所以128MB是在“减少元数据压力”和“保证并行粒度”之间取得的一个平衡点这也是实践中沉淀出来的经验值。一个直观的计算1TB的文件在128MB块大小下会被切成8192个block分布在若干节点上。一个扫描任务可以拆成8192个Map任务由集群并行执行。这就是分布式系统“化整为零、并行处理”的核心动作。3.2 3副本背后容错从“单机冗余”变成“集群冗余”有了数据块下一步就是决定每个块存几份。HDFS默认的副本数是3。3副本的主要目的是容错但副本放在哪里其实也大有讲究这里面就涉及机架感知的概念。默认的副本放置策略大致是这样第一个副本放在客户端所在节点如果客户端不在集群内则随机选一个节点第二个副本放在与第一个副本不同机架的节点上第三个副本放在与第二个副本同机架的另一台节点上。这样设计的好处是即使整个机架断电或者网络故障数据仍然有至少一个副本在其他机架同时读取数据时可以通过机架感知选择“距离最近”的副本减少跨机架的网络传输。3副本带来的收益很直接。一是容错任何一台节点宕机数据仍然完整。二是有益于计算本地性Map任务可以优先调度到存有本地副本的节点上减少网络读。三是自动恢复当某个副本丢失时NameNode会检测到副本数低于预期然后在其他节点上重新复制一份。实际部署中有一些团队会把副本数调成2甚至调回1主要为了节省存储空间。但我建议在测试环境可以这样做生产环境至少要3份因为副本数不只是容错保障它也是数据恢复、任务调度的基础资源。3.3 NameNode肩负的“目录”职责单点怎么破分布式存储把数据块“打散”到了各个节点那么问题来了谁来记录“哪个块在哪台机器上”这在HDFS里由NameNode负责。NameNode维护整个文件系统的元数据包括目录结构、文件与块的映射、块与DataNode的映射。你可以把它理解成一本总目录数据本身分散在各地但“哪里有什么”这个信息集中在一处。这种主从架构的好处是设计简单清晰客户端只需要问NameNode一次“这个文件的数据在哪个DataNode”然后直接和DataNode通信读写数据整个链路高效直接。但代价是NameNode成了整个HDFS的单点如果它挂了整个文件系统就无法访问了。应对单点Hadoop经历了几个阶段的演进早期靠SecondaryNameNode定期合并edits日志但它只是辅助节点不是热备后来引入了基于QJM的NameNode高可用方案配合ZooKeeper进行故障自动转移。所谓Hadoop和ZooKeeper整合实战通常指的就是搭建NameNode HA让Active NameNode和Standby NameNode共享同一份编辑日志ZooKeeper负责监控状态并触发自动切换。这里我想特别纠正一个常见的误区SecondaryNameNode不是NameNode的热备节点它的职责是定期把edits日志和fsimage合并减轻NameNode的内存压力和恢复时间。很多人误以为有了SecondaryNameNodeNameNode挂了就能顶上真不是这样。所以生产环境要认真考虑HA的搭建而HA必然离不开ZooKeeper的协调。4. 亲手感受分布式从伪分布式搭建到集群迁移实录4.1 为什么要先从“伪分布式”开始纸上谈兵再多不如亲手把环境搭一遍。在只有一台电脑或者一台虚拟机的情况下你也可以完整地体验Hadoop的分布式逻辑办法就是搭“伪分布式”。所谓伪分布式pseudo-distributed就是一台机器上同时运行NameNode、DataNode、ResourceManager、NodeManager等所有守护进程。虽然它们跑在同一台机器上但进程是分离的配置和启动流程与真实集群完全一致。因为数据块、副本、心跳、任务调度这些机制仍然在工作你可以完整看到分布式系统内部是怎么协作的只是都在本机完成。我建议想学Hadoop的人都要亲手搭一遍伪分布式这比自己只读文档要有效得多。而且从单机伪分布式扩展到多节点集群时你会发现80%的配置思路是完全相通的到时候只是加节点、改hostname、配SSH的事情。4.2 核心配置文件的四个关键角色Hadoop的配置文件集中在etc/hadoop目录下常用的有四个各自管一摊core-site.xml定义文件系统方案和默认端口核心配置是fs.defaultFS和hadoop.tmp.dir。hdfs-site.xmlHDFS相关配置比如副本数、NameNode数据目录、DataNode数据目录、Web UI端口。mapred-site.xmlMapReduce框架配置核心是设置运行框架为YARN。yarn-site.xmlYARN资源调度配置包括ResourceManager地址和NodeManager的辅助服务。下面是我在搭建Hadoop 3.x伪分布式时使用过的一份最小配置可以直接参考!-- core-site.xml -- configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property property namehadoop.tmp.dir/name value/opt/hadoop/tmp/value /property /configuration!-- hdfs-site.xml -- configuration property namedfs.replication/name value1/value /property property namedfs.namenode.name.dir/name valuefile:///opt/hadoop/namenode/value /property property namedfs.datanode.data.dir/name valuefile:///opt/hadoop/datanode/value /property property namedfs.http.address/name value0.0.0.0:9870/value /property /configuration!-- mapred-site.xml -- configuration property namemapreduce.framework.name/name valueyarn/value /property /configuration!-- yarn-site.xml -- configuration property nameyarn.nodemanager.aux-services/name valuemapreduce_shuffle/value /property property nameyarn.resourcemanager.hostname/name valuelocalhost/value /property /configuration需要注意伪分布式环境下dfs.replication要设置为1。如果保持默认3DataNode会尝试在“其他节点”上创建副本而集群只有一台机器这个操作必然失败数据写入会一直卡住。这是伪分布式新手最容易踩的坑之一。4.3 启动与验证怎么确认它真的“动”起来了启动之前先确认几件事JDK版本是不是8以上JAVA_HOME有没有正确设置SSH能不能免密登录本机Hadoop启动脚本依赖SSH。然后执行NameNode格式化hdfs namenode -format格式化就是初始化NameNode的元数据目录。注意这句话执行几次要心里有数——每次格式化会生成新的clusterID如果DataNode和NameNode的clusterID不一致DataNode会启动失败。很多人在练习环境反复格式化导致DataNode起不来本质就是这个原因。格式化之后分别启动HDFS和YARNstart-dfs.sh start-yarn.sh然后执行jps你会看到5个关键进程NameNode DataNode SecondaryNameNode ResourceManager NodeManager这5个进程齐全就说明伪分布式基本启动成功了。接着访问http://localhost:9870在HDFS Web UI里可以看到你的NameNode状态、数据节点列表以及文件块信息。往HDFS上传一个文件然后在Web UI里查看块分布你会直观地看到数据被切成了128MB或更小的块并且分布在“数据节点”上。这个体验会比你读十篇文章都管用。4.4 迁移到真集群时踩过的坑伪分布式跑通之后很多人会像我一样热血上头直接上三节点集群然后就开始踩坑。我把几个典型问题整理成了一张速查表。常见现象可能原因排查思路DataNode起不来日志里clusterID不匹配多次格式化NameNode元数据ID不一致清空NameNode和DataNode的数据目录后重新格式化节点间无法通信SecondaryNameNode连接失败未配置hostname或hosts文件检查/etc/hosts让每个节点能通过主机名互相访问DataNode连不上NameNode9000端口被防火墙拦截检查防火墙和fs.defaultFS配置提交任务报NoClassDefFoundError: org/apache/hadoop/cryptoHadoop或相关组件版本不一致classpath缺失检查hadoop classpath输出确认相关jar包在classpath中NodeManager启动成功但无法注册ResourceManager地址配置不对检查yarn-site.xml中的yarn.resourcemanager.hostname其中NoClassDefFoundError这个问题我在配置Hive和Tez时遇到过好几次不只是Hadoop本身的问题。这类报错往往不是真的有类缺失而是版本冲突或者classpath没有把对应jar包包含进来。排查思路很简单先找到报错类的完整路径再在Hadoop安装目录下的jar包里确认是否存在然后用export HADOOP_CLASSPATH补上或者检查各组件版本是否匹配。另外不同Hadoop版本的少量API有差异如果你拿高版本的jar包跑到低版本的集群上也会出现类似问题。从伪分布式扩展到真实集群时还有一个操作细节每台节点都要有自己的hostname而且/etc/hosts里要写清楚所有节点的映射互相信任的SSH免密也要提前配好。我在第一次搭三节点集群时就是漏了其中一台机器的hosts配置导致DataNode之间一直找不到对方花了很长时间才发现是这么个不起眼的问题。分布式环境的很多坑都藏在基础配置里。5. 把“为什么”变成面试与架构选型里的底气5.1 面试官为什么爱问这个问题“Hadoop为什么设计为分布式系统”这个问题在大数据方向的面试中出现频率极高。面试官问它不是想要你背一个“因为数据大所以要分布式”的标准答案而是想考察三件事第一你有没有真正理解Hadoop的适用场景第二你能不能把一个技术决策背后的权衡讲清楚第三你在做架构选型时能不能判断某个场景到底该用集中式还是分布式。很多人回答这个问题时只强调“分布式好、能扩展、能容错”听起来像在背书。但如果你能结合HDFS的块设计、副本机制、数据本地性再补充分布式带来的协调成本、一致性问题这套回答的层次就会明显不一样因为它展示了你既能看到优点也能看到代价。5.2 一个能踩分的回答框架如果是我来回答“Hadoop为什么设计为分布式系统”我会按下面这个顺序展开先说业务背景大数据时代的核心矛盾是数据规模与单机处理能力不匹配数据量达到TB甚至PB级后单台机器在容量、吞吐、算力、成本四个方面都会遇到瓶颈。再说分布式如何化解矛盾把文件切块分散到多台节点把计算任务拆成并行子任务用横向扩展替代纵向升级用多副本容错替代单点依赖用数据本地性减少网络传输。然后转折说代价分布式引入了节点协调、网络通信、状态一致性和运维复杂度不是所有场景都适用比如低延迟事务还是集中式更合适。最后回归本质所以Hadoop选择分布式本质上是数据规模驱动的必然结果是架构对物理规律和成本约束的妥协与设计。这个框架的好处是“有观点、有层次、有完整度”。面试官听完至少能确认你不仅知道“是什么”还理解“为什么”。5.3 反向思考什么时候别用Hadoop既然讲了这么多分布式的好话最后也泼一盆冷水。Hadoop并适合所有场景尤其不适合以下三类第一类是小数据量场景。几GB甚至几十GB的数据用一台好一点的服务器加一个关系型数据库完全能搞定何必搭集群复杂通报分布式系统的启动和调度开销反而会让简单任务变慢。第二类是强实时场景。Hadoop的核心计算模型是批处理数据写入HDFS后再通过MapReduce或Spark进行分析整个过程分钟级别的延迟。如果业务需要毫秒级响应比如在线推荐、风控实时拦截这个架构就不对口应该考虑OLTP数据库或流式处理引擎。第三类是强事务、强一致性挑战极大的场景。HDFS先写后读的一致性模型用来做“一次写入、多次读取”的存储类任务很顺手但如果你需要复杂事务、行级更新或者对一致性要求极其严格那就不是Hadoop的设计目标。我见过一些团队为了“用大数据而用大数据”把只有几百GB的MySQL日志库搬到Hadoop集群里结果运维成本上去了查询反而更慢。技术选型这件事没有最好的系统只有最匹配的场景。理解Hadoop为什么是分布式的另一面就是理解什么时候它不应该出现。最后再分享一点个人经验。我在多个环境里搭过Hadoop集群从单机伪分布式到三节点、五节点踩过的坑说多不多说少不少。如果你想少走弯路我觉得最实用的一条建议是先在一台虚拟机里把伪分布式真实跑通上传几个文件看一眼Web UI上的块分布再扩展到多节点集群。这样你对“数据被切散、副本被复制、任务被调度”这些概念会有切身体感后面排查问题也会顺手很多。分布式系统的“优势”从来不是凭空给的它靠的是把每台机器可能的失败提前设计进系统的日常运行里。这也是它和集中式系统对比时最本质的魅力——它从不对单台机器抱有不切实际的期望。
返回列表