1
00:00:00,810 --> 00:00:02,890
Welcome to Mining of Massive Datasets.

2
00:00:02,890 --> 00:00:06,480
I'm Anand Rajaraman and
today's topic is Map-Reduce.

3
00:00:06,480 --> 00:00:10,100
In the last few years Map-Reduce has
emerged as a leading paradigm for

4
00:00:10,100 --> 00:00:12,250
mining really massive data sets.

5
00:00:12,250 --> 00:00:15,690
But before we get into Map-Reduce proper,
let's spend a few minutes trying to

6
00:00:15,690 --> 00:00:17,580
understand why we need
Map-Reduce in the first place.

7
00:00:17,580 --> 00:00:19,450
Let's start with the basics.

8
00:00:21,460 --> 00:00:24,950
Now we're all familiar with the basic
computational model of CPU and

9
00:00:24,950 --> 00:00:26,350
memory, right?

10
00:00:26,350 --> 00:00:31,190
The algorithm runs on the CPU, and
accesses data that's in memory.

11
00:00:31,190 --> 00:00:34,784
Now we may need to bring the data
in from disk into memory, but

12
00:00:34,784 --> 00:00:37,691
once the data is in memory,
fits in there fully.

13
00:00:37,691 --> 00:00:39,911
So you don't need to access disk again,
and

14
00:00:39,911 --> 00:00:43,240
the algorithm just runs in
the data that's on memory.

15
00:00:43,240 --> 00:00:47,273
Now there's a familiar model that we use
to implement all kinds of algorithms, and

16
00:00:47,273 --> 00:00:49,084
machined learning, and statistics.

17
00:00:49,084 --> 00:00:51,850
And pretty much everything else.

18
00:00:51,850 --> 00:00:52,520
All right?

19
00:00:52,520 --> 00:00:54,630
Now, what happened to the data is so

20
00:00:54,630 --> 00:00:57,820
big, that it can't all fit
in memory at the same time.

21
00:00:57,820 --> 00:00:59,460
That's where data mining comes in.

22
00:00:59,460 --> 00:01:02,410
And classical data mining algorithms.

23
00:01:02,410 --> 00:01:05,230
Look at the disk in addition
to looking at CPU and memory.

24
00:01:05,230 --> 00:01:06,670
So the data's on disk,

25
00:01:06,670 --> 00:01:09,440
you can only bring in a portion of
the data into memory at a time.

26
00:01:10,510 --> 00:01:14,930
And you can process it in batches, and
you know, write back results to disk.

27
00:01:14,930 --> 00:01:17,700
And this is the realm of
classical data mining algorithms.

28
00:01:17,700 --> 00:01:20,190
But sometimes even this is not sufficient.

29
00:01:20,190 --> 00:01:21,009
Let's look at an example.

30
00:01:23,210 --> 00:01:28,260
So think about Google, crawling and
indexing the web, right?

31
00:01:28,260 --> 00:01:32,510
Let's say,
google has crawled 10 billion web pages.

32
00:01:32,510 --> 00:01:37,860
And let's further say, that the average
size of a web page is 20 KB.

33
00:01:37,860 --> 00:01:41,450
Now, these are representative
numbers from real life.

34
00:01:41,450 --> 00:01:44,780
Now if you take ten billion webpages,
each of 20 KB,

35
00:01:44,780 --> 00:01:49,480
you have, total data set size of 200 TB.

36
00:01:49,480 --> 00:01:52,300
Now, when you have 200 TB,
let's assume that they're using

37
00:01:52,300 --> 00:01:55,260
the classical computational model,
classical data mining model.

38
00:01:55,260 --> 00:01:57,740
And all this data is stored
on a single disk, and

39
00:01:57,740 --> 00:01:59,860
we have read tend to be
processed inside a CPU.

40
00:02:01,240 --> 00:02:05,330
Now the fundamental limitation
here is the bandwidth,

41
00:02:05,330 --> 00:02:08,030
the data bandwidth between the disk and
the CPU.

