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

2
00:00:05,636 --> 00:00:10,100
Okay, so Dynamo from Amazon in [SOUND] 
2007, which is a paper and again a few 

3
00:00:10,100 --> 00:00:16,078
years later, it's been released as a 
cloud service called DynamoDB. 

4
00:00:16,078 --> 00:00:20,460
Okay, so here we're looking you know, 
scale, 2000 of nodes. 

5
00:00:20,460 --> 00:00:26,240
you can look up things by primary index, 
and basically nothing else, just a key 

6
00:00:26,240 --> 00:00:32,340
value store, just like nimcash, okay? 
Right. 

7
00:00:32,340 --> 00:00:35,120
So ,what are some of the tricks it let 
to? 

8
00:00:35,120 --> 00:00:36,908
Well, so some key features are that it 
has. 

9
00:00:36,908 --> 00:00:39,273
It's sort of some, so talking to, in 
terms of DynamoDB, which is the 

10
00:00:39,273 --> 00:00:42,130
implementation you can now go and use, 
and pay for. 

11
00:00:42,130 --> 00:00:45,668
One of the neat things here, is that it 
offers a service level agreement on 

12
00:00:45,668 --> 00:00:49,586
performance. 
And so, at the 99th percentile, you know, 

13
00:00:49,586 --> 00:00:55,243
they promise to respond within 300 
milliseconds for 99.9% of its request. 

14
00:00:55,243 --> 00:00:59,397
And the reason they do this on the 99th 
percentile as opposed to some sort of 

15
00:00:59,397 --> 00:01:03,551
notion of the average the mean or the 
medians, or the mean and the median, is 

16
00:01:03,551 --> 00:01:11,096
that would artificially penalize the 
people who are using it heavily, right. 

17
00:01:11,096 --> 00:01:15,140
They would get a disproportionate number 
of failed requests. 

18
00:01:15,140 --> 00:01:20,420
It would be easy to satisfy the average 
by only focusing only on the lightweight 

19
00:01:20,420 --> 00:01:25,280
users for example, okay. 
So Dynamo, the system is a distributive 

20
00:01:25,280 --> 00:01:30,082
hash table, that's what DHT stands for. 
And each key is stored at, or sorry, each 

21
00:01:30,082 --> 00:01:35,114
value is stored at locations, multiple 
locations for replication purposes, and 

22
00:01:35,114 --> 00:01:41,594
its up for replication factor of N. 
And so it at location K, K plus 1 all the 

23
00:01:41,594 --> 00:01:46,880
way up to K plus N minus 1. 
And they achieve eventual consistency 

24
00:01:46,880 --> 00:01:52,512
through vector clots, which I'll describe 
in the next couple of slides. 

25
00:01:52,512 --> 00:01:57,248
And so reconciliation of potential 
conflicts when things are being read and 

26
00:01:57,248 --> 00:02:00,716
written. 
Happens at read time, which is another 

27
00:02:00,716 --> 00:02:04,271
maybe interesting feature of Dynamo. 
Okay, so, rights never fail and they site 

28
00:02:04,271 --> 00:02:07,970
in the paper, that the reason for this is 
poor customer experience, alright. 

29
00:02:07,970 --> 00:02:12,269
So, if you're sort of typing into your 
Google Doc, well, [LAUGH] [INAUDIBLE]. 

30
00:02:12,269 --> 00:02:15,965
if you're billing application, a web 
application where you say you update your 

31
00:02:15,965 --> 00:02:19,605
status on some social networking site, 
and it comes back with an error message 

32
00:02:19,605 --> 00:02:26,080
and says, Sorry couldn't commit, you 
know, somebody else was editing the same. 

33
00:02:26,080 --> 00:02:28,290
Or you're editing the same status from 
somewhere else. 

34
00:02:28,290 --> 00:02:35,668
Their claim is that that's more 
disruptive than getting the, you know, 

35
00:02:35,668 --> 00:02:43,405
the wrong read, which seems reasonable to 
me, okay. 

36
00:02:43,405 --> 00:02:46,349
All right, and so, conflict resolution 
for many applications may be the most 

37
00:02:46,349 --> 00:02:51,397
recent write is the one that wins. 
Or you can actually have the application 

38
00:02:51,397 --> 00:02:53,991
controlled. 
In some cases, you may even sort of go 

39
00:02:53,991 --> 00:02:57,660
back and ask the user to resolve the 
conflict manually, okay. 

40
00:02:57,660 --> 00:03:01,690
Okay, so the goal with Vector Clocks is 
to detect conflicts in a concurrent 

41
00:03:01,690 --> 00:03:05,256
read-write scenario. 
But not to necessarily do anything about 

42
00:03:05,256 --> 00:03:08,247
them automatically. 
Okay, so, In this scheme every data item 

43
00:03:08,247 --> 00:03:11,728
is associated with a list of server 
timestamp pairs that indicates its 

44
00:03:11,728 --> 00:03:17,725
version history. 
And so, in this example, [SOUND] sum 

45
00:03:17,725 --> 00:03:27,470
value D was read by a client and D1 was 
written back at the server called SX. 

