1
23:59:59,500 --> 00:00:05,259
[MUSIC]. 

2
00:00:05,259 --> 00:00:09,859
Okay. 
So, the third influential system that 

3
00:00:09,859 --> 00:00:16,122
Rick mentioned in the paper is BigTable 
from Google which is a paper from 2006. 

4
00:00:16,122 --> 00:00:21,320
And so, here we're looking at primary 
index look up, secondary index look up. 

5
00:00:21,320 --> 00:00:24,680
Transactions are also at this sort of 
scale of individual record. 

6
00:00:24,680 --> 00:00:28,964
joint analytic is not supported by 
BigTable directly, but in the, in H base 

7
00:00:28,964 --> 00:00:32,937
which is the open source implementation 
of it. 

8
00:00:32,937 --> 00:00:36,353
And in, well, in, in Google's 
implementation as well, it was sort of 

9
00:00:36,353 --> 00:00:41,610
designed to be compatible with MapReduce. 
And so, you can run MapReduce on, over 

10
00:00:41,610 --> 00:00:45,323
here with the same data that's stored in 
BigTable. 

11
00:00:45,323 --> 00:00:49,105
So, there kind of complementary, and then 
there some notion of integrity 

12
00:00:49,105 --> 00:00:54,200
constraints or schema here, and we'll 
talk about how to implement it. 

13
00:00:54,200 --> 00:00:57,564
There's no views that I can see and 
there's no sort of language level or 

14
00:00:57,564 --> 00:01:01,630
alge, or algebra level for manipulating 
these things. 

15
00:01:01,630 --> 00:01:05,400
There all sort of NoSQL style micro 
interactions with individual records and 

16
00:01:05,400 --> 00:01:08,654
cells, okay? 
So, this is a paper in OSDI for 2006, and 

17
00:01:08,654 --> 00:01:12,285
some overlap with the authors of the 
MapReduce paper. 

18
00:01:12,285 --> 00:01:16,654
And it was sort of designed from the 
start to be complementary of MapReduce. 

19
00:01:16,654 --> 00:01:20,619
So, if you can remember what was one of 
the main things that was missing from 

20
00:01:20,619 --> 00:01:24,520
MapReduce s, or a few things that are 
missing. 

21
00:01:24,520 --> 00:01:30,000
Well, in particular, you couldn't look 
things up by index. 

22
00:01:30,000 --> 00:01:34,742
You couldn't get these little sort of low 
latency accesses. 

23
00:01:34,742 --> 00:01:38,042
So, for example, you want to find all the 
records, you know, given a big data set, 

24
00:01:38,042 --> 00:01:42,720
you want to find all the records in some 
other data set that correspond to. 

25
00:01:42,720 --> 00:01:45,460
You want to do some sort of a join. 
The best you could do is process, you had 

26
00:01:45,460 --> 00:01:48,360
to touch every single record. 
There was no way to zoom in to adjust the 

27
00:01:48,360 --> 00:01:52,500
right one you wanted, okay? 
And so, BigTable provides that fast 

28
00:01:52,500 --> 00:01:58,836
key-based look up, but you could still 
process the overall data as a big set of 

29
00:01:58,836 --> 00:02:05,083
key value records with MapReduce. 
Fine. 

30
00:02:05,083 --> 00:02:07,891
So, the data model here is a sparse, 
distributed, persistent, 

31
00:02:07,891 --> 00:02:11,822
multidimensional, sorted map. 
And what they mean here is that, you can 

32
00:02:11,822 --> 00:02:15,350
basically access any cell in a big table 
by giving a row ID, a column, a column 

33
00:02:15,350 --> 00:02:20,258
name and a time stamp. 
the time stamp isn't really describing 

34
00:02:20,258 --> 00:02:23,650
this English description here, it's for 
versioning. 

35
00:02:23,650 --> 00:02:27,013
So, when you have, after you have 
updates, you'll, you'll keep track of 

36
00:02:27,013 --> 00:02:32,414
past versions of the same cell, okay? 
And so, if you provide these three 

37
00:02:32,414 --> 00:02:39,010
parameters, data table will return you a 
string quickly, all right? 

38
00:02:39,010 --> 00:02:43,924
So each row is data's all sorted 
lexigraphically by the row key, which is 

39
00:02:43,924 --> 00:02:49,780
this row ID in this, in this bit, right? 
So this is a, some sort of primary key, 

