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

2
00:00:05,414 --> 00:00:07,745
Okay, so we start with the same schematic 
that we were looking at when we were 

3
00:00:07,745 --> 00:00:11,472
talking about Map Reduce and Scalability. 
Where we take a big data set and break it 

4
00:00:11,472 --> 00:00:14,735
into chunks, and send those chunks to 
different machines. 

5
00:00:14,735 --> 00:00:18,554
Okay, and here we are replicating this 
chunk to three different machines, which 

6
00:00:18,554 --> 00:00:21,917
is the same thing we did for the Hadoop 
file system, for fault tolerance 

7
00:00:21,917 --> 00:00:25,666
purposes. 
Where you know if this machine dies we 

8
00:00:25,666 --> 00:00:28,994
still have two copies of the data to draw 
from, and we do this with every chunk, 

9
00:00:28,994 --> 00:00:34,060
alright. 
But the two questions, the two 

10
00:00:34,060 --> 00:00:38,088
requirements we need to speak to hear is 
we need to ensure high availability. 

11
00:00:38,088 --> 00:00:41,352
So, that when something goes wrong the 
data still available, and we also want to 

12
00:00:41,352 --> 00:00:44,664
support updates in this context which is 
different than what we were talking about 

13
00:00:44,664 --> 00:00:48,389
before. 
So instead of just read performance or 

14
00:00:48,389 --> 00:00:52,476
fault tolerance in the content of reads, 
we also want to make changes to this data 

15
00:00:52,476 --> 00:00:56,108
now. 
And have them propagate to both other 

16
00:00:56,108 --> 00:01:01,016
replicas, and in some cases to other 
consumers of that change, right. 

17
00:01:01,016 --> 00:01:05,480
There might be other blocks of data that 
that are referred to as the same 

18
00:01:05,480 --> 00:01:11,190
information, I'll give you an example on 
the next slide, okay. 

19
00:01:11,190 --> 00:01:14,850
So, imagine a social networking 
application where people are updating 

20
00:01:14,850 --> 00:01:18,510
their status and their friends get to, 
you friends get to see your status 

21
00:01:18,510 --> 00:01:23,094
updates, okay. 
And so the right operation here is Sue 

22
00:01:23,094 --> 00:01:29,090
updates her own status and the question 
we ask is, of her friends, what happens? 

23
00:01:29,090 --> 00:01:30,730
Who sees the new one, who sees the old 
one? 

24
00:01:30,730 --> 00:01:33,680
How do these, how does this status change 
propagate? 

25
00:01:33,680 --> 00:01:37,320
And the answer to this question from a 
database perspective was, well look you 

26
00:01:37,320 --> 00:01:41,100
know everyone must see the new change or 
no one does. 

27
00:01:41,100 --> 00:01:45,589
Right, either the transaction commits and 
all copies of the data everywhere are 

28
00:01:45,589 --> 00:01:50,646
synchronized simultaneously. 
And further anybody attempting to read 

29
00:01:50,646 --> 00:01:54,678
the value in new media state is able to 
read only the old value, or is, or has to 

30
00:01:54,678 --> 00:01:59,153
wait until the transaction commits, 
right. 

31
00:01:59,153 --> 00:02:03,374
Which could be an arbitrarily, a pretty 
long time, deadlocks can happen which is 

32
00:02:03,374 --> 00:02:07,886
why I said arbitrarily, fine. 
So, that's the answer given by databases, 

33
00:02:07,886 --> 00:02:11,267
everything synchronous, everything must 
be updated, It's either all or nothing, 

34
00:02:11,267 --> 00:02:14,032
okay. 
And the noSQL system just sort of make 

35
00:02:14,032 --> 00:02:17,360
this observation, I said well look for 
really large applications, we simply 

36
00:02:17,360 --> 00:02:21,860
can't afford to wait arbitrarily long for 
this to happen, right? 

37
00:02:21,860 --> 00:02:25,766
I mean, you need status updates to be 
able to commit and respond, so the user 