42
00:02:08,030 --> 00:02:12,450
The data has to be read from
the disk into the CPU, and

43
00:02:12,450 --> 00:02:17,070
the disk read bandwidth for most modern
SATA disk representative number.

44
00:02:17,070 --> 00:02:19,430
Is around 50MB a second.

45
00:02:19,430 --> 00:02:22,130
So, so we can read data at 50MB a second.

46
00:02:22,130 --> 00:02:25,680
How long does it take to
read 200TB at 50MB a second?

47
00:02:25,680 --> 00:02:27,050
Can do some simple math, and

48
00:02:27,050 --> 00:02:31,020
the answer is 4 million seconds
which is more than 46 days.

49
00:02:31,020 --> 00:02:32,890
Remember, this is an awfully long time,
and

50
00:02:32,890 --> 00:02:35,940
is just the time to read
the data into memory.

51
00:02:35,940 --> 00:02:39,830
To do something useful with the data,
it's going to take even longer.

52
00:02:39,830 --> 00:02:41,870
Right, so clearly this is unacceptable.

53
00:02:41,870 --> 00:02:44,330
You can't take four to six
days just to read the data.

54
00:02:44,330 --> 00:02:45,640
So you need a better solution.

55
00:02:45,640 --> 00:02:50,610
Now the obvious thing that you think of is
that it can split the data into chunks.

56
00:02:50,610 --> 00:02:53,310
And you can have multiple disks and CPUs.

57
00:02:53,310 --> 00:02:56,090
you, you stripe the data
across multiple disks.

58
00:02:56,090 --> 00:03:00,150
And you can read it, and, and
process it in parallel in multiple CPUs.

59
00:03:00,150 --> 00:03:02,460
That will cut down, this time by a lot.

60
00:03:02,460 --> 00:03:06,668
For example, if you had a 1,000 disks and
CPUs, in four thousa-,

61
00:03:06,668 --> 00:03:08,080
4 million seconds.

62
00:03:08,080 --> 00:03:12,240
And we were completely in parallel, in 4
million seconds, you could do the job in,

63
00:03:13,600 --> 00:03:17,940
4 million by 1,000,
which is 4,000 seconds.

64
00:03:17,940 --> 00:03:22,610
And that's just about an hour which is,
which is very acceptable time.

65
00:03:22,610 --> 00:03:23,490
Right?
So

66
00:03:23,490 --> 00:03:27,250
this is the fundamental idea behind
the idea of cluster computing.

67
00:03:27,250 --> 00:03:28,100
Right?
And this is,

68
00:03:28,100 --> 00:03:30,280
this tiered architecture
that has emerged for

69
00:03:30,280 --> 00:03:32,140
cluster computing is something like this.

70
00:03:32,140 --> 00:03:36,087
You have the racks consisting
of commodity Linux nodes.

71
00:03:36,087 --> 00:03:39,610
As you go with commodity Linux
nodes because they are very cheap.

72
00:03:39,610 --> 00:03:44,428
And you can, you can buy thousands and
thousands of them and, and rack them up.

73
00:03:44,428 --> 00:03:47,170
you, you have many of these racks.

74
00:03:47,170 --> 00:03:53,760
Each rack has 16 to 64 of these
commodity Linux nodes and

75
00:03:53,760 --> 00:03:56,602
these nodes are connected by a switch.

76
00:03:56,602 --> 00:03:59,851
and, the, the, the switch in a rack
is typically a gigabit switch.

77
00:03:59,851 --> 00:04:05,120
So there's 1 Gbps bandwidth
between any pair of nodes in rack.

78
00:04:06,160 --> 00:04:08,950
Of course 16 to 64 nodes
is not sufficient.

79
00:04:08,950 --> 00:04:12,460
So you have multiple racks, and all the,

80
00:04:12,460 --> 00:04:15,390
the racks themselves are connected
by backbone switches.

81
00:04:15,390 --> 00:04:16,700
And the backbones is,