46
00:03:27,470 --> 00:03:32,970
And so what SX does is append this fact 
to this vector clock. 

47
00:03:32,970 --> 00:03:38,350
So, at timestamp one, server SX, you 
know, created a change. 

48
00:03:38,350 --> 00:03:42,370
Then some other client reads D1 and 
writes back D2. 

49
00:03:42,370 --> 00:03:51,492
And, you know, you might want to, append 
to the vector clock, both values. 

50
00:03:51,492 --> 00:03:57,830
But, you know, this change descends from, 
D2 descends from D1. 

51
00:03:57,830 --> 00:04:02,118
It was handled by the same server, and 
so, you can garbage collect this part of 

52
00:04:02,118 --> 00:04:06,042
the vector clock right? 
So, it's the same server with a higher 

53
00:04:06,042 --> 00:04:08,772
timestamp, means that the old version 
where the older timestamp is not needed 

54
00:04:08,772 --> 00:04:11,572
anymore. 
Okay, and since there were no other 

55
00:04:11,572 --> 00:04:16,074
conflicts to work on, okay. 
But now, independently two different 

56
00:04:16,074 --> 00:04:20,620
clients read D2, and write back different 
values. 

57
00:04:20,620 --> 00:04:23,870
One writes back D4, and one writes back 
D3, and these two requests were handled 

58
00:04:23,870 --> 00:04:27,853
by different servers, SY and SZ. 
And so these facts get recorded in the 

59
00:04:27,853 --> 00:04:33,477
vector clocks since they're different. 
Now, the contexts here will, this, we 

60
00:04:33,477 --> 00:04:39,708
call this sort of vector clock contacts. 
The contacts will reflect this fact when 

61
00:04:39,708 --> 00:04:42,076
the next read comes in, it will see that, 
oh wait, there's a conflict, because 

62
00:04:42,076 --> 00:04:44,902
there's the same timestamp but two 
different servers. 

63
00:04:44,902 --> 00:04:48,466
And you can either ask the client what to 
do or you apply some of your instinct 

64
00:04:48,466 --> 00:04:53,425
where the later one run, because these 
might not be timestamps like integers. 

65
00:04:53,425 --> 00:04:57,269
They could be sort of the actual clock 
time stamps, in which case you make an 

66
00:04:57,269 --> 00:05:01,180
arbitrary decision, just pick it and go, 
okay. 

67
00:05:01,180 --> 00:05:05,946
So, that's how vector clocks work. 
So, in the example, just to run through 

68
00:05:05,946 --> 00:05:09,028
what we just saw. 
A client writes D1 to server SX and 

69
00:05:09,028 --> 00:05:13,033
creates this value. 
Another client writes D, reads D1 and 

70
00:05:13,033 --> 00:05:17,685
writes back D2, also handled by SX, and 
D1 was garbage collected. 

71
00:05:17,685 --> 00:05:23,205
Then separate clients read D2 and write 
back D3 and D4, and two different 

72
00:05:23,205 --> 00:05:27,640
servers, SY and SZ. 
And then, another client reads D3 and D4, 

73
00:05:27,640 --> 00:05:30,240
and notice, it then finds that there's a 
con, the system reports that there's a 

74
00:05:30,240 --> 00:05:34,935
conflict to be handled. 
Okay, so let's practice with these. 

75
00:05:34,935 --> 00:05:39,022
Here's, [SOUND] two different vector 
clocks and you, figure out whether 

76
00:05:39,022 --> 00:05:43,490
there's a con, whether they represent a 
conflict or not. 

77
00:05:43,490 --> 00:05:47,458
So, in this case we have server SX with a 
timestamp of 3, and on this data server 

78
00:05:47,458 --> 00:05:52,235
Sx with a timestamp of 3. 
And then each one is a different server 

79
00:05:52,235 --> 00:05:56,080
with different timestamps. 
So, is there a conflict? 

80
00:05:56,080 --> 00:06:00,184
Well, yeah, there is, because on one 
version path, SY made a series of 

81
00:06:00,184 --> 00:06:05,980
changes, and on another version path, SZ 
made a series of changes. 

82
00:06:05,980 --> 00:06:08,185
And they didn't talk to each other, 
because they don't reflect each others' 

83
00:06:08,185 --> 00:06:09,922
changes. 
So, yes, there is. 

84
00:06:09,922 --> 00:06:15,020
And on this one it have the same server 
at a later timestamp. 

85
00:06:15,020 --> 00:06:19,274
So, is there a conflict here? 
Well no, because they weren't handled by 

86
00:06:19,274 --> 00:06:22,864
different servers. 
So, really just this one subsumes that 

87
00:06:22,864 --> 00:06:28,149
one, and we're okay. 
So, on this one we have server SX with 3, 

88
00:06:28,149 --> 00:06:35,900
server SX with 3, server SY with 6, 
server SY with 6 so they agree so far. 

89
00:06:35,900 --> 00:06:41,340
And this was an extra change of S, of 
server SZ, with the timestamp of 2. 

90
00:06:41,340 --> 00:06:46,224
And so, no, there's no conflict here, 
because they agree wherever there's, on, 