38
00:02:25,766 --> 00:02:30,882
can go on and do other things, right. 
They can't, sort of, look at a, at an 

39
00:02:30,882 --> 00:02:34,905
hourglass while the synchronization is 
still ocurring. 

40
00:02:34,905 --> 00:02:38,325
You know, and then further the 
observation is, well maybe it doesn't 

41
00:02:38,325 --> 00:02:41,917
matter anyway. 
I mean is it really that important that, 

42
00:02:41,917 --> 00:02:46,057
you know, here, if, if Sue's friend Joe, 
sees the new status while Kai still sees 

43
00:02:46,057 --> 00:02:51,946
the old status, maybe who cares, right? 
As long as Kai eventually sees the new 

44
00:02:51,946 --> 00:02:57,150
status, maybe that's good enough, okay. 
And so these observations suggested 

45
00:02:57,150 --> 00:03:01,770
moving in a different area of the design 
space in sort of high scalability, high 

46
00:03:01,770 --> 00:03:07,973
availability and, consistency, 
application consistency. 

47
00:03:07,973 --> 00:03:12,943
and that motivated and, and those, that 
space of systems started to be associated 

48
00:03:12,943 --> 00:03:17,773
with kind of anti-database, right, took a 
very different approach than databases 

49
00:03:17,773 --> 00:03:22,610
data, and so in turn NoSQL came into 
play. 

50
00:03:22,610 --> 00:03:26,636
It's actually unfortunate that the you 
know, the name that stuck was NoSQL, 

51
00:03:26,636 --> 00:03:30,716
because it doesn't have a whole lot to do 
with SQL. 

52
00:03:30,716 --> 00:03:34,681
It has more to do with the transaction 
processing side of the databases, which 

53
00:03:34,681 --> 00:03:39,134
is not all that relevant for SQL, right? 
I mean the model of transactions, the 

54
00:03:39,134 --> 00:03:41,795
sequences of reads and writes nothing to 
do with the query language. 

55
00:03:41,795 --> 00:03:46,501
But, hey, that's what stuck. 
Now, I, I don't mean to that the term 

56
00:03:46,501 --> 00:03:51,120
NoSQL only suggests these transaction 
models. 

57
00:03:51,120 --> 00:03:53,240
It also sort of suggests a weaker data 
model and so on. 

58
00:03:53,240 --> 00:03:56,548
And we'll talk a little bit about that. 
But I want this point to come across, 

59
00:03:56,548 --> 00:04:00,548
because this is one of the key ideas. 
Okay, how did databases solve this 

60
00:04:00,548 --> 00:04:05,972
problem, or why did they take so long? 
Well there's a, protocol called Two-Phase 

61
00:04:05,972 --> 00:04:11,952
Commit that's fair, fairly standard in 
these situations for synchronous 

62
00:04:11,952 --> 00:04:15,870
processing. 
And so, the motivation for why you want 

63
00:04:15,870 --> 00:04:18,940
two-phase commit goes like this. 
If you want to have a bunch of replicas 

64
00:04:18,940 --> 00:04:22,125
or other kinds of subordinates, anybody 
that wants to see the change, you know 

65
00:04:22,125 --> 00:04:27,007
your, the, the, you make your status 
updated and your friends need to see it. 

66
00:04:27,007 --> 00:04:30,031
The server's holding those different 
friends need to be told of the change, 

67
00:04:30,031 --> 00:04:32,767
and so if you just go ahead and tell 
them, say look I made this change go 

68
00:04:32,767 --> 00:04:37,479
ahead and update your internal state to 
reflect Sue's new status. 

69
00:04:38,680 --> 00:04:41,842
Then you can have you know, some of them 
report back success, but one of them 

70
00:04:41,842 --> 00:04:44,683
could fail. 
But now you're in trouble right, because 

71
00:04:44,683 --> 00:04:47,620
this one has the old value, because it 
failed for some reason. 