82
00:04:16,700 --> 00:04:22,520
is a higher bandwidth switch can do
two to ten gigabits between racks.

83
00:04:22,520 --> 00:04:26,000
Right?
So so we have 16 to 64 nodes in a rack.

84
00:04:26,000 --> 00:04:30,580
And then you, you rack up multiple racks,
and, and you get a data center.

85
00:04:30,580 --> 00:04:34,140
So this is the standard classical
architecture that has emerged over

86
00:04:34,140 --> 00:04:35,638
the last few years.

87
00:04:35,638 --> 00:04:40,832
For you know, for storing and
mining very large data sets.

88
00:04:40,832 --> 00:04:44,350
Now once you have this kind of cluster
this doesn't solve the problem completely.

89
00:04:44,350 --> 00:04:47,230
Because cluster computing comes
with it's own challenges.

90
00:04:49,422 --> 00:04:53,963
But before we get there, let's get us,
you know, ideal of the scale, right?

91
00:04:53,963 --> 00:04:58,253
In 2011 somebody estimated that
Google had a million machines,

92
00:04:58,253 --> 00:05:00,490
million nodes like this.

93
00:05:00,490 --> 00:05:03,670
In stacked up you know,
is, is somewhat like this.

94
00:05:03,670 --> 00:05:08,120
So, so it gives, so that gives you a sense
of the scale of modern data centers and,

95
00:05:08,120 --> 00:05:09,790
and, and clusters, right?

96
00:05:09,790 --> 00:05:11,500
So here's, here's a picture.

97
00:05:11,500 --> 00:05:15,240
This is what,
it looks like inside a data center.

98
00:05:15,240 --> 00:05:18,610
So the, the, what you see there is,
is the back up racks, and

99
00:05:18,610 --> 00:05:21,090
you can see the connections,
between, between the racks.

100
00:05:22,460 --> 00:05:26,630
Now, once you have such a big cluster,

101
00:05:26,630 --> 00:05:29,320
you actually have to do
computations on the cluster.

102
00:05:29,320 --> 00:05:29,940
Right?

103
00:05:29,940 --> 00:05:33,390
And clustered computing comes
with its own, challenges.

104
00:05:33,390 --> 00:05:37,420
The first and the most major
challenge is that nodes can fail.

105
00:05:37,420 --> 00:05:37,940
Right?

106
00:05:37,940 --> 00:05:41,140
Now a single,
node doesn't fail that often.

107
00:05:41,140 --> 00:05:41,940
Right?
If you,

108
00:05:41,940 --> 00:05:43,980
if you just connect, the next node and

109
00:05:43,980 --> 00:05:48,000
let it stay up, it can probably stay
up for, three years without failing.

110
00:05:48,000 --> 00:05:49,940
Three years is about a 1,000 days.

111
00:05:49,940 --> 00:05:54,160
So that's, you know, once in a 1,000
days failure isn't such a big deal.

112
00:05:54,160 --> 00:05:57,160
But now imagine that you have
a 1,000 servers in a cluster.

113
00:05:57,160 --> 00:06:02,970
And in your, and if you assume that these,
servers fail, independent of each other.

114
00:06:02,970 --> 00:06:05,500
You're going to get
approximately one failure a day.

115
00:06:05,500 --> 00:06:06,840
Which is, still isn't such a big deal.

116
00:06:06,840 --> 00:06:08,360
You can probably deal with it.

117
00:06:08,360 --> 00:06:11,710
But now imagine something on the scale
of Google which has a million servers,

118
00:06:11,710 --> 00:06:12,710
in its cluster.

119
00:06:12,710 --> 00:06:15,990
So if you have a million servers, you're
going to get a 1,000 failures per day.

120
00:06:15,990 --> 00:06:18,127
Now a 1,000 failures per day is a lot and

121
00:06:18,127 --> 00:06:21,380
you need some kind of infrastructure
to deal with that kind of failure rate.

