1
00:00:00,290 --> 00:00:02,250
Welcome back to Mining
of Massive Datasets.

2
00:00:02,250 --> 00:00:06,190
In, the previous section will be
studied the Map-Reduce model and

3
00:00:06,190 --> 00:00:08,890
how to solve some simple
problems using Map-Reduce.

4
00:00:08,890 --> 00:00:11,740
In this section, we're going to go under
the hood of a Map-Reduce system and

5
00:00:11,740 --> 00:00:13,190
understand how it actually works.

6
00:00:16,850 --> 00:00:20,983
Just to refresh your
memory a Map-Reduce system

7
00:00:20,983 --> 00:00:24,660
has simple Map-Reduce
system has three steps.

8
00:00:24,660 --> 00:00:29,250
In the Map step,
you take a Big a document which is

9
00:00:29,250 --> 00:00:35,300
divided into chunks and
you run a Map process on each chunk.

10
00:00:35,300 --> 00:00:39,680
And the map process go through
each record in that chunk and

11
00:00:39,680 --> 00:00:44,430
it outputs an intermediate key value pairs
for each vector in that, in that chunk.

12
00:00:45,810 --> 00:00:51,210
In the second set step which is
a group by step you group by key.

13
00:00:51,210 --> 00:00:55,670
You you bring together all the values for,
for the same key.

14
00:00:57,330 --> 00:00:59,490
And in the third step, is a reduce step.

15
00:01:00,580 --> 00:01:06,170
You apply a reducer to each
intermediate key value pair set.

16
00:01:06,170 --> 00:01:07,650
And you create a final output.

17
00:01:08,960 --> 00:01:14,680
Now, here's a schematic of how it actually
works in in a distributed system.

18
00:01:14,680 --> 00:01:19,410
The previous schematic was how it
worked in a, in a centralized system.

19
00:01:19,410 --> 00:01:23,270
In a distributed system, you actually,
have multiple nodes and map and

20
00:01:23,270 --> 00:01:27,190
reduced tasks are running in
pattern on multiple nodes.

21
00:01:27,190 --> 00:01:34,880
So the here are the few chunks of the
file, input file might be on on node 1.

22
00:01:34,880 --> 00:01:37,870
Few chunks on node 2 and
a few chunks on node 3.

23
00:01:37,870 --> 00:01:41,634
And you have map tasks running on,
on each of those nodes.

24
00:01:41,634 --> 00:01:45,969
And and producing producing it to be
intermediate key value pairs on each of

25
00:01:45,969 --> 00:01:46,713
those nodes.

26
00:01:47,850 --> 00:01:48,918
And once the,

27
00:01:48,918 --> 00:01:54,750
once the intermediate key value pairs
are produced, the underlying system

28
00:01:54,750 --> 00:01:58,830
the Map-Reduce system uses a partitioning
function which is just a hash function.

29
00:01:59,930 --> 00:02:03,850
So the the the Map-Reduce system

30
00:02:03,850 --> 00:02:08,180
applies a hash function to
each intermediate key value.

31
00:02:08,180 --> 00:02:12,200
And the has function will tell
the Map-Reduce system which,

32
00:02:12,200 --> 00:02:15,290
reduce node to send
that key value pair to.

33
00:02:15,290 --> 00:02:19,860
Right, this ensures that all all the,
the same key values,

34
00:02:19,860 --> 00:02:24,740
whether they are map task 1, 2, or 3 end
up being sent to the same reduce task.

35
00:02:24,740 --> 00:02:27,839
Right?
So, in this case the key key 4.

36
00:02:27,839 --> 00:02:32,110
Regardless of where it started from,
whether at 1, 2, or 3.

37
00:02:32,110 --> 00:02:34,110
Always end up at reduce task 1.

38
00:02:34,110 --> 00:02:38,690
And the key,
key 1 always ends up at reduce task 2.

39
00:02:38,690 --> 00:02:43,400
Now, once once the reduce
task has a reduce task has

40
00:02:43,400 --> 00:02:48,540
received input from all
from all the map tasks.

41
00:02:48,540 --> 00:02:52,680
All the map tasks have completed,
then you can start the reduced tasks.

42
00:02:52,680 --> 00:02:55,770
And the, the reduced tasks first job is,

43
00:02:55,770 --> 00:03:00,190
is to sort, it's input, and
group it together by key.

44
00:03:01,410 --> 00:03:05,180
And so in this case, there are three
values associated with the key key, key 4,

45
00:03:05,180 --> 00:03:07,170
they're all grouped together.

46
00:03:07,170 --> 00:03:10,980
And once that is done, the reduce task
then, works the reduce function which is

47
00:03:10,980 --> 00:03:18,320
provided by the programmer on each each
such group and creates the final output.