72
00:04:47,620 --> 00:04:50,244
Either you didn't hear back from the 
server at all or it said look something's 

73
00:04:50,244 --> 00:04:53,016
wrong with my disk. 
I can't do it, so it responds with a 

74
00:04:53,016 --> 00:04:55,678
failure, regardless. 
And these two have already successfully 

75
00:04:55,678 --> 00:04:58,009
applied the, I'm going to put a check 
mark, have already successfully applied 

76
00:04:58,009 --> 00:05:00,920
the transaction. 
And so now you're in an inconsistent 

77
00:05:00,920 --> 00:05:05,730
state. 
Subordinate 3 has the old value, and 

78
00:05:05,730 --> 00:05:12,609
these guys have the new values, okay? 
So, how you solve this problem, is 

79
00:05:12,609 --> 00:05:16,415
two-phase commit. 
So the first phase here is, the 

80
00:05:16,415 --> 00:05:20,502
coordinator sends a prepare to commit 
message, and the subordinates make sure, 

81
00:05:20,502 --> 00:05:24,772
they take action to make sure that they 
can commit that transaction when asked no 

82
00:05:24,772 --> 00:05:30,414
matter what. 
And so typically this means writing to a 

83
00:05:30,414 --> 00:05:34,860
log the, the, the information related to 
the transaction. 

84
00:05:34,860 --> 00:05:38,764
So, that even if the power goes out, when 
they wake backup they can pull it from 

85
00:05:38,764 --> 00:05:42,010
the log, okay. 
And the subordinates reply with a yes, 

86
00:05:42,010 --> 00:05:44,450
I'm ready to commit. 
And then in phase 2, if all subordinates 

87
00:05:44,450 --> 00:05:47,540
say they're ready, then you'll go ahead 
and send the commit message. 

88
00:05:47,540 --> 00:05:51,635
And if anyone failed, if, if rather 
instead anyone failed then you send back 

89
00:05:51,635 --> 00:05:56,050
an abort message and need to be just 
clean up, okay. 

90
00:05:56,050 --> 00:05:58,875
So this is fine. 
and here's the schematic of it and step 

91
00:05:58,875 --> 00:06:02,587
one, they say prepare, these guys all 
write ahead to the log, and say I'm about 

92
00:06:02,587 --> 00:06:07,120
to write, I'm going to commit this 
transaction. 

93
00:06:07,120 --> 00:06:08,819
They response with yes, I'm ready to do 
so. 

94
00:06:10,400 --> 00:06:14,080
The coordinator comes back with commit, 
and then finally all the work is done. 

95
00:06:14,080 --> 00:06:19,390
And I'm not going to show the schematic 
for what happens in a failure, then 

96
00:06:19,390 --> 00:06:25,060
essentially the coordinator needs to 
watch out for it and send back an abort 

97
00:06:25,060 --> 00:06:35,392
if something had gone wrong, okay. 
Okay, so there's a couple of problems 

98
00:06:35,392 --> 00:06:38,816
with this. 
One is there's some dependencies on the 

99
00:06:38,816 --> 00:06:42,568
coordinator here, that if the coordinator 
fails at the wrong time, things can go 

100
00:06:42,568 --> 00:06:47,685
kind of screwy. 
And a fully distributed protocol for 

101
00:06:47,685 --> 00:06:53,005
ensuring mutual commitment of 
transactions, or other kinds of 

102
00:06:53,005 --> 00:06:59,144
operations can be achieved. 
And one, one of the most successful and 

103
00:06:59,144 --> 00:07:02,430
popular methods of doing this is an 
algorithm called Paxos, that we're not 

104
00:07:02,430 --> 00:07:06,624
going to talk about in detail. 
But you're going to see that term, if you 

105
00:07:06,624 --> 00:07:09,718
look in some of the reading for the NoSQL 
systems, okay. 

106
00:07:09,718 --> 00:07:14,134
So think two-phase commit on a local 
cluster for a database, think Paxos for a 