40
00:02:49,780 --> 00:02:52,880
and it, you know, in sort of relational 
language or just a key in, in kind of a 

41
00:02:52,880 --> 00:02:58,074
NoSQL framework. 
And then this key range, I say that that 

42
00:02:58,074 --> 00:03:04,238
the integers, the contiguous subranges of 
this set of keys, will be assigned to a 

43
00:03:04,238 --> 00:03:11,678
tablet, all right? 
Okay, so this is a little different way 

44
00:03:11,678 --> 00:03:16,670
of dividing up the data that we've seen 
in the past, in at least one system. 

45
00:03:16,670 --> 00:03:20,382
We talked about a parallel database model 
we happened to use the example from 

46
00:03:20,382 --> 00:03:24,030
Teradata, and so how did they breakup 
data. 

47
00:03:24,030 --> 00:03:29,535
Well, they did it by hashing, right? 
So, every individual record would be sent 

48
00:03:29,535 --> 00:03:34,292
to a server according to a hash function, 
which you can generally just think of as 

49
00:03:34,292 --> 00:03:39,702
sort of a round robin. 
The point is that two keys that are next 

50
00:03:39,702 --> 00:03:45,084
to each other in space, so sort of time 
stamp 5 pm, and time stamp 5:01. 

51
00:03:45,084 --> 00:03:49,180
There's no reason to believe that 5:00 
and 5:01 are going to be on the server in 

52
00:03:49,180 --> 00:03:52,370
Teradata's model. 
Here they are. 

53
00:03:52,370 --> 00:03:56,165
So, what are the pros and cons of this? 
Well, if you're going to typically access 

54
00:03:56,165 --> 00:03:59,276
a whole range of keys at once, it's 
pretty nice to be able to, you know, when 

55
00:03:59,276 --> 00:04:02,639
you get one. 
You get the others, too, sort of for 

56
00:04:02,639 --> 00:04:05,692
free, because you're, you're pulling them 
all back. 

57
00:04:05,692 --> 00:04:11,540
However If one particular key range is 
much more popular than the others, just 

58
00:04:11,540 --> 00:04:18,937
by using the time example again. 
the most recent data perhaps is the most 

59
00:04:18,937 --> 00:04:22,590
popular. 
And so, if all the requests are going to 

60
00:04:22,590 --> 00:04:26,446
that, one key range. 
Then, you got a bunch of idle servers 

61
00:04:26,446 --> 00:04:30,736
hosting all the other tablets that are 
corresponding to older times and all the 

62
00:04:30,736 --> 00:04:37,066
requests are going to this one tablet. 
And so for that reason, Teradata sort of 

63
00:04:37,066 --> 00:04:42,038
chooses to hash everything. 
So that on every request, all the servers 

64
00:04:42,038 --> 00:04:46,729
may have to be accessed, but that's good 
for scalability. 

65
00:04:46,729 --> 00:04:50,872
Okay, so pros and cons. 
Alright, so the tablet here is the unit 

66
00:04:50,872 --> 00:04:52,530
of distribution and load balancing 
fluency. 

67
00:04:52,530 --> 00:04:55,218
And so, they'll move tablets between 
servers as things start to get 

68
00:04:55,218 --> 00:04:58,145
unbalanced, right? 
A key, if you're, if you're key range is, 

69
00:04:58,145 --> 00:05:00,682
you know, January, February, March, 
April, May. 

70
00:05:00,682 --> 00:05:03,517
And there's a whole lot of data coming in 
from March, they'll split that into 

71
00:05:03,517 --> 00:05:06,240
multiple tablets and start, and start 
moving. 

72
00:05:06,240 --> 00:05:10,937
They're moving those tablets around 
between servers in order to balance 

73
00:05:10,937 --> 00:05:14,340
things. 
Okay, so within a single table, you can 

74
00:05:14,340 --> 00:05:18,036
have these groups of columns called 
Column Families. 

75
00:05:18,036 --> 00:05:22,175
And the column names have the family 
right in there as a qualifier. 

76
00:05:22,175 --> 00:05:26,719
And this family is the basic unit of axis 
control so you can provide permissions on 

77
00:05:26,719 --> 00:05:30,824
a group of columns. 
memory accounting in that they are sort 

78
00:05:30,824 --> 00:05:34,440
of allocated as a, as a unit in memory, 
and then disk accounting. 