91
00:06:46,224 --> 00:06:51,040
on, this is just an extra change on top 
of this one. 

92
00:06:51,040 --> 00:06:55,106
So this guy wins, okay. 
In this next one, server SX with 

93
00:06:55,106 --> 00:07:01,110
timestamp of 3, server SX with a 
timestamp of 3. 

94
00:07:01,110 --> 00:07:05,388
Server SY with a time stamp of 10, and 
this has servers with a timestamp of 6, 

95
00:07:05,388 --> 00:07:10,860
and then some later change at SZ. 
So, is there a conflict here? 

96
00:07:10,860 --> 00:07:15,944
Well there is because this one is later 
in time on SY, right, 10 then this one 

97
00:07:15,944 --> 00:07:20,560
is. 
But then, if that was all there was. 

98
00:07:20,560 --> 00:07:23,208
If it wasn't for this guy, if it wasn't 
for this guy here. 

99
00:07:23,208 --> 00:07:29,397
I'd be okay, we would just pick this one, 
because it's later. 

100
00:07:29,397 --> 00:07:35,400
But, because this one's here, we now have 
some changes at SZ. 

101
00:07:35,400 --> 00:07:37,701
And some changes SY they were both 
forward in time from the latest point 

102
00:07:37,701 --> 00:07:40,600
that we've agreed on. 
And so we don't know how to resolve that. 

103
00:07:40,600 --> 00:07:45,592
And so yes there is a conflict. 
And them similarly here, SX in timestamp 

104
00:07:45,592 --> 00:07:50,052
with 3 and SX in timestamp with 3. 
Well here we have SY and SY 10 is the 

105
00:07:50,052 --> 00:07:53,160
same as the last one. 
But here instead of 6, it's 20. 

106
00:07:53,160 --> 00:07:58,500
And so it's later than this 10, and then 
further we have a change at SZ. 

107
00:07:58,500 --> 00:08:03,702
And so is there a conflict here? 
well, no, because this one is strictly 

108
00:08:03,702 --> 00:08:07,065
later than that one is. 
On all the servers that they share it has 

109
00:08:07,065 --> 00:08:10,900
later timestamps. 
So it's strictly subsumitive. 

110
00:08:10,900 --> 00:08:18,188
And so no, there is no confluence. 
Ok, so those are, in the last segment we 

111
00:08:18,188 --> 00:08:22,780
talked about consistent hashing, in this 
segment we talked about vector clocks. 

112
00:08:22,780 --> 00:08:26,630
These are two little gadgets to be 
familiar with. 

113
00:08:26,630 --> 00:08:29,880
Because they come up time and again in 
these no sequel systems, and in other 

114
00:08:29,880 --> 00:08:33,941
systems in general, okay. 
All right, so Dynamo also talks about a 

115
00:08:33,941 --> 00:08:37,270
way to parameterize the level of 
consistency. 

116
00:08:37,270 --> 00:08:40,990
And this comes up, occasionally in 
papers, so I just want to make sure 

117
00:08:40,990 --> 00:08:46,360
you're exposed to it. 
So, the idea here is that you have two 

118
00:08:46,360 --> 00:08:51,162
parameters, R and W. 
And R is the minimum number of nodes that 

119
00:08:51,162 --> 00:08:55,920
need to participate in a successful read. 
Okay and W is the minimum number of nodes 

120
00:08:55,920 --> 00:08:59,480
that are needed to participate in a 
successful write. 

121
00:08:59,480 --> 00:09:04,444
So, this is, sort of, how many replicas 
you write to, and how many replicas need 

122
00:09:04,444 --> 00:09:11,800
to respond from a pool to know that you, 
sort of, have all the information, right? 

123
00:09:11,800 --> 00:09:14,057
Because if everybody is updating 
everything all the time, there might be 

124
00:09:14,057 --> 00:09:16,536
that you have 20 different servers, they 
all have a different version of the data 

125
00:09:16,536 --> 00:09:20,456
you're trying to read. 
And so the question is how many of these 

126
00:09:20,456 --> 00:09:25,145
you need to sample before you feel like 
you have the right one? 

127
00:09:25,145 --> 00:09:29,049
So for a replication factor of N, if R 
plus W, is greater than N, then you can 

128
00:09:29,049 --> 00:09:32,953
claim consistency, but where, often you 
want to set R plus W less than N, in 

129
00:09:32,953 --> 00:09:39,352
order to achieve lower latency. 
So, you don't want to have to actually 

130
00:09:39,352 --> 00:09:43,816
contact, you know too many servers in 
order to satisfy some read request or 

131
00:09:43,816 --> 00:09:49,307
write request, okay. 
So, if you see that notation this is what 

132
00:09:49,307 --> 00:09:52,147
it means. 
But I'm not going to describe too much 

133
00:09:52,147 --> 00:09:55,539
more about it, because I think this sort 
of, you know, the formula falls over a 

134
00:09:55,539 --> 00:09:59,110
little bit under a little bit of 
scrutiny. 

135
00:09:59,110 --> 00:10:03,457
But that's what they're talking about 
when you see this, discussion on say, 

136
00:10:03,457 --> 00:10:05,610
blog posts, okay. 