48
00:03:18,320 --> 00:03:19,490
Okay.

49
00:03:19,490 --> 00:03:24,320
So remember, the programmer provides
two functions, Map and Reduce, and

50
00:03:24,320 --> 00:03:26,260
specifies the input file.

51
00:03:26,260 --> 00:03:29,610
The Map-Reduce environment take,
has to take care of a bunch of things.

52
00:03:29,610 --> 00:03:32,360
It takes care of
Partitioning the input data.

53
00:03:32,360 --> 00:03:35,320
Scheduling the program's
execution on a set of machines.

54
00:03:35,320 --> 00:03:39,270
Figuring out where the map tasks run,
where the reduce tasks run, and so on.

55
00:03:40,430 --> 00:03:42,870
It performs a gr,
the intermediate group by step.

56
00:03:44,210 --> 00:03:47,960
And while all this is going
on some nodes may fail.

57
00:03:47,960 --> 00:03:52,090
And the environment make sure
that the node failures are hidden

58
00:03:52,090 --> 00:03:54,370
from the from the program.

59
00:03:54,370 --> 00:03:58,300
And finally the Map-Reduced Environment
also Manages all

60
00:03:58,300 --> 00:04:00,461
the required inter-machine communication.

61
00:04:00,461 --> 00:04:03,035
[SOUND] So, we're going to take,

62
00:04:03,035 --> 00:04:07,894
take a look at exactly how what's,
what's going on in a.

63
00:04:07,894 --> 00:04:11,700
So, let's look at the data flow that's
associated with with, with map reduce.

64
00:04:12,790 --> 00:04:15,680
Now the the input and

65
00:04:15,680 --> 00:04:20,470
the final output of a Map-Reduced program
are stored on the distributed file system.

66
00:04:21,520 --> 00:04:25,770
And the scheduler tries
to schedule the map task

67
00:04:25,770 --> 00:04:29,030
close to the physical storage
location of the import data.

68
00:04:29,030 --> 00:04:34,580
What that means is that recall
the input data is, is a file.

69
00:04:34,580 --> 00:04:36,420
And the file is divided into chunks.

70
00:04:36,420 --> 00:04:39,996
And there are replicas of the chunks
on different chunk servers.

71
00:04:39,996 --> 00:04:43,190
The Map-Reduce system try to schedule each

72
00:04:43,190 --> 00:04:47,860
map task on a chunk server that holds
a copy of the corresponding chunk.

73
00:04:47,860 --> 00:04:49,310
So, there's no actual copy.

74
00:04:50,440 --> 00:04:58,220
A data copy associated with the map
step of the Map-Reduce program.

75
00:04:58,220 --> 00:05:02,380
Now, the intermediate results are, are at
least not stored in the distributed file

76
00:05:02,380 --> 00:05:06,525
system but stored in the local file
system of the map and reduce workers.

77
00:05:06,525 --> 00:05:08,390
what, what are intermediate results?

78
00:05:08,390 --> 00:05:11,820
Intermediate results, intermediate results
could be the output of a map step.

79
00:05:13,880 --> 00:05:16,600
An intermediate result
could be something that,

80
00:05:16,600 --> 00:05:20,150
that limited why, why in the process
of computing the reduce.

81
00:05:20,150 --> 00:05:24,070
Now why, why are such debated results not
stored in the distributed file system?

82
00:05:24,070 --> 00:05:27,770
It turns out that there's some
overhead to storing data in

83
00:05:27,770 --> 00:05:29,450
the distributed file system.

84
00:05:29,450 --> 00:05:32,880
Remember there are multiple replicas
of the data that need to be made.

85
00:05:32,880 --> 00:05:35,910
And so there's a lot of copying.

86
00:05:35,910 --> 00:05:38,240
And network shuffling involved in,

87
00:05:38,240 --> 00:05:41,230
in storing new data in
the distributed file system.

88
00:05:41,230 --> 00:05:44,770
So, whenever possible, intermediate
results are actually stored in

89
00:05:44,770 --> 00:05:47,400
the local file system of the Map and
Reduced workers,

90
00:05:47,400 --> 00:05:52,020
ended up being stored in the distributed
file system to avoid more network traffic.

91
00:05:55,780 --> 00:06:01,100
And finally, as you'll see in
future examples the output

92
00:06:01,100 --> 00:06:05,080
of a Map-Reduce task is often being
the input to another Map-Reduce task.

93
00:06:10,440 --> 00:06:16,080
Now, the master node takes care of all the
coordination aspects of a Map-Reduce job.

94
00:06:16,080 --> 00:06:21,400
The master node keeps, you know,
associates a task status with each task.

95
00:06:21,400 --> 00:06:24,490
A task to see the map tasker reduce task.

96
00:06:24,490 --> 00:06:27,740
And each task has has a status flag.