79
00:05:34,440 --> 00:05:37,105
So, they moved around on disk as a unit 
as well, okay? 

80
00:05:37,105 --> 00:05:40,940
And so, during this point, the typically 
all columns in the family are the same 

81
00:05:40,940 --> 00:05:45,872
type, which I find a little unusual. 
Because they sort of talk about being the 

82
00:05:45,872 --> 00:05:49,624
basic unit of access control and suggests 
that there's, you know, things that go 

83
00:05:49,624 --> 00:05:53,199
together. 
for access control sort of social 

84
00:05:53,199 --> 00:05:56,615
security number and employee ID or 
something may or may not be the same 

85
00:05:56,615 --> 00:05:59,660
type. 
So, there's sort of a logical grouping 

86
00:05:59,660 --> 00:06:02,680
requirement that they seem to be trying 
to meet. 

87
00:06:02,680 --> 00:06:05,440
But then there at the same time, they 
have to be the same type which is very 

88
00:06:05,440 --> 00:06:09,634
technical reasons, especially because 
they want to compress these things. 

89
00:06:09,634 --> 00:06:12,783
So, if you have a whole bunch of integers 
it's easier to compressed, and you have a 

90
00:06:12,783 --> 00:06:16,429
mix of integers and strings. 
So, I think they're trying to kill too 

91
00:06:16,429 --> 00:06:20,273
many birds with one stone here, alright? 
And then, each cell a can be versioned, 

92
00:06:20,273 --> 00:06:23,345
which is the third part of that of that 
key look up, row ID, column name and time 

93
00:06:23,345 --> 00:06:26,633
stamp. 
And each new version increments that time 

94
00:06:26,633 --> 00:06:30,332
stamp and say hey, you here, you can 
enact different kinds of policies. 

95
00:06:30,332 --> 00:06:35,142
Where you only keep the latest inversions 
or you keep only the versions since a 

96
00:06:35,142 --> 00:06:40,596
given, a given time stamp, right? 
So, how these tablets are managed is a 

97
00:06:40,596 --> 00:06:44,330
master will assign the tablets to tablet 
servers. 

98
00:06:44,330 --> 00:06:47,594
And the tablet server handles reads and 
writes from the tablets it controls, 

99
00:06:47,594 --> 00:06:50,262
okay? 
And so, clients communicate directly with 

100
00:06:50,262 --> 00:06:53,364
the tablet server as oppose to having to 
go through the master every time, which 

101
00:06:53,364 --> 00:06:58,260
is helps through scalability, okay? 
And when, when a tablets are to get too 

102
00:06:58,260 --> 00:07:01,330
big, it will split it and load balance 
it, alright? 

103
00:07:01,330 --> 00:07:11,462
So, the metadata keeping track of where 
tablets are located is organize itself in 

104
00:07:11,462 --> 00:07:18,702
another tablet. 
So, there is a root tablet here that 

105
00:07:18,702 --> 00:07:26,858
describes each record in here describes a 
group of records. 

106
00:07:26,858 --> 00:07:32,339
A group of location records in a, you 
know, bigger table, and then each one of 

107
00:07:32,339 --> 00:07:40,335
these metadata tablets gives the location 
of a particular user table, okay? 

108
00:07:40,335 --> 00:07:43,599
And so, this is how you sort of keep 
track hierarchically of where everything 

109
00:07:43,599 --> 00:07:47,451
is at one time. 
and so chubby that they mention in the 

110
00:07:47,451 --> 00:07:52,790
paper is a distributed lock service for 
controlling access to things. 

111
00:07:52,790 --> 00:07:55,760
I'm not going to talk too much about it, 
okay? 

112
00:07:55,760 --> 00:08:01,740
So how, how are reads and writes handled 
in this system? 

113
00:08:01,740 --> 00:08:05,940
Well, there's a table in memory that 
stores a sequence of updates as they 

114
00:08:05,940 --> 00:08:10,430
occur, okay? 
And, a right operation is lo, is, you 

115
00:08:10,430 --> 00:08:15,737
know, adds a record into the memory, 
memory resident table, but it's also 

116
00:08:15,737 --> 00:08:22,172
written to a log for fault tolerance 
purposes, okay? 

117
00:08:22,172 --> 00:08:24,532
So, if this [UNKNOWN] ever goes down, it 
reads the tablet log and you can 