122
00:06:22,430 --> 00:06:25,950
Your failures on that scale
introduce two kinds of problems.

123
00:06:25,950 --> 00:06:28,610
The first problem is that if, you know,
if nodes are going to fail and

124
00:06:28,610 --> 00:06:30,840
you're going to store
your data on these nodes.

125
00:06:30,840 --> 00:06:34,060
How do you keep the data and
store persistently?

126
00:06:34,060 --> 00:06:35,120
What does this mean?

127
00:06:35,120 --> 00:06:37,430
Persistence means that
once you store the data,

128
00:06:37,430 --> 00:06:39,430
you're guaranteed you can read it again.

129
00:06:39,430 --> 00:06:43,460
But if the node in which you stored the
data fails, then you can't read the data.

130
00:06:43,460 --> 00:06:44,640
You might even lose the data.

131
00:06:44,640 --> 00:06:48,050
So how do you keep the data
stored persistently if like,

132
00:06:48,050 --> 00:06:49,350
these nodes can fail.

133
00:06:49,350 --> 00:06:53,450
Now the second problem is
is is one of availability.

134
00:06:53,450 --> 00:06:58,520
So, let's say you're running one of the
computations, and this computation is, a,

135
00:06:58,520 --> 00:07:01,850
you know,
analyzing massive amounts of data.

136
00:07:01,850 --> 00:07:03,510
And it's chugging through
the computation and

137
00:07:03,510 --> 00:07:06,180
it's going, you know,
run half way through the computation.

138
00:07:06,180 --> 00:07:10,690
And, you know, at this critical point,
a couple of nodes fail, right?

139
00:07:10,690 --> 00:07:14,050
And that node had data that is
necessary for the computation.

140
00:07:14,050 --> 00:07:15,720
Now how we deal with this problem.

141
00:07:15,720 --> 00:07:17,590
Now in the first place you
may have to go back and

142
00:07:17,590 --> 00:07:19,590
restart the computation all over again.

143
00:07:19,590 --> 00:07:22,060
But if you restart it now and, and, and

144
00:07:22,060 --> 00:07:25,860
the computation turns again when
the computation is running.

145
00:07:25,860 --> 00:07:30,690
So kind of need an infrastructure that
can hide these kinds of node failures and

146
00:07:30,690 --> 00:07:34,410
let the computation go to go to
completion even if nodes fail.

147
00:07:35,990 --> 00:07:39,800
The second challenge of
cluster computing is that

148
00:07:39,800 --> 00:07:42,450
the network itself can
become a bottleneck.

149
00:07:42,450 --> 00:07:45,981
Now remember,
there is this 1 Gbps network bandwidth.

150
00:07:45,981 --> 00:07:49,253
That is available between
individual nodes in a rack and

151
00:07:49,253 --> 00:07:52,923
a smaller bandwidth that's
available between individual racks.

152
00:07:52,923 --> 00:07:55,953
Though if you have 10 TB of data,
and you have to move it

153
00:07:55,953 --> 00:08:00,100
across a 1 Gbps network connection,
that takes approximately a day.

154
00:08:00,100 --> 00:08:02,450
You can do the math and figure that out.

155
00:08:02,450 --> 00:08:07,110
You know a complex computation might
need to move a lot of data, and

156
00:08:07,110 --> 00:08:08,420
that can slow the computation down.

157
00:08:08,420 --> 00:08:12,270
So you need a framework that you know,
doesn't move data around so

158
00:08:12,270 --> 00:08:13,680
much while it's doing computation.

159
00:08:15,670 --> 00:08:19,990
The third problem is that distributed
programming can be really really hard.

160
00:08:19,990 --> 00:08:24,370
Even sophisticated programmers find
it hard to write distributed programs

161
00:08:24,370 --> 00:08:28,240
correctly and avoid race conditions and
various kinds of complications.

162
00:08:28,240 --> 00:08:31,770
So here's a simple problem that
hides most of the complexity of

163
00:08:31,770 --> 00:08:33,190
distributed programming.