107
00:07:14,134 --> 00:07:18,680
distributed sort of peer-to-peer kind of 
protocol. 

108
00:07:18,680 --> 00:07:21,464
And just briefly, what Paxos is 
essentially doing is it's a voting 

109
00:07:21,464 --> 00:07:23,998
scheme. 
So, people sort of vote on, you know, the 

110
00:07:23,998 --> 00:07:26,988
individual servers, well to, self 
determine whether or not they're supposed 

111
00:07:26,988 --> 00:07:31,942
to commit the transaction or not. 
And a, the details can get a little bit 

112
00:07:31,942 --> 00:07:37,110
subtle, but overall it's pretty simple 
given the nature, given the difficulty of 

113
00:07:37,110 --> 00:07:42,070
the task involved. 
Okay, so fine, that's one problem. 

114
00:07:42,070 --> 00:07:45,255
The other problem is just, with the Paxos 
sort of shared, is that, this can take 

115
00:07:45,255 --> 00:07:48,391
awhile, right, if subordinates don't 
respond promptly, he might be waiting 

116
00:07:48,391 --> 00:07:51,542
around. 
If things fail multiple times, and you 

117
00:07:51,542 --> 00:07:54,624
just sort of abort and retry transactions 
if the application layer things can go 

118
00:07:54,624 --> 00:07:59,100
slow. 
when there's, it doesn't necessarily 

119
00:07:59,100 --> 00:08:02,870
scale when there's thousands or millions 
of the subordinates are needed to do 

120
00:08:02,870 --> 00:08:08,243
this, you're kind of dead in the water. 
So, other protocols that I'm not going to 

121
00:08:08,243 --> 00:08:11,955
talk about in too much detail include 
multi-version, but you will see in some 

122
00:08:11,955 --> 00:08:16,719
of the papers mentioned, multi-version 
concurrency control. 

123
00:08:16,719 --> 00:08:20,359
Where each write creates a new version of 
the data item, and the legality of the 

124
00:08:20,359 --> 00:08:24,279
read is determined by checking the 
timestamp of the read transaction versus 

125
00:08:24,279 --> 00:08:30,600
the current timestamp of the version that 
you're trying to read, okay? 

126
00:08:30,600 --> 00:08:35,022
And if it's been updated since the time 
you're supposed to be reading it, then 

127
00:08:35,022 --> 00:08:41,140
you know, prior to MVCC, all you could do 
is abort the ter, abort the read. 

128
00:08:41,140 --> 00:08:43,275
And say, look, you, you're looking at 
dirty data, you're done. 

129
00:08:43,275 --> 00:08:46,967
And but with multi-version concurrency 
control, you can actually keep multiple 

130
00:08:46,967 --> 00:08:50,607
versions around, and redirect the read to 
the potentially to the prior version that 

131
00:08:50,607 --> 00:08:55,424
is correct. 
Okay, and thereby avoid avoid aborting a 

132
00:08:55,424 --> 00:09:00,414
certain transactions, fine. 
So, that mechanism still has the 

133
00:09:00,414 --> 00:09:05,990
dependency on a coordinator role to 
administer the time stamp. 

134
00:09:05,990 --> 00:09:11,838
A fully distributive scheme, where the 
decision to go forward of the transaction 

135
00:09:11,838 --> 00:09:20,006
or to avoid a transaction is made. 
The revoting scheme among peers is Paxos. 

136
00:09:20,006 --> 00:09:24,290
And Paxos is very successful and very 
widely applied. 

137
00:09:24,290 --> 00:09:26,666
And you'll see it mentioned in some of 
the noSQL papers if you take the time to 

138
00:09:26,666 --> 00:09:29,150
read them, and they're on the reading 
list. 

139
00:09:29,150 --> 00:09:31,650
And so this relieves the dependency on 
having a central coordinator. 

140
00:09:31,650 --> 00:09:37,500
but is still synchronous, and still has 
the potential for deadlock. 