118
00:08:24,532 --> 00:08:29,440
reconstruct what's going on. 
And then read operations are served by 

119
00:08:29,440 --> 00:08:34,592
reading these SSTable files. 
You have actual data itself, but then 

120
00:08:34,592 --> 00:08:38,750
also by applying the updates from the 
memtable on the fly, right? 

121
00:08:38,750 --> 00:08:42,105
So it needs, it needs a stream, it says 
here's the value and then here's the 

122
00:08:42,105 --> 00:08:45,735
stream of updates. 
I can do apply that value to get the true 

123
00:08:45,735 --> 00:08:48,722
value, okay? 
And then, there's two, so this, so this 

124
00:08:48,722 --> 00:08:52,077
is fine but, but what happens when the 
memtable gets bigger and bigger and 

125
00:08:52,077 --> 00:08:55,950
bigger? 
Well, there's two kinds of events that 

126
00:08:55,950 --> 00:08:59,530
occur to, you know, do the bookkeeping 
here. 

127
00:08:59,530 --> 00:09:04,185
So, one is a minor compaction. 
And this is when the memtable gets big, 

128
00:09:04,185 --> 00:09:08,910
it gets written out into an SS, into a 
new SS Table file and the changes are 

129
00:09:08,910 --> 00:09:13,904
merged okay? 
And then a major compaction is, take all 

130
00:09:13,904 --> 00:09:18,500
the SS Tables and rewrite them all into 
one big one. 

131
00:09:18,500 --> 00:09:22,020
They may be split into multiple files, 
and also clean up any deletes that have 

132
00:09:22,020 --> 00:09:24,875
occurred. 
So, deletes are just appended as 

133
00:09:24,875 --> 00:09:28,670
instructions but aren't necessarily, 
doesn't actually remove anything, so they 

134
00:09:28,670 --> 00:09:33,513
are sort of garbage collected, okay? 
So in this way, you can keep sort of the 

135
00:09:33,513 --> 00:09:37,986
read throughput pretty high. 
And for, to keep this upkeep going on in 

136
00:09:37,986 --> 00:09:41,745
the background, alright? 
So, those are a host of other tricks 

137
00:09:41,745 --> 00:09:44,990
here, too, where they, can do various 
forms of compression, specify by 

138
00:09:44,990 --> 00:09:48,870
compliance, which can be specified by the 
clients. 

139
00:09:48,870 --> 00:09:53,110
There's some different ways of doing it. 
They use bloom filters to speed up 

140
00:09:53,110 --> 00:09:56,977
existence test. 
So, if I give you row ID, a column ID and 

141
00:09:56,977 --> 00:10:01,527
a time stamp and say find me this value, 
what these bloom filters allow you to do 

142
00:10:01,527 --> 00:10:08,880
are, is to very quickly determine whether 
that does not exist in the system. 

143
00:10:08,880 --> 00:10:12,780
So, these bloom filter data structures 
are pretty cool. 

144
00:10:12,780 --> 00:10:14,956
And I'm going to walk through them in 
this course in a, in a couple of weeks, 

145
00:10:14,956 --> 00:10:18,325
okay? 
So, they help you quickly determine 

146
00:10:18,325 --> 00:10:23,850
whether that key does not exist in the 
system, it avoids disk accesses during, 

147
00:10:23,850 --> 00:10:28,428
during reads. 
Alright, and then there's locality groups 

148
00:10:28,428 --> 00:10:32,100
which you can define another layer of 
organization on top of families. 

149
00:10:32,100 --> 00:10:35,461
And these are groups of column families 
that tend to be accessed together. 

150
00:10:35,461 --> 00:10:37,693
Fine. 
And then, another trick here is, is to 

151
00:10:37,693 --> 00:10:42,038
make sure that the SS Tables, these disk 
chunks are immutable. 

152
00:10:42,038 --> 00:10:46,446
They never get written indirectly. 
The only time they get written is when 

153
00:10:46,446 --> 00:10:52,210
these major compactions happen and the 
whole thing is sort of reorganized. 

154
00:10:52,210 --> 00:10:56,170
And so, that means that the only writable 
data structure is this memtable. 

155
00:10:56,170 --> 00:11:02,085
And so, the amount of concurrency control 
to keep things un, remains, remains 

156
00:11:02,085 --> 00:11:05,146
pretty simple, okay? 