164
00:08:33,190 --> 00:08:35,960
And, and makes it easy to write you know,

165
00:08:35,960 --> 00:08:38,920
algorithms that can mine
very massive data sets.

166
00:08:38,920 --> 00:08:43,950
So we look at three problems
that you know that we face when,

167
00:08:43,950 --> 00:08:45,800
when we're dealing with cluster computing.

168
00:08:45,800 --> 00:08:51,070
And, Map-Reduce addresses all
three of these challenges.

169
00:08:51,070 --> 00:08:51,730
Right?
First of all,

170
00:08:51,730 --> 00:08:54,810
the first problem that we saw was that,
was one of persistence and

171
00:08:54,810 --> 00:08:56,540
availability of nodes can fade.

172
00:08:56,540 --> 00:09:00,230
The Map-Reduce model addresses this
problem by storing data redundantly on

173
00:09:00,230 --> 00:09:01,000
multiple nodes.

174
00:09:01,000 --> 00:09:04,700
The same data is stored on multiple
nodes so that even if you lose one of

175
00:09:04,700 --> 00:09:06,680
those nodes, the data is still
available on another node.

176
00:09:07,820 --> 00:09:11,970
The second problem that we saw
was one of network bottlenecks.

177
00:09:11,970 --> 00:09:13,986
And this happens when you
move around data a lot.

178
00:09:13,986 --> 00:09:18,970
What the Map-Reduce model does is it
moves the computation close to the data.

179
00:09:18,970 --> 00:09:21,910
And avoids copying data
around the network.

180
00:09:21,910 --> 00:09:24,520
And this minimizes the network
bottle neck problem.

181
00:09:24,520 --> 00:09:27,170
And thirdly,
the Map-Reduce model also provides a very

182
00:09:27,170 --> 00:09:32,070
simple programming model that hides
the complexity of all the online magic.

183
00:09:32,070 --> 00:09:34,870
So let's look at each of
these pieces in turn.

184
00:09:34,870 --> 00:09:38,160
The first piece is the redundant
storage infrastructure.

185
00:09:38,160 --> 00:09:41,900
Now redundant storage is provided by
what's called a distributed file system.

186
00:09:41,900 --> 00:09:46,950
Now distributed file system is a file
system that stores data you know,

187
00:09:46,950 --> 00:09:50,550
across a cluster, but
stores each piece of data multiple times.

188
00:09:50,550 --> 00:09:54,140
So, the distributed file system
provides a global file namespace.

189
00:09:54,140 --> 00:09:56,120
It provides redundancy and availability.

190
00:09:56,120 --> 00:09:59,220
There are multiple implementations
of distributed file systems.

191
00:09:59,220 --> 00:10:03,490
Google's GFS is or Google File System,
or GFS is one example.

192
00:10:03,490 --> 00:10:06,370
Hadoop's HDFS is another example.

193
00:10:06,370 --> 00:10:09,269
And these are the two most popular
distributed file systems out there.

194
00:10:12,230 --> 00:10:16,730
Our typical usage pattern that these
distributed file systems are optimized for

195
00:10:16,730 --> 00:10:18,278
is huge files.

196
00:10:18,278 --> 00:10:20,820
That are in the 100s to, of GB to TB.

197
00:10:21,850 --> 00:10:24,070
But the,
even though the files are really huge,

198
00:10:24,070 --> 00:10:26,380
the data is very rarely updated in place.

199
00:10:26,380 --> 00:10:30,830
Right, once, once data is written you
know it's, it's very, very often.

200
00:10:30,830 --> 00:10:33,450
But when it's updated,
it's updated through appends.

201
00:10:33,450 --> 00:10:35,430
It's never updated in place.

202
00:10:35,430 --> 00:10:41,198
And for example let, let,
imagine the Google scenario once again.

203
00:10:41,198 --> 00:10:46,272
When Google encounters a new webpage it,
it adds the webpage to a depository.

