hdfs_design.html 38 KB


  1. <!DOCTYPE html PUBLIC "-//W3C//DTD HTML 4.01 Transitional//EN" "http://www.w3.org/TR/html4/loose.dtd">
  2. <html>
  3. <head>
  4. <META http-equiv="Content-Type" content="text/html; charset=UTF-8">
  5. <meta content="Apache Forrest" name="Generator">
  6. <meta name="Forrest-version" content="0.8">
  7. <meta name="Forrest-skin-name" content="pelt">
  8. <title>
  9. Hadoop分布式文件系统:架构和设计
  10. </title>
  11. <link type="text/css" href="skin/basic.css" rel="stylesheet">
  12. <link media="screen" type="text/css" href="skin/screen.css" rel="stylesheet">
  13. <link media="print" type="text/css" href="skin/print.css" rel="stylesheet">
  14. <link type="text/css" href="skin/profile.css" rel="stylesheet">
  15. <script src="skin/getBlank.js" language="javascript" type="text/javascript"></script><script src="skin/getMenu.js" language="javascript" type="text/javascript"></script><script src="skin/fontsize.js" language="javascript" type="text/javascript"></script>
  16. <link rel="shortcut icon" href="images/favicon.ico">
  17. </head>
  18. <body onload="init()">
  19. <script type="text/javascript">ndeSetTextSize();</script>
  20. <div id="top">
  21. <!--+
  22. |breadtrail
  23. +-->
  24. <div class="breadtrail">
  25. <a href="http://www.apache.org/">Apache</a> &gt; <a href="http://hadoop.apache.org/">Hadoop</a> &gt; <a href="http://hadoop.apache.org/core/">Core</a><script src="skin/breadcrumbs.js" language="JavaScript" type="text/javascript"></script>
  26. </div>
  27. <!--+
  28. |header
  29. +-->
  30. <div class="header">
  31. <!--+
  32. |start group logo
  33. +-->
  34. <div class="grouplogo">
  35. <a href="http://hadoop.apache.org/"><img class="logoImage" alt="Hadoop" src="images/hadoop-logo.jpg" title="Apache Hadoop"></a>
  36. </div>
  37. <!--+
  38. |end group logo
  39. +-->
  40. <!--+
  41. |start Project Logo
  42. +-->
  43. <div class="projectlogo">
  44. <a href="http://hadoop.apache.org/core/"><img class="logoImage" alt="Hadoop" src="images/core-logo.gif" title="Scalable Computing Platform"></a>
  45. </div>
  46. <!--+
  47. |end Project Logo
  48. +-->
  49. <!--+
  50. |start Search
  51. +-->
  52. <div class="searchbox">
  53. <form action="http://www.google.com/search" method="get" class="roundtopsmall">
  54. <input value="hadoop.apache.org" name="sitesearch" type="hidden"><input onFocus="getBlank (this, 'Search the site with google');" size="25" name="q" id="query" type="text" value="Search the site with google">&nbsp;
  55. <input name="Search" value="Search" type="submit">
  56. </form>
  57. </div>
  58. <!--+
  59. |end search
  60. +-->
  61. <!--+
  62. |start Tabs
  63. +-->
  64. <ul id="tabs">
  65. <li>
  66. <a class="unselected" href="http://hadoop.apache.org/core/">项目</a>
  67. </li>
  68. <li>
  69. <a class="unselected" href="http://wiki.apache.org/hadoop">维基</a>
  70. </li>
  71. <li class="current">
  72. <a class="selected" href="index.html">Hadoop 0.18文档</a>
  73. </li>
  74. </ul>
  75. <!--+
  76. |end Tabs
  77. +-->
  78. </div>
  79. </div>
  80. <div id="main">
  81. <div id="publishedStrip">
  82. <!--+
  83. |start Subtabs
  84. +-->
  85. <div id="level2tabs"></div>
  86. <!--+
  87. |end Endtabs
  88. +-->
  89. <script type="text/javascript"><!--
  90. document.write("Last Published: " + document.lastModified);
  91. // --></script>
  92. </div>
  93. <!--+
  94. |breadtrail
  95. +-->
  96. <div class="breadtrail">
  97. &nbsp;
  98. </div>
  99. <!--+
  100. |start Menu, mainarea
  101. +-->
  102. <!--+
  103. |start Menu
  104. +-->
  105. <div id="menu">
  106. <div onclick="SwitchMenu('menu_selected_1.1', 'skin/')" id="menu_selected_1.1Title" class="menutitle" style="background-image: url('skin/images/chapter_open.gif');">文档</div>
  107. <div id="menu_selected_1.1" class="selectedmenuitemgroup" style="display: block;">
  108. <div class="menuitem">
  109. <a href="index.html">概述</a>
  110. </div>
  111. <div class="menuitem">
  112. <a href="quickstart.html">快速入门</a>
  113. </div>
  114. <div class="menuitem">
  115. <a href="cluster_setup.html">集群搭建</a>
  116. </div>
  117. <div class="menupage">
  118. <div class="menupagetitle">HDFS构架设计</div>
  119. </div>
  120. <div class="menuitem">
  121. <a href="hdfs_user_guide.html">HDFS使用指南</a>
  122. </div>
  123. <div class="menuitem">
  124. <a href="hdfs_permissions_guide.html">HDFS权限指南</a>
  125. </div>
  126. <div class="menuitem">
  127. <a href="hdfs_quota_admin_guide.html">HDFS配额管理指南</a>
  128. </div>
  129. <div class="menuitem">
  130. <a href="commands_manual.html">命令手册</a>
  131. </div>
  132. <div class="menuitem">
  133. <a href="hdfs_shell.html">FS Shell使用指南</a>
  134. </div>
  135. <div class="menuitem">
  136. <a href="distcp.html">DistCp使用指南</a>
  137. </div>
  138. <div class="menuitem">
  139. <a href="mapred_tutorial.html">Map-Reduce教程</a>
  140. </div>
  141. <div class="menuitem">
  142. <a href="native_libraries.html">Hadoop本地库</a>
  143. </div>
  144. <div class="menuitem">
  145. <a href="streaming.html">Streaming</a>
  146. </div>
  147. <div class="menuitem">
  148. <a href="hadoop_archives.html">Hadoop Archives</a>
  149. </div>
  150. <div class="menuitem">
  151. <a href="hod.html">Hadoop On Demand</a>
  152. </div>
  153. <div class="menuitem">
  154. <a href="http://hadoop.apache.org/core/docs/r0.18.2/api/index.html">API参考</a>
  155. </div>
  156. <div class="menuitem">
  157. <a href="http://hadoop.apache.org/core/docs/r0.18.2/jdiff/changes.html">API Changes</a>
  158. </div>
  159. <div class="menuitem">
  160. <a href="http://wiki.apache.org/hadoop/">维基</a>
  161. </div>
  162. <div class="menuitem">
  163. <a href="http://wiki.apache.org/hadoop/FAQ">常见问题</a>
  164. </div>
  165. <div class="menuitem">
  166. <a href="http://hadoop.apache.org/core/mailing_lists.html">邮件列表</a>
  167. </div>
  168. <div class="menuitem">
  169. <a href="http://hadoop.apache.org/core/docs/r0.18.2/releasenotes.html">发行说明</a>
  170. </div>
  171. <div class="menuitem">
  172. <a href="http://hadoop.apache.org/core/docs/r0.18.2/changes.html">变更日志</a>
  173. </div>
  174. </div>
  175. <div id="credit"></div>
  176. <div id="roundbottom">
  177. <img style="display: none" class="corner" height="15" width="15" alt="" src="skin/images/rc-b-l-15-1body-2menu-3menu.png"></div>
  178. <!--+
  179. |alternative credits
  180. +-->
  181. <div id="credit2"></div>
  182. </div>
  183. <!--+
  184. |end Menu
  185. +-->
  186. <!--+
  187. |start content
  188. +-->
  189. <div id="content">
  190. <div title="Portable Document Format" class="pdflink">
  191. <a class="dida" href="hdfs_design.pdf"><img alt="PDF -icon" src="skin/images/pdfdoc.gif" class="skin"><br>
  192. PDF</a>
  193. </div>
  194. <h1>
  195. Hadoop分布式文件系统:架构和设计
  196. </h1>
  197. <div id="minitoc-area">
  198. <ul class="minitoc">
  199. <li>
  200. <a href="#%E5%BC%95%E8%A8%80"> 引言 </a>
  201. </li>
  202. <li>
  203. <a href="#%E5%89%8D%E6%8F%90%E5%92%8C%E8%AE%BE%E8%AE%A1%E7%9B%AE%E6%A0%87"> 前提和设计目标 </a>
  204. <ul class="minitoc">
  205. <li>
  206. <a href="#%E7%A1%AC%E4%BB%B6%E9%94%99%E8%AF%AF"> 硬件错误 </a>
  207. </li>
  208. <li>
  209. <a href="#%E6%B5%81%E5%BC%8F%E6%95%B0%E6%8D%AE%E8%AE%BF%E9%97%AE"> 流式数据访问 </a>
  210. </li>
  211. <li>
  212. <a href="#%E5%A4%A7%E8%A7%84%E6%A8%A1%E6%95%B0%E6%8D%AE%E9%9B%86"> 大规模数据集 </a>
  213. </li>
  214. <li>
  215. <a href="#%E7%AE%80%E5%8D%95%E7%9A%84%E4%B8%80%E8%87%B4%E6%80%A7%E6%A8%A1%E5%9E%8B"> 简单的一致性模型 </a>
  216. </li>
  217. <li>
  218. <a href="#%E2%80%9C%E7%A7%BB%E5%8A%A8%E8%AE%A1%E7%AE%97%E6%AF%94%E7%A7%BB%E5%8A%A8%E6%95%B0%E6%8D%AE%E6%9B%B4%E5%88%92%E7%AE%97%E2%80%9D"> &ldquo;移动计算比移动数据更划算&rdquo; </a>
  219. </li>
  220. <li>
  221. <a href="#%E5%BC%82%E6%9E%84%E8%BD%AF%E7%A1%AC%E4%BB%B6%E5%B9%B3%E5%8F%B0%E9%97%B4%E7%9A%84%E5%8F%AF%E7%A7%BB%E6%A4%8D%E6%80%A7"> 异构软硬件平台间的可移植性 </a>
  222. </li>
  223. </ul>
  224. </li>
  225. <li>
  226. <a href="#Namenode+%E5%92%8C+Datanode"> Namenode 和 Datanode </a>
  227. </li>
  228. <li>
  229. <a href="#%E6%96%87%E4%BB%B6%E7%B3%BB%E7%BB%9F%E7%9A%84%E5%90%8D%E5%AD%97%E7%A9%BA%E9%97%B4+%28namespace%29"> 文件系统的名字空间 (namespace) </a>
  230. </li>
  231. <li>
  232. <a href="#%E6%95%B0%E6%8D%AE%E5%A4%8D%E5%88%B6"> 数据复制 </a>
  233. <ul class="minitoc">
  234. <li>
  235. <a href="#%E5%89%AF%E6%9C%AC%E5%AD%98%E6%94%BE%3A+%E6%9C%80%E6%9C%80%E5%BC%80%E5%A7%8B%E7%9A%84%E4%B8%80%E6%AD%A5"> 副本存放: 最最开始的一步 </a>
  236. </li>
  237. <li>
  238. <a href="#%E5%89%AF%E6%9C%AC%E9%80%89%E6%8B%A9"> 副本选择 </a>
  239. </li>
  240. <li>
  241. <a href="#%E5%AE%89%E5%85%A8%E6%A8%A1%E5%BC%8F"> 安全模式 </a>
  242. </li>
  243. </ul>
  244. </li>
  245. <li>
  246. <a href="#%E6%96%87%E4%BB%B6%E7%B3%BB%E7%BB%9F%E5%85%83%E6%95%B0%E6%8D%AE%E7%9A%84%E6%8C%81%E4%B9%85%E5%8C%96"> 文件系统元数据的持久化 </a>
  247. </li>
  248. <li>
  249. <a href="#%E9%80%9A%E8%AE%AF%E5%8D%8F%E8%AE%AE"> 通讯协议 </a>
  250. </li>
  251. <li>
  252. <a href="#%E5%81%A5%E5%A3%AE%E6%80%A7"> 健壮性 </a>
  253. <ul class="minitoc">
  254. <li>
  255. <a href="#%E7%A3%81%E7%9B%98%E6%95%B0%E6%8D%AE%E9%94%99%E8%AF%AF%EF%BC%8C%E5%BF%83%E8%B7%B3%E6%A3%80%E6%B5%8B%E5%92%8C%E9%87%8D%E6%96%B0%E5%A4%8D%E5%88%B6"> 磁盘数据错误,心跳检测和重新复制 </a>
  256. </li>
  257. <li>
  258. <a href="#%E9%9B%86%E7%BE%A4%E5%9D%87%E8%A1%A1"> 集群均衡 </a>
  259. </li>
  260. <li>
  261. <a href="#%E6%95%B0%E6%8D%AE%E5%AE%8C%E6%95%B4%E6%80%A7"> 数据完整性 </a>
  262. </li>
  263. <li>
  264. <a href="#%E5%85%83%E6%95%B0%E6%8D%AE%E7%A3%81%E7%9B%98%E9%94%99%E8%AF%AF"> 元数据磁盘错误 </a>
  265. </li>
  266. <li>
  267. <a href="#%E5%BF%AB%E7%85%A7"> 快照 </a>
  268. </li>
  269. </ul>
  270. </li>
  271. <li>
  272. <a href="#%E6%95%B0%E6%8D%AE%E7%BB%84%E7%BB%87"> 数据组织 </a>
  273. <ul class="minitoc">
  274. <li>
  275. <a href="#%E6%95%B0%E6%8D%AE%E5%9D%97"> 数据块 </a>
  276. </li>
  277. <li>
  278. <a href="#Staging"> Staging </a>
  279. </li>
  280. <li>
  281. <a href="#%E6%B5%81%E6%B0%B4%E7%BA%BF%E5%A4%8D%E5%88%B6"> 流水线复制 </a>
  282. </li>
  283. </ul>
  284. </li>
  285. <li>
  286. <a href="#%E5%8F%AF%E8%AE%BF%E9%97%AE%E6%80%A7"> 可访问性 </a>
  287. <ul class="minitoc">
  288. <li>
  289. <a href="#DFSShell"> DFSShell </a>
  290. </li>
  291. <li>
  292. <a href="#DFSAdmin"> DFSAdmin </a>
  293. </li>
  294. <li>
  295. <a href="#%E6%B5%8F%E8%A7%88%E5%99%A8%E6%8E%A5%E5%8F%A3"> 浏览器接口 </a>
  296. </li>
  297. </ul>
  298. </li>
  299. <li>
  300. <a href="#%E5%AD%98%E5%82%A8%E7%A9%BA%E9%97%B4%E5%9B%9E%E6%94%B6"> 存储空间回收 </a>
  301. <ul class="minitoc">
  302. <li>
  303. <a href="#%E6%96%87%E4%BB%B6%E7%9A%84%E5%88%A0%E9%99%A4%E5%92%8C%E6%81%A2%E5%A4%8D"> 文件的删除和恢复 </a>
  304. </li>
  305. <li>
  306. <a href="#%E5%87%8F%E5%B0%91%E5%89%AF%E6%9C%AC%E7%B3%BB%E6%95%B0"> 减少副本系数 </a>
  307. </li>
  308. </ul>
  309. </li>
  310. <li>
  311. <a href="#%E5%8F%82%E8%80%83%E8%B5%84%E6%96%99"> 参考资料 </a>
  312. </li>
  313. </ul>
  314. </div>
  315. <a name="N10014"></a><a name="%E5%BC%95%E8%A8%80"></a>
  316. <h2 class="h3"> 引言 </h2>
  317. <div class="section">
  318. <p>
  319. Hadoop分布式文件系统(<acronym title="Hadoop分布式文件系统">HDFS</acronym>)被设计成适合运行在通用硬件(commodity hardware)上的分布式文件系统。它和现有的分布式文件系统有很多共同点。但同时,它和其他的分布式文件系统的区别也是很明显的。HDFS是一个高度容错性的系统,适合部署在廉价的机器上。HDFS能提供高吞吐量的数据访问,非常适合大规模数据集上的应用。HDFS放宽了一部分POSIX约束,来实现流式读取文件系统数据的目的。HDFS在最开始是作为Apache Nutch搜索引擎项目的基础架构而开发的。HDFS是Apache Hadoop Core项目的一部分。这个项目的地址是<a href="http://hadoop.apache.org/core/">http://hadoop.apache.org/core/</a>。
  320. </p>
  321. </div>
  322. <a name="N10026"></a><a name="%E5%89%8D%E6%8F%90%E5%92%8C%E8%AE%BE%E8%AE%A1%E7%9B%AE%E6%A0%87"></a>
  323. <h2 class="h3"> 前提和设计目标 </h2>
  324. <div class="section">
  325. <a name="N1002C"></a><a name="%E7%A1%AC%E4%BB%B6%E9%94%99%E8%AF%AF"></a>
  326. <h3 class="h4"> 硬件错误 </h3>
  327. <p>
  328. 硬件错误是常态而不是异常。HDFS可能由成百上千的服务器所构成,每个服务器上存储着文件系统的部分数据。我们面对的现实是构成系统的组件数目是巨大的,而且任一组件都有可能失效,这意味着总是有一部分HDFS的组件是不工作的。因此错误检测和快速、自动的恢复是HDFS最核心的架构目标。
  329. </p>
  330. <a name="N10036"></a><a name="%E6%B5%81%E5%BC%8F%E6%95%B0%E6%8D%AE%E8%AE%BF%E9%97%AE"></a>
  331. <h3 class="h4"> 流式数据访问 </h3>
  332. <p>
  333. 运行在HDFS上的应用和普通的应用不同,需要流式访问它们的数据集。HDFS的设计中更多的考虑到了数据批处理,而不是用户交互处理。比之数据访问的低延迟问题,更关键的在于数据访问的高吞吐量。POSIX标准设置的很多硬性约束对HDFS应用系统不是必需的。为了提高数据的吞吐量,在一些关键方面对POSIX的语义做了一些修改。
  334. </p>
  335. <a name="N10040"></a><a name="%E5%A4%A7%E8%A7%84%E6%A8%A1%E6%95%B0%E6%8D%AE%E9%9B%86"></a>
  336. <h3 class="h4"> 大规模数据集 </h3>
  337. <p>
  338. 运行在HDFS上的应用具有很大的数据集。HDFS上的一个典型文件大小一般都在G字节至T字节。因此,HDFS被调节以支持大文件存储。它应该能提供整体上高的数据传输带宽,能在一个集群里扩展到数百个节点。一个单一的HDFS实例应该能支撑数以千万计的文件。
  339. </p>
  340. <a name="N1004A"></a><a name="%E7%AE%80%E5%8D%95%E7%9A%84%E4%B8%80%E8%87%B4%E6%80%A7%E6%A8%A1%E5%9E%8B"></a>
  341. <h3 class="h4"> 简单的一致性模型 </h3>
  342. <p>
  343. HDFS应用需要一个&ldquo;一次写入多次读取&rdquo;的文件访问模型。一个文件经过创建、写入和关闭之后就不需要改变。这一假设简化了数据一致性问题,并且使高吞吐量的数据访问成为可能。Map/Reduce应用或者网络爬虫应用都非常适合这个模型。目前还有计划在将来扩充这个模型,使之支持文件的附加写操作。
  344. </p>
  345. <a name="N10058"></a><a name="%E2%80%9C%E7%A7%BB%E5%8A%A8%E8%AE%A1%E7%AE%97%E6%AF%94%E7%A7%BB%E5%8A%A8%E6%95%B0%E6%8D%AE%E6%9B%B4%E5%88%92%E7%AE%97%E2%80%9D"></a>
  346. <h3 class="h4"> &ldquo;移动计算比移动数据更划算&rdquo; </h3>
  347. <p>
  348. 一个应用请求的计算,离它操作的数据越近就越高效,在数据达到海量级别的时候更是如此。因为这样就能降低网络阻塞的影响,提高系统数据的吞吐量。将计算移动到数据附近,比之将数据移动到应用所在显然更好。HDFS为应用提供了将它们自己移动到数据附近的接口。
  349. </p>
  350. <a name="N10062"></a><a name="%E5%BC%82%E6%9E%84%E8%BD%AF%E7%A1%AC%E4%BB%B6%E5%B9%B3%E5%8F%B0%E9%97%B4%E7%9A%84%E5%8F%AF%E7%A7%BB%E6%A4%8D%E6%80%A7"></a>
  351. <h3 class="h4"> 异构软硬件平台间的可移植性 </h3>
  352. <p>
  353. HDFS在设计的时候就考虑到平台的可移植性。这种特性方便了HDFS作为大规模数据应用平台的推广。
  354. </p>
  355. </div>
  356. <a name="N1006D"></a><a name="Namenode+%E5%92%8C+Datanode"></a>
  357. <h2 class="h3"> Namenode 和 Datanode </h2>
  358. <div class="section">
  359. <p>
  360. HDFS采用master/slave架构。一个HDFS集群是由一个Namenode和一定数目的Datanodes组成。Namenode是一个中心服务器,负责管理文件系统的名字空间(namespace)以及客户端对文件的访问。集群中的Datanode一般是一个节点一个,负责管理它所在节点上的存储。HDFS暴露了文件系统的名字空间,用户能够以文件的形式在上面存储数据。从内部看,一个文件其实被分成一个或多个数据块,这些块存储在一组Datanode上。Namenode执行文件系统的名字空间操作,比如打开、关闭、重命名文件或目录。它也负责确定数据块到具体Datanode节点的映射。Datanode负责处理文件系统客户端的读写请求。在Namenode的统一调度下进行数据块的创建、删除和复制。
  361. </p>
  362. <div id="" style="text-align: center;">
  363. <img id="" class="figure" alt="HDFS 架构" src="images/hdfsarchitecture.gif"></div>
  364. <p>
  365. Namenode和Datanode被设计成可以在普通的商用机器上运行。这些机器一般运行着GNU/Linux操作系统(<acronym title="操作系统">OS</acronym>)。HDFS采用Java语言开发,因此任何支持Java的机器都可以部署Namenode或Datanode。由于采用了可移植性极强的Java语言,使得HDFS可以部署到多种类型的机器上。一个典型的部署场景是一台机器上只运行一个Namenode实例,而集群中的其它机器分别运行一个Datanode实例。这种架构并不排斥在一台机器上运行多个Datanode,只不过这样的情况比较少见。
  366. </p>
  367. <p>
  368. 集群中单一Namenode的结构大大简化了系统的架构。Namenode是所有HDFS元数据的仲裁者和管理者,这样,用户数据永远不会流过Namenode。
  369. </p>
  370. </div>
  371. <a name="N10089"></a><a name="%E6%96%87%E4%BB%B6%E7%B3%BB%E7%BB%9F%E7%9A%84%E5%90%8D%E5%AD%97%E7%A9%BA%E9%97%B4+%28namespace%29"></a>
  372. <h2 class="h3"> 文件系统的名字空间 (namespace) </h2>
  373. <div class="section">
  374. <p>
  375. HDFS支持传统的层次型文件组织结构。用户或者应用程序可以创建目录,然后将文件保存在这些目录里。文件系统名字空间的层次结构和大多数现有的文件系统类似:用户可以创建、删除、移动或重命名文件。当前,HDFS不支持用户磁盘配额和访问权限控制,也不支持硬链接和软链接。但是HDFS架构并不妨碍实现这些特性。
  376. </p>
  377. <p>
  378. Namenode负责维护文件系统的名字空间,任何对文件系统名字空间或属性的修改都将被Namenode记录下来。应用程序可以设置HDFS保存的文件的副本数目。文件副本的数目称为文件的副本系数,这个信息也是由Namenode保存的。
  379. </p>
  380. </div>
  381. <a name="N10096"></a><a name="%E6%95%B0%E6%8D%AE%E5%A4%8D%E5%88%B6"></a>
  382. <h2 class="h3"> 数据复制 </h2>
  383. <div class="section">
  384. <p>
  385. HDFS被设计成能够在一个大集群中跨机器可靠地存储超大文件。它将每个文件存储成一系列的数据块,除了最后一个,所有的数据块都是同样大小的。为了容错,文件的所有数据块都会有副本。每个文件的数据块大小和副本系数都是可配置的。应用程序可以指定某个文件的副本数目。副本系数可以在文件创建的时候指定,也可以在之后改变。HDFS中的文件都是一次性写入的,并且严格要求在任何时候只能有一个写入者。
  386. </p>
  387. <p>
  388. Namenode全权管理数据块的复制,它周期性地从集群中的每个Datanode接收心跳信号和块状态报告(Blockreport)。接收到心跳信号意味着该Datanode节点工作正常。块状态报告包含了一个该Datanode上所有数据块的列表。
  389. </p>
  390. <div id="" style="text-align: center;">
  391. <img id="" class="figure" alt="HDFS Datanodes" src="images/hdfsdatanodes.gif"></div>
  392. <a name="N100A6"></a><a name="%E5%89%AF%E6%9C%AC%E5%AD%98%E6%94%BE%3A+%E6%9C%80%E6%9C%80%E5%BC%80%E5%A7%8B%E7%9A%84%E4%B8%80%E6%AD%A5"></a>
  393. <h3 class="h4"> 副本存放: 最最开始的一步 </h3>
  394. <p>
  395. 副本的存放是HDFS可靠性和性能的关键。优化的副本存放策略是HDFS区分于其他大部分分布式文件系统的重要特性。这种特性需要做大量的调优,并需要经验的积累。HDFS采用一种称为机架感知(rack-aware)的策略来改进数据的可靠性、可用性和网络带宽的利用率。目前实现的副本存放策略只是在这个方向上的第一步。实现这个策略的短期目标是验证它在生产环境下的有效性,观察它的行为,为实现更先进的策略打下测试和研究的基础。
  396. </p>
  397. <p>
  398. 大型HDFS实例一般运行在跨越多个机架的计算机组成的集群上,不同机架上的两台机器之间的通讯需要经过交换机。在大多数情况下,同一个机架内的两台机器间的带宽会比不同机架的两台机器间的带宽大。
  399. </p>
  400. <p>
  401. 通过一个<a href="cluster_setup.html#Hadoop%E7%9A%84%E6%9C%BA%E6%9E%B6%E6%84%9F%E7%9F%A5">机架感知</a>的过程,Namenode可以确定每个Datanode所属的机架id。一个简单但没有优化的策略就是将副本存放在不同的机架上。这样可以有效防止当整个机架失效时数据的丢失,并且允许读数据的时候充分利用多个机架的带宽。这种策略设置可以将副本均匀分布在集群中,有利于当组件失效情况下的负载均衡。但是,因为这种策略的一个写操作需要传输数据块到多个机架,这增加了写的代价。
  402. </p>
  403. <p>
  404. 在大多数情况下,副本系数是3,HDFS的存放策略是将一个副本存放在本地机架的节点上,一个副本放在同一机架的另一个节点上,最后一个副本放在不同机架的节点上。这种策略减少了机架间的数据传输,这就提高了写操作的效率。机架的错误远远比节点的错误少,所以这个策略不会影响到数据的可靠性和可用性。于此同时,因为数据块只放在两个(不是三个)不同的机架上,所以此策略减少了读取数据时需要的网络传输总带宽。在这种策略下,副本并不是均匀分布在不同的机架上。三分之一的副本在一个节点上,三分之二的副本在一个机架上,其他副本均匀分布在剩下的机架中,这一策略在不损害数据可靠性和读取性能的情况下改进了写的性能。
  405. </p>
  406. <p>
  407. 当前,这里介绍的默认副本存放策略正在开发的过程中。
  408. </p>
  409. <a name="N100C0"></a><a name="%E5%89%AF%E6%9C%AC%E9%80%89%E6%8B%A9"></a>
  410. <h3 class="h4"> 副本选择 </h3>
  411. <p>
  412. 为了降低整体的带宽消耗和读取延时,HDFS会尽量让读取程序读取离它最近的副本。如果在读取程序的同一个机架上有一个副本,那么就读取该副本。如果一个HDFS集群跨越多个数据中心,那么客户端也将首先读本地数据中心的副本。
  413. </p>
  414. <a name="N100CA"></a><a name="%E5%AE%89%E5%85%A8%E6%A8%A1%E5%BC%8F"></a>
  415. <h3 class="h4"> 安全模式 </h3>
  416. <p>
  417. Namenode启动后会进入一个称为安全模式的特殊状态。处于安全模式的Namenode是不会进行数据块的复制的。Namenode从所有的 Datanode接收心跳信号和块状态报告。块状态报告包括了某个Datanode所有的数据块列表。每个数据块都有一个指定的最小副本数。当Namenode检测确认某个数据块的副本数目达到这个最小值,那么该数据块就会被认为是副本安全(safely replicated)的;在一定百分比(这个参数可配置)的数据块被Namenode检测确认是安全之后(加上一个额外的30秒等待时间),Namenode将退出安全模式状态。接下来它会确定还有哪些数据块的副本没有达到指定数目,并将这些数据块复制到其他Datanode上。
  418. </p>
  419. </div>
  420. <a name="N100D5"></a><a name="%E6%96%87%E4%BB%B6%E7%B3%BB%E7%BB%9F%E5%85%83%E6%95%B0%E6%8D%AE%E7%9A%84%E6%8C%81%E4%B9%85%E5%8C%96"></a>
  421. <h2 class="h3"> 文件系统元数据的持久化 </h2>
  422. <div class="section">
  423. <p>
  424. Namenode上保存着HDFS的名字空间。对于任何对文件系统元数据产生修改的操作,Namenode都会使用一种称为EditLog的事务日志记录下来。例如,在HDFS中创建一个文件,Namenode就会在Editlog中插入一条记录来表示;同样地,修改文件的副本系数也将往Editlog插入一条记录。Namenode在本地操作系统的文件系统中存储这个Editlog。整个文件系统的名字空间,包括数据块到文件的映射、文件的属性等,都存储在一个称为FsImage的文件中,这个文件也是放在Namenode所在的本地文件系统上。
  425. </p>
  426. <p>
  427. Namenode在内存中保存着整个文件系统的名字空间和文件数据块映射(Blockmap)的映像。这个关键的元数据结构设计得很紧凑,因而一个有4G内存的Namenode足够支撑大量的文件和目录。当Namenode启动时,它从硬盘中读取Editlog和FsImage,将所有Editlog中的事务作用在内存中的FsImage上,并将这个新版本的FsImage从内存中保存到本地磁盘上,然后删除旧的Editlog,因为这个旧的Editlog的事务都已经作用在FsImage上了。这个过程称为一个检查点(checkpoint)。在当前实现中,检查点只发生在Namenode启动时,在不久的将来将实现支持周期性的检查点。
  428. </p>
  429. <p>
  430. Datanode将HDFS数据以文件的形式存储在本地的文件系统中,它并不知道有关HDFS文件的信息。它把每个HDFS数据块存储在本地文件系统的一个单独的文件中。Datanode并不在同一个目录创建所有的文件,实际上,它用试探的方法来确定每个目录的最佳文件数目,并且在适当的时候创建子目录。在同一个目录中创建所有的本地文件并不是最优的选择,这是因为本地文件系统可能无法高效地在单个目录中支持大量的文件。当一个Datanode启动时,它会扫描本地文件系统,产生一个这些本地文件对应的所有HDFS数据块的列表,然后作为报告发送到Namenode,这个报告就是块状态报告。
  431. </p>
  432. </div>
  433. <a name="N100E5"></a><a name="%E9%80%9A%E8%AE%AF%E5%8D%8F%E8%AE%AE"></a>
  434. <h2 class="h3"> 通讯协议 </h2>
  435. <div class="section">
  436. <p>
  437. 所有的HDFS通讯协议都是建立在TCP/IP协议之上。客户端通过一个可配置的<acronym title="Transmission Control Protocol">TCP</acronym>端口连接到Namenode,通过ClientProtocol协议与Namenode交互。而Datanode使用DatanodeProtocol协议与Namenode交互。一个远程过程调用(<acronym title="Remote Procedure Call">RPC</acronym>)模型被抽象出来封装ClientProtocol和Datanodeprotocol协议。在设计上,Namenode不会主动发起RPC,而是响应来自客户端或 Datanode 的RPC请求。
  438. </p>
  439. </div>
  440. <a name="N100F7"></a><a name="%E5%81%A5%E5%A3%AE%E6%80%A7"></a>
  441. <h2 class="h3"> 健壮性 </h2>
  442. <div class="section">
  443. <p>
  444. HDFS的主要目标就是即使在出错的情况下也要保证数据存储的可靠性。常见的三种出错情况是:Namenode出错, Datanode出错和网络割裂(network partitions)。
  445. </p>
  446. <a name="N10100"></a><a name="%E7%A3%81%E7%9B%98%E6%95%B0%E6%8D%AE%E9%94%99%E8%AF%AF%EF%BC%8C%E5%BF%83%E8%B7%B3%E6%A3%80%E6%B5%8B%E5%92%8C%E9%87%8D%E6%96%B0%E5%A4%8D%E5%88%B6"></a>
  447. <h3 class="h4"> 磁盘数据错误,心跳检测和重新复制 </h3>
  448. <p>
  449. 每个Datanode节点周期性地向Namenode发送心跳信号。网络割裂可能导致一部分Datanode跟Namenode失去联系。Namenode通过心跳信号的缺失来检测这一情况,并将这些近期不再发送心跳信号Datanode标记为宕机,不会再将新的<acronym title="Input/Output">IO</acronym>请求发给它们。任何存储在宕机Datanode上的数据将不再有效。Datanode的宕机可能会引起一些数据块的副本系数低于指定值,Namenode不断地检测这些需要复制的数据块,一旦发现就启动复制操作。在下列情况下,可能需要重新复制:某个Datanode节点失效,某个副本遭到损坏,Datanode上的硬盘错误,或者文件的副本系数增大。
  450. </p>
  451. <a name="N1010E"></a><a name="%E9%9B%86%E7%BE%A4%E5%9D%87%E8%A1%A1"></a>
  452. <h3 class="h4"> 集群均衡 </h3>
  453. <p>
  454. HDFS的架构支持数据均衡策略。如果某个Datanode节点上的空闲空间低于特定的临界点,按照均衡策略系统就会自动地将数据从这个Datanode移动到其他空闲的Datanode。当对某个文件的请求突然增加,那么也可能启动一个计划创建该文件新的副本,并且同时重新平衡集群中的其他数据。这些均衡策略目前还没有实现。
  455. </p>
  456. <a name="N10118"></a><a name="%E6%95%B0%E6%8D%AE%E5%AE%8C%E6%95%B4%E6%80%A7"></a>
  457. <h3 class="h4"> 数据完整性 </h3>
  458. <p>
  459. <!-- XXX "checksum checking" sounds funny -->
  460. 从某个Datanode获取的数据块有可能是损坏的,损坏可能是由Datanode的存储设备错误、网络错误或者软件bug造成的。HDFS客户端软件实现了对HDFS文件内容的校验和(checksum)检查。当客户端创建一个新的HDFS文件,会计算这个文件每个数据块的校验和,并将校验和作为一个单独的隐藏文件保存在同一个HDFS名字空间下。当客户端获取文件内容后,它会检验从Datanode获取的数据跟相应的校验和文件中的校验和是否匹配,如果不匹配,客户端可以选择从其他Datanode获取该数据块的副本。
  461. </p>
  462. <a name="N10124"></a><a name="%E5%85%83%E6%95%B0%E6%8D%AE%E7%A3%81%E7%9B%98%E9%94%99%E8%AF%AF"></a>
  463. <h3 class="h4"> 元数据磁盘错误 </h3>
  464. <p>
  465. FsImage和Editlog是HDFS的核心数据结构。如果这些文件损坏了,整个HDFS实例都将失效。因而,Namenode可以配置成支持维护多个FsImage和Editlog的副本。任何对FsImage或者Editlog的修改,都将同步到它们的副本上。这种多副本的同步操作可能会降低Namenode每秒处理的名字空间事务数量。然而这个代价是可以接受的,因为即使HDFS的应用是数据密集的,它们也非元数据密集的。当Namenode重启的时候,它会选取最近的完整的FsImage和Editlog来使用。
  466. </p>
  467. <p>
  468. Namenode是HDFS集群中的单点故障(single point of failure)所在。如果Namenode机器故障,是需要手工干预的。目前,自动重启或在另一台机器上做Namenode故障转移的功能还没实现。
  469. </p>
  470. <a name="N10131"></a><a name="%E5%BF%AB%E7%85%A7"></a>
  471. <h3 class="h4"> 快照 </h3>
  472. <p>
  473. 快照支持某一特定时刻的数据的复制备份。利用快照,可以让HDFS在数据损坏时恢复到过去一个已知正确的时间点。HDFS目前还不支持快照功能,但计划在将来的版本进行支持。
  474. </p>
  475. </div>
  476. <a name="N1013C"></a><a name="%E6%95%B0%E6%8D%AE%E7%BB%84%E7%BB%87"></a>
  477. <h2 class="h3"> 数据组织 </h2>
  478. <div class="section">
  479. <a name="N10144"></a><a name="%E6%95%B0%E6%8D%AE%E5%9D%97"></a>
  480. <h3 class="h4"> 数据块 </h3>
  481. <p>
  482. HDFS被设计成支持大文件,适用HDFS的是那些需要处理大规模的数据集的应用。这些应用都是只写入数据一次,但却读取一次或多次,并且读取速度应能满足流式读取的需要。HDFS支持文件的&ldquo;一次写入多次读取&rdquo;语义。一个典型的数据块大小是64MB。因而,HDFS中的文件总是按照64M被切分成不同的块,每个块尽可能地存储于不同的Datanode中。
  483. </p>
  484. <a name="N1014E"></a><a name="Staging"></a>
  485. <h3 class="h4"> Staging </h3>
  486. <p>
  487. 客户端创建文件的请求其实并没有立即发送给Namenode,事实上,在刚开始阶段HDFS客户端会先将文件数据缓存到本地的一个临时文件。应用程序的写操作被透明地重定向到这个临时文件。当这个临时文件累积的数据量超过一个数据块的大小,客户端才会联系Namenode。Namenode将文件名插入文件系统的层次结构中,并且分配一个数据块给它。然后返回Datanode的标识符和目标数据块给客户端。接着客户端将这块数据从本地临时文件上传到指定的Datanode上。当文件关闭时,在临时文件中剩余的没有上传的数据也会传输到指定的Datanode上。然后客户端告诉Namenode文件已经关闭。此时Namenode才将文件创建操作提交到日志里进行存储。如果Namenode在文件关闭前宕机了,则该文件将丢失。
  488. </p>
  489. <p>
  490. 上述方法是对在HDFS上运行的目标应用进行认真考虑后得到的结果。这些应用需要进行文件的流式写入。如果不采用客户端缓存,由于网络速度和网络堵塞会对吞估量造成比较大的影响。这种方法并不是没有先例的,早期的文件系统,比如<acronym title="Andrew File System">AFS</acronym>,就用客户端缓存来提高性能。为了达到更高的数据上传效率,已经放松了POSIX标准的要求。
  491. </p>
  492. <a name="N10161"></a><a name="%E6%B5%81%E6%B0%B4%E7%BA%BF%E5%A4%8D%E5%88%B6"></a>
  493. <h3 class="h4"> 流水线复制 </h3>
  494. <p>
  495. 当客户端向HDFS文件写入数据的时候,一开始是写到本地临时文件中。假设该文件的副本系数设置为3,当本地临时文件累积到一个数据块的大小时,客户端会从Namenode获取一个Datanode列表用于存放副本。然后客户端开始向第一个Datanode传输数据,第一个Datanode一小部分一小部分(4 KB)地接收数据,将每一部分写入本地仓库,并同时传输该部分到列表中第二个Datanode节点。第二个Datanode也是这样,一小部分一小部分地接收数据,写入本地仓库,并同时传给第三个Datanode。最后,第三个Datanode接收数据并存储在本地。因此,Datanode能流水线式地从前一个节点接收数据,并在同时转发给下一个节点,数据以流水线的方式从前一个Datanode复制到下一个。
  496. </p>
  497. </div>
  498. <a name="N1016C"></a><a name="%E5%8F%AF%E8%AE%BF%E9%97%AE%E6%80%A7"></a>
  499. <h2 class="h3"> 可访问性 </h2>
  500. <div class="section">
  501. <p>
  502. HDFS给应用提供了多种访问方式。用户可以通过<a href="http://hadoop.apache.org/core/docs/current/api/">Java API</a>接口访问,也可以通过C语言的封装API访问,还可以通过浏览器的方式访问HDFS中的文件。通过<acronym title="Web-based Distributed Authoring and Versioning">WebDAV</acronym>协议访问的方式正在开发中。
  503. </p>
  504. <a name="N10181"></a><a name="DFSShell"></a>
  505. <h3 class="h4"> DFSShell </h3>
  506. <p>
  507. HDFS以文件和目录的形式组织用户数据。它提供了一个命令行的接口(DFSShell)让用户与HDFS中的数据进行交互。命令的语法和用户熟悉的其他shell(例如 bash, csh)工具类似。下面是一些动作/命令的示例:
  508. </p>
  509. <table class="ForrestTable" cellspacing="1" cellpadding="4">
  510. <tr>
  511. <th colspan="1" rowspan="1"> 动作 </th><th colspan="1" rowspan="1"> 命令 </th>
  512. </tr>
  513. <tr>
  514. <td colspan="1" rowspan="1"> 创建一个名为 <span class="codefrag">/foodir</span> 的目录 </td> <td colspan="1" rowspan="1"> <span class="codefrag">bin/hadoop dfs -mkdir /foodir</span> </td>
  515. </tr>
  516. <tr>
  517. <td colspan="1" rowspan="1"> 创建一个名为 <span class="codefrag">/foodir</span> 的目录 </td> <td colspan="1" rowspan="1"> <span class="codefrag">bin/hadoop dfs -mkdir /foodir</span> </td>
  518. </tr>
  519. <tr>
  520. <td colspan="1" rowspan="1"> 查看名为 <span class="codefrag">/foodir/myfile.txt</span> 的文件内容 </td> <td colspan="1" rowspan="1"> <span class="codefrag">bin/hadoop dfs -cat /foodir/myfile.txt</span> </td>
  521. </tr>
  522. </table>
  523. <p>
  524. DFSShell 可以用在那些通过脚本语言和文件系统进行交互的应用程序上。
  525. </p>
  526. <a name="N101D6"></a><a name="DFSAdmin"></a>
  527. <h3 class="h4"> DFSAdmin </h3>
  528. <p>
  529. DFSAdmin 命令用来管理HDFS集群。这些命令只有HDSF的管理员才能使用。下面是一些动作/命令的示例:
  530. </p>
  531. <table class="ForrestTable" cellspacing="1" cellpadding="4">
  532. <tr>
  533. <th colspan="1" rowspan="1"> 动作 </th><th colspan="1" rowspan="1"> 命令 </th>
  534. </tr>
  535. <tr>
  536. <td colspan="1" rowspan="1"> 将集群置于安全模式 </td> <td colspan="1" rowspan="1"> <span class="codefrag">bin/hadoop dfsadmin -safemode enter</span> </td>
  537. </tr>
  538. <tr>
  539. <td colspan="1" rowspan="1"> 显示Datanode列表 </td> <td colspan="1" rowspan="1"> <span class="codefrag">bin/hadoop dfsadmin -report</span> </td>
  540. </tr>
  541. <tr>
  542. <td colspan="1" rowspan="1"> 使Datanode节点 <span class="codefrag">datanodename</span>退役</td><td colspan="1" rowspan="1"> <span class="codefrag">bin/hadoop dfsadmin -decommission datanodename</span> </td>
  543. </tr>
  544. </table>
  545. <a name="N10221"></a><a name="%E6%B5%8F%E8%A7%88%E5%99%A8%E6%8E%A5%E5%8F%A3"></a>
  546. <h3 class="h4"> 浏览器接口 </h3>
  547. <p>
  548. 一个典型的HDFS安装会在一个可配置的TCP端口开启一个Web服务器用于暴露HDFS的名字空间。用户可以用浏览器来浏览HDFS的名字空间和查看文件的内容。
  549. </p>
  550. </div>
  551. <a name="N1022C"></a><a name="%E5%AD%98%E5%82%A8%E7%A9%BA%E9%97%B4%E5%9B%9E%E6%94%B6"></a>
  552. <h2 class="h3"> 存储空间回收 </h2>
  553. <div class="section">
  554. <a name="N10232"></a><a name="%E6%96%87%E4%BB%B6%E7%9A%84%E5%88%A0%E9%99%A4%E5%92%8C%E6%81%A2%E5%A4%8D"></a>
  555. <h3 class="h4"> 文件的删除和恢复 </h3>
  556. <p>
  557. 当用户或应用程序删除某个文件时,这个文件并没有立刻从HDFS中删除。实际上,HDFS会将这个文件重命名转移到<span class="codefrag">/trash</span>目录。只要文件还在<span class="codefrag">/trash</span>目录中,该文件就可以被迅速地恢复。文件在<span class="codefrag">/trash</span>中保存的时间是可配置的,当超过这个时间时,Namenode就会将该文件从名字空间中删除。删除文件会使得该文件相关的数据块被释放。注意,从用户删除文件到HDFS空闲空间的增加之间会有一定时间的延迟。</p>
  558. <p>
  559. 只要被删除的文件还在<span class="codefrag">/trash</span>目录中,用户就可以恢复这个文件。如果用户想恢复被删除的文件,他/她可以浏览<span class="codefrag">/trash</span>目录找回该文件。<span class="codefrag">/trash</span>目录仅仅保存被删除文件的最后副本。<span class="codefrag">/trash</span>目录与其他的目录没有什么区别,除了一点:在该目录上HDFS会应用一个特殊策略来自动删除文件。目前的默认策略是删除<span class="codefrag">/trash</span>中保留时间超过6小时的文件。将来,这个策略可以通过一个被良好定义的接口配置。
  560. </p>
  561. <a name="N10257"></a><a name="%E5%87%8F%E5%B0%91%E5%89%AF%E6%9C%AC%E7%B3%BB%E6%95%B0"></a>
  562. <h3 class="h4"> 减少副本系数 </h3>
  563. <p>
  564. 当一个文件的副本系数被减小后,Namenode会选择过剩的副本删除。下次心跳检测时会将该信息传递给Datanode。Datanode遂即移除相应的数据块,集群中的空闲空间加大。同样,在调用<span class="codefrag">setReplication</span> API结束和集群中空闲空间增加间会有一定的延迟。</p>
  565. </div>
  566. <a name="N10265"></a><a name="%E5%8F%82%E8%80%83%E8%B5%84%E6%96%99"></a>
  567. <h2 class="h3"> 参考资料 </h2>
  568. <div class="section">
  569. <p>
  570. HDFS Java API:
  571. <a href="http://hadoop.apache.org/core/docs/current/api/">
  572. http://hadoop.apache.org/core/docs/current/api/
  573. </a>
  574. </p>
  575. <p>
  576. HDFS 源代码:
  577. <a href="http://hadoop.apache.org/core/version_control.html">
  578. http://hadoop.apache.org/core/version_control.html
  579. </a>
  580. </p>
  581. </div>
  582. <p align="right">
  583. <font size="-2">by&nbsp;Dhruba Borthakur</font>
  584. </p>
  585. </div>
  586. <!--+
  587. |end content
  588. +-->
  589. <div class="clearboth">&nbsp;</div>
  590. </div>
  591. <div id="footer">
  592. <!--+
  593. |start bottomstrip
  594. +-->
  595. <div class="lastmodified">
  596. <script type="text/javascript"><!--
  597. document.write("Last Published: " + document.lastModified);
  598. // --></script>
  599. </div>
  600. <div class="copyright">
  601. Copyright &copy;
  602. 2007 <a href="http://www.apache.org/licenses/">The Apache Software Foundation.</a>
  603. </div>
  604. <!--+
  605. |end bottomstrip
  606. +-->
  607. </div>
  608. </body>
  609. </html>