141
00:09:37,500 --> 00:09:39,878
And, can take some amount of time to 
reach consensus, to be able to know 

142
00:09:39,878 --> 00:09:42,880
what's going on and what kinds of 
failures are happening. 

143
00:09:42,880 --> 00:09:48,130
And so it's difficult to guarantee very 
high performance of very low latency 

144
00:09:48,130 --> 00:09:53,539
response times, alright. 
So, then the term eventual consistency 

145
00:09:53,539 --> 00:09:58,357
was originally defined, not so much in 
the context of its utility, and allowing 

146
00:09:58,357 --> 00:10:04,758
systems to scale to very large levels. 
But just in this argument that the right, 

147
00:10:04,758 --> 00:10:09,222
the only players in distributed systems 
that could make the appropriate decision 

148
00:10:09,222 --> 00:10:14,250
about how to handle conflict where 
applications themselves. 

149
00:10:14,250 --> 00:10:17,727
So, it's a version of this in to in 
argument that you may or may not come 

150
00:10:17,727 --> 00:10:23,430
across in the context of networking. 
And so this was a paper in 1995 by Doug 

151
00:10:23,430 --> 00:10:28,429
Terry, where this term was coined. 
And so he says, you know, we believe that 

152
00:10:28,429 --> 00:10:31,267
applications must be aware that they may 
have read weekly consistent data, and 

153
00:10:31,267 --> 00:10:33,718
that the right operations may conflict 
with those of other users and 

154
00:10:33,718 --> 00:10:37,424
applications. 
And that applications must be involved in 

155
00:10:37,424 --> 00:10:40,414
the detection and resolution conflicts, 
since these naturally depend on the 

156
00:10:40,414 --> 00:10:44,149
semantics of the application. 
And so we'll make the argument in a few 

157
00:10:44,149 --> 00:10:47,185
couple of segments, but I'm not, I'm not 
sure I totally agree with these assert, 

158
00:10:47,185 --> 00:10:50,822
assertions. 
That it actually is better for the system 

159
00:10:50,822 --> 00:10:54,715
to take care of this when it can. 
But what I wanted to do is let you know 

160
00:10:54,715 --> 00:10:57,900
that this is where the term comes from, 
as oppose to the NoSQL system in the last 

161
00:10:57,900 --> 00:11:02,778
10 years or so, which it really would, 
would increase in popularity. 

162
00:11:02,778 --> 00:11:05,518
Okay, so what does it mean? 
Well, what it means is, that in the 

163
00:11:05,518 --> 00:11:08,224
absence of updates, all replicas will 
eventually converge towards identical 

164
00:11:08,224 --> 00:11:11,714
copies, right? 
So, as long as things don't continuously 

165
00:11:11,714 --> 00:11:16,970
change, as changes settle down, we'll all 
eventually see the same value, right? 

166
00:11:16,970 --> 00:11:19,170
All your friends will see your status, 
alright? 

167
00:11:19,170 --> 00:11:23,820
They won't be permanently stuck looking 
at an old one. 

168
00:11:23,820 --> 00:11:27,159
But, you know, what the application sees 
in the meantime, what's one of your 

169
00:11:27,159 --> 00:11:30,657
friends, which, which status one of your 
friends might be looking at, is really 

170
00:11:30,657 --> 00:11:34,208
sensitive to the internal details of 
whatever application you're building and 

171
00:11:34,208 --> 00:11:39,120
is difficult, and therefore is difficult 
to predict. 

172
00:11:39,120 --> 00:11:42,333
Okay, and so for this reason it's, it's a 
little bit difficult to reason very 

173
00:11:42,333 --> 00:11:46,520
precisely or formally about what eventual 
consistency means. 

174
00:11:46,520 --> 00:11:50,532
Because it is so dependent on particular 
limitation details that are themselves 

175
00:11:50,532 --> 00:11:54,334
difficult to formalize, okay. 
And in general, contrast as we've only 