204
00:10:46,272 --> 00:10:47,351
Doesn't ever go and

205
00:10:47,351 --> 00:10:51,170
update the content of the webpage
that it already has crawled, right?

206
00:10:51,170 --> 00:10:55,660
So a typical usage pattern
consists of writing the data once,

207
00:10:55,660 --> 00:10:58,240
reading it multiple times and
appending to it occasionally.

208
00:10:59,380 --> 00:11:03,090
Lets go into the hood of a distributed
file system to see how it actually works.

209
00:11:03,090 --> 00:11:06,010
Data is kept in chunks that
are spread across machines.

210
00:11:06,010 --> 00:11:09,900
So if you take any file,
the file is divided into chunks, and

211
00:11:09,900 --> 00:11:12,120
these chunks are spread
across multiple machines.

212
00:11:12,120 --> 00:11:16,480
So the machines themselves are called
chunk servers in this context.

213
00:11:16,480 --> 00:11:18,010
So here's, here's an example.

214
00:11:18,010 --> 00:11:22,468
There are multiple
multiple chunks servers.

215
00:11:22,468 --> 00:11:25,150
Chunk server 1, 2, 3, and 4.

216
00:11:25,150 --> 00:11:30,060
And here's the file 1.

217
00:11:30,060 --> 00:11:38,091
And file 1 is divided into six chunks in
this case, C0, C1, C2, C3, C4 and C5.

218
00:11:38,091 --> 00:11:42,810
And these chunks as you can see four of
the chunks happen to be on Chunk server 1.

219
00:11:42,810 --> 00:11:47,730
One of them is on Chunks server 2 and,
one of them is on Chunks server 3.

220
00:11:47,730 --> 00:11:49,710
Now this is not sufficient.

221
00:11:49,710 --> 00:11:55,480
You actually have to store multiple
copies of each of these chunks and so

222
00:11:55,480 --> 00:11:59,960
we replicate these chunks so
here copy, here is a copy of C1.

223
00:12:01,080 --> 00:12:04,360
On Chunk server 2,
a copy of C2 in Chunk server 3, and so on.

224
00:12:04,360 --> 00:12:07,830
So each chunk,
in this case is replicated twice.

225
00:12:09,100 --> 00:12:12,980
And if you notice carefully
you'll see that replicas of

226
00:12:12,980 --> 00:12:15,640
a chunk are never on
the same chunk server.

227
00:12:15,640 --> 00:12:18,216
They're always on different chunks of, so

228
00:12:18,216 --> 00:12:22,698
C1 has one replica on Chunk server 1 and
one on Chunk server 2.

229
00:12:22,698 --> 00:12:27,220
C0 has one on Chunk server 1, and
one on Chunk server N, and so on.

230
00:12:28,960 --> 00:12:32,740
And here is here is another file, D.

231
00:12:32,740 --> 00:12:35,680
D has two chunks, D0 and D1.

232
00:12:35,680 --> 00:12:38,040
And that's replicated twice.

233
00:12:38,040 --> 00:12:42,279
And so and so that's stored on
different chunks server [INAUDIBLE].

234
00:12:46,730 --> 00:12:51,850
Now so, so
you serve you serve from chunk files and

235
00:12:51,850 --> 00:12:54,860
store them on, on these,
on these chunk servers.

236
00:12:54,860 --> 00:12:59,096
Now we turn some of the chunk servers,
also act as compute servers.

237
00:12:59,096 --> 00:13:03,430
And when, whenever your
computation has to access data.

238
00:13:03,430 --> 00:13:06,750
That computation is actually
scheduled on the chunk server that

239
00:13:06,750 --> 00:13:09,060
actually contains the data.

240
00:13:09,060 --> 00:13:13,120
This way you avoid moving data to
where the computation needs to run,

241
00:13:13,120 --> 00:13:16,850
but instead you move the computation
to where the data is.

242
00:13:16,850 --> 00:13:22,890
And that's how you put a wide under
the city data movement in the system.