97
00:06:27,740 --> 00:06:31,570
And the status flag can either be idle,
in progress, or completed.

98
00:06:32,854 --> 00:06:40,210
The master schedules idle tasks
whenever workers become available.

99
00:06:40,210 --> 00:06:44,940
Whenever, there is a free a node that
is tha, that's available for, for

100
00:06:44,940 --> 00:06:46,420
scheduling tasks.

101
00:06:46,420 --> 00:06:51,270
The master goes through it's queue of idle
tasks, and schedules an idle task on that,

102
00:06:51,270 --> 00:06:51,770
on that worker.

103
00:06:53,402 --> 00:06:58,787
When the, when a map task completes,
it sends the the master the location and

104
00:06:58,787 --> 00:07:03,020
sizes of it's the R intermediate
files that it, that creates.

105
00:07:03,020 --> 00:07:05,750
Now, why, R intermediate files?

106
00:07:05,750 --> 00:07:09,620
There's one intermediate file
that's created for each reducer.

107
00:07:10,990 --> 00:07:15,950
Because the data, the output of the mapper
has to be shipped to each of the reducers,

108
00:07:15,950 --> 00:07:18,470
depending on the, on the key value.

109
00:07:18,470 --> 00:07:21,450
And so there R intermediate files,
one for each reducer.

110
00:07:21,450 --> 00:07:24,985
So, whenever, a map task completes,
it let it's, it's,

111
00:07:24,985 --> 00:07:27,550
it's stores the R intermediate files.

112
00:07:27,550 --> 00:07:28,710
On it's local file system,

113
00:07:28,710 --> 00:07:32,930
and it let's the master know what
the names of those files are.

114
00:07:32,930 --> 00:07:35,520
The master pushes this inf,
information to the reducers.

115
00:07:36,630 --> 00:07:42,250
Once the reducers know that all
the mappers map tasks are completed,

116
00:07:42,250 --> 00:07:47,840
then they copy the intermediate
file from each of the map tasks.

117
00:07:47,840 --> 00:07:49,349
And then they can proceed with their work.

118
00:07:50,870 --> 00:07:52,330
Now, the master also per,

119
00:07:52,330 --> 00:07:56,560
periodically pings the workers,
to detect whether a worker has failed.

120
00:07:56,560 --> 00:07:59,310
And if a worker has failed,
the master has to do something.

121
00:07:59,310 --> 00:08:00,860
And we're going to,
see what that something is.

122
00:08:05,420 --> 00:08:12,120
If a map worker fails, then the,
all the map tasks that were scheduled.

123
00:08:12,120 --> 00:08:15,630
On that on that map
worker may have failed.

124
00:08:15,630 --> 00:08:20,380
So, the the tricky thing is that
the output of a map task is written to

125
00:08:20,380 --> 00:08:23,050
the local file system of the,
of the map worker.

126
00:08:23,050 --> 00:08:26,970
So, if a map worker fails,
then the node fails.

127
00:08:26,970 --> 00:08:30,470
Then all intermediate output created
by all the map tasks that have

128
00:08:30,470 --> 00:08:33,320
ran on that worker, are lost.

129
00:08:33,320 --> 00:08:38,830
And so the, what the master does,
is that it resets to idle,

130
00:08:38,830 --> 00:08:44,090
the status of every task that was either
completed or in progress on that worker.

131
00:08:44,090 --> 00:08:48,410
Right, and so all those tasks need to be,
eventually be done, and

132
00:08:48,410 --> 00:08:51,209
they will eventually be rescheduled
on other workers in the course.

133
00:08:54,100 --> 00:08:56,660
If a reduced worker
fails on the other hand,

134
00:08:56,660 --> 00:08:58,770
only the in progress
tasks are set to idle.

135
00:08:58,770 --> 00:09:01,660
The tasks that are actually been
completed by the reduced worker,

136
00:09:01,660 --> 00:09:03,190
don't need to be set to idle.

137
00:09:03,190 --> 00:09:06,450
Because, the output of the reduced
worker is a final output, and

138
00:09:06,450 --> 00:09:08,550
it's written to
the distribute file system.

139
00:09:08,550 --> 00:09:11,020
And not to the local file
system of the reduced worker.

140
00:09:11,020 --> 00:09:14,330
Since, the output is written to
the distributed file system.

141
00:09:14,330 --> 00:09:16,810
The output is not lost even
if the reduce worker fails.

142
00:09:16,810 --> 00:09:20,150
So, only in-progress tasks
need to be set to idle.

143
00:09:20,150 --> 00:09:22,970
While completed tasks
don't need to be redone.

144
00:09:22,970 --> 00:09:24,140
Right?
And so, the and

145
00:09:24,140 --> 00:09:28,400
once again the Idle reduce tasks will be
restarted on other workers eventually.