176
00:11:54,334 --> 00:11:56,980
been talking about relational databases 
and things like Paxos, where they 

177
00:11:56,980 --> 00:12:00,160
guarantee strong consistency, but there 
maybe deadlocks. 

178
00:12:00,160 --> 00:12:05,035
And so it's, you can prove that no system 
can be free of deadlocks and guarantee 

179
00:12:05,035 --> 00:12:09,158
consistency. 
And so, relational databases like Paxos 

180
00:12:09,158 --> 00:12:12,692
give up on this liveness property, 
meaning that they'll, they might allow 

181
00:12:12,692 --> 00:12:18,892
deadlocks in favor of strong consistency. 
Now, they've, you can show that the cases 

182
00:12:18,892 --> 00:12:23,637
where deadlocks can occur, can be made 
sort of rare, through different design 

183
00:12:23,637 --> 00:12:30,238
decisions, but they can still happen. 
Okay, fine, so visually consistent models 

184
00:12:30,238 --> 00:12:34,459
say, we can't afford the cost of waiting 
for these protocols to run and more over, 

185
00:12:34,459 --> 00:12:40,195
they might not be necessary in certain 
application context. 

186
00:12:40,195 --> 00:12:43,661
Alright, so where we are now is we're 
looking at this column. 

187
00:12:43,661 --> 00:12:47,315
And I've already sort of marked this up a 
little bit, but what these words now 

188
00:12:47,315 --> 00:12:51,027
mean, and we'll talk about these a little 
more when we talk about a few of these 

189
00:12:51,027 --> 00:12:57,487
systems, is the scope of where strongly 
consistent transactions are supported. 

190
00:12:57,487 --> 00:13:01,105
And so the scope here of a single record, 
means that I can update a multiple fields 

191
00:13:01,105 --> 00:13:04,669
in one record, and either all the changes 
will occur, or none of them will occur, 

192
00:13:04,669 --> 00:13:08,984
okay? 
by the way I filtered this list only 

193
00:13:08,984 --> 00:13:13,473
include noSQL systems, so relational 
databases support this across arbitrary 

194
00:13:13,473 --> 00:13:17,855
records, right. 
You can have, you could update a record 

195
00:13:17,855 --> 00:13:21,885
over here and update a record over there 
and call that one transaction, and the 

196
00:13:21,885 --> 00:13:27,890
system will only see both of those 
changes or neither of those changes. 

197
00:13:27,890 --> 00:13:31,952
And that's what's not supported with 
these NoSQL systems. 

198
00:13:31,952 --> 00:13:34,712
So, within interview record is supported, 
within some of these systems nothing is 

199
00:13:34,712 --> 00:13:37,618
supported. 
you can't, there's no guarantees at all 

200
00:13:37,618 --> 00:13:40,034
really. 
And what this EC means is eventually 

201
00:13:40,034 --> 00:13:43,479
consistent, so its not really strongly 
consistent its not a transaction, but 

202
00:13:43,479 --> 00:13:48,029
they do have eventually consisting 
guarantees at the record level. 

203
00:13:48,029 --> 00:13:52,307
And that's what all these systems sort of 
guarantee, and then this system Megastore 

204
00:13:52,307 --> 00:13:56,460
that's based on big table from Google. 
It's also a Google system. 

205
00:13:56,460 --> 00:14:01,212
defines the notion of entity groups, and 
this is a set of related records for 

206
00:14:01,212 --> 00:14:07,060
which transactions are strongly 
consistent for that group, okay. 

207
00:14:07,060 --> 00:14:09,985
So, this is a little bit better than just 
one individual record as the one, only 

208
00:14:09,985 --> 00:14:13,000
guarantee we can give, and it's a little 
bit less than any arbitrary record in the 

209
00:14:13,000 --> 00:14:17,994
database. 
It's predefined into the groups that 

210
00:14:17,994 --> 00:14:22,260
allow transactions. 
Okay, so this is sort of a compromise, 