243
00:13:22,890 --> 00:13:25,540
This isn't clear when you look
at look at some examples.

244
00:13:28,920 --> 00:13:33,710
So the sum of this,
each file is split into contiguous chunks.

245
00:13:33,710 --> 00:13:38,850
And the chunks are typically
16 to 64 MB in in size.

246
00:13:38,850 --> 00:13:40,830
On each chunk is replicated,

247
00:13:40,830 --> 00:13:44,580
in our example we saw each
chunk replicated twice.

248
00:13:44,580 --> 00:13:46,790
But it could be 2x or 3x replication.

249
00:13:46,790 --> 00:13:47,990
3x is the most common.

250
00:13:49,210 --> 00:13:54,641
And we saw that the chunks were actually
kept on different chunk servers.

251
00:13:54,641 --> 00:13:59,080
But, but when you replicate 3x, you know,
the system usually makes an effort.

252
00:13:59,080 --> 00:14:04,240
To keep at least one replica in
a entirely different rack if possible and

253
00:14:04,240 --> 00:14:04,970
why do we do that?

254
00:14:04,970 --> 00:14:07,893
We do that because it's you know,

255
00:14:07,893 --> 00:14:12,069
the most common scenario is
that a single node can fail.

256
00:14:12,069 --> 00:14:15,212
But it's also possible that
the switch on a rack can fail, and

257
00:14:15,212 --> 00:14:19,568
when the switch on a rack fails,
the entire rack becomes inaccessible.

258
00:14:19,568 --> 00:14:23,489
And then if you have all the chunks for
a, for in all the replicas of a chunk in

259
00:14:23,489 --> 00:14:26,492
one rack then that whole chunk
can become inaccessible.

260
00:14:26,492 --> 00:14:30,674
So if you keep replicas of a chunk
on different racks then even if

261
00:14:30,674 --> 00:14:34,050
a switch fails then it can
still access that chunk.

262
00:14:34,050 --> 00:14:36,720
Right so
the system tries to make sure that,

263
00:14:36,720 --> 00:14:39,359
that the replicas of a chunk
are actually kept on different racks.

264
00:14:42,240 --> 00:14:45,840
The second component of a distributed
file system is, is a master node.

265
00:14:45,840 --> 00:14:50,110
Now the master node is also known as the,
it's called a master node in

266
00:14:50,110 --> 00:14:53,710
the Google file system, it's a called
a Name Node in Hadoop's HDFS.

267
00:14:53,710 --> 00:14:58,970
But the master node stores metadata
about where the files are stored.

268
00:14:58,970 --> 00:14:59,680
And for

269
00:14:59,680 --> 00:15:05,210
example, if my you know, it'll know that
file one is divided into six chunks.

270
00:15:05,210 --> 00:15:08,260
And here is, here are the locations
of each of the six chunks, and

271
00:15:08,260 --> 00:15:10,160
here are the locations of the replicas.

272
00:15:10,160 --> 00:15:14,170
And the master node itself may be
replicated because otherwise it

273
00:15:14,170 --> 00:15:15,610
might become a single point of failure.

274
00:15:17,170 --> 00:15:20,550
The final component of a distributed
file system is a client library.

275
00:15:20,550 --> 00:15:24,070
Now, when the, when a client,
or, or an algorithm that needs to

276
00:15:24,070 --> 00:15:28,420
access the data tries to access a file
it goes through the client library.

277
00:15:28,420 --> 00:15:30,990
The client library talks to the master and

278
00:15:30,990 --> 00:15:34,188
finds the chunk servers that
actually store the chunks.

279
00:15:34,188 --> 00:15:40,211
And once that's done the client is
directly connected to the chunk servers.

280
00:15:40,211 --> 00:15:43,060
Where it can access the data without
going through the master nodes.

281
00:15:43,060 --> 00:15:46,364
So the data access actually happens
in peer-to-peer fashion without going

282
00:15:46,364 --> 00:15:47,391
through the master node