146
00:09:29,430 --> 00:09:31,230
What happens if the master fails?

147
00:09:31,230 --> 00:09:35,490
If the master node fails,
then the map reduce tas, task is aborted.

148
00:09:35,490 --> 00:09:36,990
The client is notified, and

149
00:09:36,990 --> 00:09:41,630
the client can then do something
like restarting the map reduce task.

150
00:09:41,630 --> 00:09:44,870
So, this is the one scenario
where the task will have to be

151
00:09:44,870 --> 00:09:46,090
restarted from scratch.

152
00:09:46,090 --> 00:09:50,590
Because, the master is typically not
applicated in the Map-Reduce system.

153
00:09:50,590 --> 00:09:54,710
So, you might think that, this is
a big deal, that that the, the master

154
00:09:54,710 --> 00:10:00,230
failure means the the map-reduce task is
aborted, and the task has to be restarted.

155
00:10:00,230 --> 00:10:03,200
But remember,
node failures are actually, rather rare.

156
00:10:03,200 --> 00:10:07,880
A node fails actually recall once every
three years, or once every 1,000 days.

157
00:10:07,880 --> 00:10:12,580
And the master is, is a single node,
and therefore, the chance of the master

158
00:10:12,580 --> 00:10:17,020
failing is actually quite, you know, it,
it, it's quite an uncommon occurrence.

159
00:10:17,020 --> 00:10:19,464
the, the,
the problem that you have with if,

160
00:10:19,464 --> 00:10:23,780
you have a multiple workers associated in,
in a map reduce task.

161
00:10:23,780 --> 00:10:28,130
It's much more likely that,
one of many workers failed,

162
00:10:28,130 --> 00:10:29,450
rather than the master failing.

163
00:10:31,100 --> 00:10:33,130
So, the final question to think about is,

164
00:10:33,130 --> 00:10:36,692
how many map and
how many reduced jobs do we need?

165
00:10:36,692 --> 00:10:44,135
[NOISE] Supposed you know, they're both
throughout M map tasks and R reduce tasks.

166
00:10:44,135 --> 00:10:46,860
Our goal is to determine M and R.

167
00:10:46,860 --> 00:10:51,680
The, this is part of the input that
given to the map reduce system to let it

168
00:10:51,680 --> 00:10:55,600
know how many tasks tasks
it needs to schedule.

169
00:10:55,600 --> 00:11:02,390
The Rule of thumb is to make M much larger
than the number of nodes in the cluster.

170
00:11:02,390 --> 00:11:07,226
You might think, that it's sufficient how
one map task per node to the cluster.

171
00:11:07,226 --> 00:11:11,560
But, in fact, it the rule of thumb is
to have one map task per DFS chunk.

172
00:11:12,780 --> 00:11:14,820
The reason for this is simple.

173
00:11:14,820 --> 00:11:19,040
Imagine, that there is one map
task per node in the cluster and

174
00:11:19,040 --> 00:11:23,620
during you know during
processing the node fails.

175
00:11:23,620 --> 00:11:27,840
If a node fails then that map
task needs to be rescheduled.

176
00:11:27,840 --> 00:11:31,150
On another node in,
in the cluster when it becomes available.

177
00:11:31,150 --> 00:11:35,870
Now in, some, since all the other
nodes are processing, you know,

178
00:11:35,870 --> 00:11:41,250
one of the map tasks has to, one of those
nodes has to complete before this map task

179
00:11:41,250 --> 00:11:47,440
can be scheduled on that node and so,
the entire computation is slowed down.

180
00:11:47,440 --> 00:11:50,780
By the time it takes to com,
you know, complete this map task.

181
00:11:50,780 --> 00:11:52,700
The failed redo the failed map task.

182
00:11:54,070 --> 00:11:58,769
So, if instead of one map task on a given
node, there are many small map tasks on

183
00:11:58,769 --> 00:12:03,613
a given node, and that node fails, then
those map tasks can then be spread across

184
00:12:03,613 --> 00:12:08,343
all the available nodes and so
the entire task will complete much faster.

185
00:12:11,950 --> 00:12:16,126
On the other hand, the number produces
R is usually smaller than M and

186
00:12:16,126 --> 00:12:20,590
is usually even smaller than the total
number of nodes in the system.

187
00:12:20,590 --> 00:12:24,210
And this because the the output file is,

188
00:12:24,210 --> 00:12:29,720
is spread across spread across R
node where R the number of reducers.

189
00:12:29,720 --> 00:12:35,630
And if it's usually convenient
to have the output spread across

190
00:12:35,630 --> 00:12:38,980
a small number of nodes rather than
across a large number of nodes.

191
00:12:38,980 --> 00:12:41,850
And so usually R is set to
a smaller value than M.