211
00:14:22,260 --> 00:14:25,469
fine. 
And then this most recent system form 

212
00:14:25,469 --> 00:14:30,780
Google Spanner offers true strong 
consistency across all the records. 

213
00:14:30,780 --> 00:14:38,980
And we'll talk about why they made that 
choice, in a little bit. 

214
00:14:38,980 --> 00:14:40,560
Okay. 
So, another concept I want you to be 

215
00:14:40,560 --> 00:14:43,660
familiar with is this so-called CAP 
theorem, from Eric Brewer in 2000, 

216
00:14:43,660 --> 00:14:48,174
followed up later by Lynch in 2002. 
Where they define these three notions, 

217
00:14:48,174 --> 00:14:50,910
consistency, availability and 
partitioning. 

218
00:14:50,910 --> 00:14:55,010
And the way this is often described is 
you have to choose two of these. 

219
00:14:55,010 --> 00:14:58,727
You know you can't get all three, you 
have to choose two or sacrifice per, 

220
00:14:58,727 --> 00:15:01,627
performance. 
But I don't really like thinking of it 

221
00:15:01,627 --> 00:15:03,546
that way. 
And Eric Brewer has also sort of 

222
00:15:03,546 --> 00:15:06,820
described that maybe that's not the right 
way to think about it. 

223
00:15:06,820 --> 00:15:09,724
And the reason is because it's not clear 
what it means to choose consistency and 

224
00:15:09,724 --> 00:15:12,340
availability at the extent of 
partitioning. 

225
00:15:12,340 --> 00:15:15,418
Okay, so what is partitioning? 
Partitioning means well, if you got a big 

226
00:15:15,418 --> 00:15:18,567
distributive system with a hundreds of 
nodes involved, hundreds of servers all 

227
00:15:18,567 --> 00:15:21,763
communicating with each other, and some 
segment of them lose communication with 

228
00:15:21,763 --> 00:15:26,429
the other servers. 
Can those two segments still make forward 

229
00:15:26,429 --> 00:15:30,610
progress in the application independently 
and sync up later? 

230
00:15:30,610 --> 00:15:34,702
Or does everything have to stop and wait, 
or certain nodes have to stop and wait, 

231
00:15:34,702 --> 00:15:38,330
in order to re, reestablish 
communication? 

232
00:15:38,330 --> 00:15:41,354
So, for example, if you have a master 
node that controls everything, and you 

233
00:15:41,354 --> 00:15:44,799
have some worker nodes that lose contact 
with the master. 

234
00:15:45,850 --> 00:15:49,165
There are many designs at which you can't 
make any forward progress, until you 

235
00:15:49,165 --> 00:15:52,220
reestablish connections with the master 
node. 

236
00:15:52,220 --> 00:15:54,236
Right, you can, you can do no useful 
work, because you're waiting on 

237
00:15:54,236 --> 00:15:56,648
communication, you're waiting on that 
last message from the, from the master to 

238
00:15:56,648 --> 00:16:00,875
tell you what to do next. 
so in those cases you've given up 

239
00:16:00,875 --> 00:16:06,142
availability, right, you go down. 
Those nodes are no longer accessible or, 

240
00:16:06,142 --> 00:16:11,224
or, or doing useful work in the context 
of a network partitioning, okay. 

241
00:16:11,224 --> 00:16:15,316
On the other hand, if you do say, well 
sure we're going to continue to do useful 

242
00:16:15,316 --> 00:16:19,870
work even independently, then it's not 
difficult to show that you can arrive at 

243
00:16:19,870 --> 00:16:25,871
an inconsistent state, right? 
Updates are coming in to this partition, 

244
00:16:25,871 --> 00:16:29,700
and updates are coming to this partition 
of the network. 

245
00:16:29,700 --> 00:16:32,430
And sometime down the road communications 
are reestablished. 

246
00:16:32,430 --> 00:16:35,716
And you find out, oops, you know, your 
replica has one value, my replica has 

247
00:16:35,716 --> 00:16:38,190
another. 
Which one's right? 

248
00:16:38,190 --> 00:16:40,206
Well, we're going to have to sort it out, 
but meanwhile we've already sort of 

249
00:16:40,206 --> 00:16:44,002
exposed these values to the application. 
So in some sense we're, we're 

250
00:16:44,002 --> 00:16:48,688
demonstrably inconsistent, okay. 
It's the point is you can't get all three 

251
00:16:48,688 --> 00:16:51,810
of these. 
All right. 

252
00:16:51,810 --> 00:16:54,190
So, you either sacrifice availability or 
you sacrifice consistency by allowing 

253
00:16:54,190 --> 00:16:57,365
things to continue working. 
So, conventional databases essentially 

254
00:16:57,365 --> 00:17:00,805
assume that there is no partitioning. 
And again, this is a function of them 

255
00:17:00,805 --> 00:17:03,570
only operating on tens of nodes at a 
time, all right? 

256
00:17:03,570 --> 00:17:11,506
They didn't go to this thousand node 
scale, or planet wide distributed 

257
00:17:11,506 --> 00:17:15,673
systems. 
and so you can kind of assume that there 

258
00:17:15,673 --> 00:17:18,280
wasn't, there wasn't really a need to 
worry about. 

259
00:17:18,280 --> 00:17:20,926
Well, what if queries are coming in to 
half my nodes and they can't talk to the 

260
00:17:20,926 --> 00:17:24,442
other half of my nodes, and so on. 
They're all sitting there in a cluster 

261
00:17:24,442 --> 00:17:27,140
that's in your data center. 
Not even in your data center, in your 

262
00:17:27,140 --> 00:17:29,720
server room to some extent. 
And so that wasn't, really wasn't an 

263
00:17:29,720 --> 00:17:31,830
issue that they were thinking about too 
much, okay. 

264
00:17:31,830 --> 00:17:35,210
And the NoSQL systems do need to worry 
about this. 

265
00:17:35,210 --> 00:17:36,850
They are very large, they are very 
distributed. 

266
00:17:36,850 --> 00:17:39,434
There are different kinds of Byzantine 
failures happening all the time, because 

267
00:17:39,434 --> 00:17:43,976
of, because of the sheer scale. 
And therefore, they choose to sacrifice 

268
00:17:43,976 --> 00:17:48,040
consistency instead of availability, 
okay. 

269
00:17:48,040 --> 00:17:52,194
And so, graphically you can look at this 
in this sort of triangle forming, you can 

270
00:17:52,194 --> 00:17:56,348
put different systems on sort of an edge 
here where relational databases assume 

271
00:17:56,348 --> 00:18:03,117
consistency and availability. 
but, but assume partitioning can never 

272
00:18:03,117 --> 00:18:05,665
happen. 
While other systems need to tolerate 

273
00:18:05,665 --> 00:18:07,880
partitioning, but give up on 
availability. 

274
00:18:07,880 --> 00:18:15,062
and they are ensuring that certain kinds 
of transactions are going to be 

275
00:18:15,062 --> 00:18:20,764
consistent, okay. 
And then other systems say, well, we're 

276
00:18:20,764 --> 00:18:26,920
going to give up on consistency. 
but you can always use [UNKNOWN] work, 

277
00:18:26,920 --> 00:18:29,724
okay, so fine. 
And really this is, the important thing 

278
00:18:29,724 --> 00:18:33,090
here is the scope of the variation I put 
in my table's sort of critical here too. 

279
00:18:33,090 --> 00:18:39,362
It's not sort of, nothing except for 
spanner over here, even tries to provide 

280
00:18:39,362 --> 00:18:46,282
global transactions like relational 
databases do, okay. 

281
00:18:46,282 --> 00:18:53,211
Alright. 
So fine, so I'll pick up here in the next 

282
00:18:53,211 --> 00:18:56,286
segment. 

