BigDataAnalytics
" 49
On
imposesa limit how many nodes can be added to the
misingon scalability. distributed shared disk system, thus compro-
3.12.7 CAP Theorem Explained
is also called the Brewers
The CAP heorem interconnected Theorem. It states that in a distributed
(a collection of nodes that share computing environment
data), it is impossible to provide the
Refer Figure 3.14. At best you can have two of the following three - following guarantees.
one must be sacrificed.
1. Consistency
2. Availability
3. Partition tolerance
[Link] CAP Theorem
Let us spendsome time understanding the earlier mentioned terms.
1. Consistency implies that every read fetches the last write.
2. Availability implies that reads and writes always succeed. In other
words, each non-failing node will
return a response in a reasonable amount of time.
3. Partition tolerance implies that the system will continue to
function when network partition occurs.
Let us try to understand this using a real-life situation.
You work for a training institute, "XYZ. The institute has 50 instructors including you. Allof you
toa training coordinator. At the end of the month, all the instructors together with the training
report
coordina
tor peruse through thetraining requests received from the various corporate houses and prepare a training
schedule for each instructor. These training schedules (one for each instructor) are shared with Amey," the
othce administrator. Each morning, you either call the office helpdesk (essentially Amey's desk) or check
in-person with Amey for your schedule for the day. In case atraining request has been cancelled or updared
(updates can be in theform of change in course, change in duration, change of the training túmings, etc.),
Amey is informed of the updates and the schedules are subsequendly updated by him.
Things were good until now. Few corporate houses were your clients and the schedules of each instructor
couid be smoothly managed wichout any major hiccups. But your training insticute has been implement
ng promotion campaigns toexpand the busines. As a result of advertising in the media and word of
mouth publicity by your existing clients, you suddenly see an upsurge in training requests from existing and
new clients. In consequence of that, more instructors have been recruited. Few trainers/consultants have
also been roped in from other training institutes to help tackle the load.
-
Consistency
CAP Availability
Partition tolerance
Figure 3.14 Brewer's CAP.
50 "
Big Data and Anab
Now when you go to Amey tocheck your schedule or call in at the
in the queue. Looking at the current state of affairs, the trainirg
helpdesk,you are prepared for a
office administrator "Joey." The helpdesk number willremain thecoordinator decides to recruit an
administrators.
same and will be shared by both the o addition
This arrangement works well for a couple of days. Then one
day...
You: Hey Amey!
Amey: Hi! How can Ihelp?
You: Ithink Iam scheduled to anchor a
training at 3:00 pm today. Can Iplease have the details?
Ameyr: Sure! Just a minute.
Amey browses through the file where he
your name at 3:00pm today and maintains the schedules. He does not see a
You: How is that posible? The
responds back, "Youdo not have any training to training scheduled agains
conduct at 3:00 pm."
said he has updated the office training coordinator called up yesterday evening to
administrators of the same. inform of the same and
Amey: Oh! Did he say which office
Amey. Hey Joey! Please check the administrator? It could have been Joey. Please check with Joey.
today? schedule for Paul here... Do yousee
something scheduled at 3:00 pm
Jocy. Sureenough! He is
Aclear case of anchoring the training for client Z" today at 3:00 pm.
with Joey and you inconsistent system!!! The updates in the schedule were
were checking for your schedule with shared by the training coordinator
You share this incident with the Amey.
addressed immediately otherwise it will training coordinator and that gets him
and shares it with both the be difficult to avoid a
chaotic thinking. The issue has to be
office administrators the
following day.
situation. He comes up with a plan
Training Coordinator: Folks, each time that
schedule, make sure that both of you update either an instructor or me
it in your calls any one of you to
get the most recent and respective files. This way the update
speaks to. consistent information irrespective of whom amongst theinstructor willalways
two of you hel sne
Joey: But that could mean adelay in
waiting in queue. answering either a phone call or sharing the
schedule with the instructor
Training Coordinator. Yes, I
Amey: There is this other understand. But there is no way that we can give incorrect
mean that we cannot problem
take any update
as well.
Suppose one of us is on leave on a particular information.
fles (my file and related calls as we will not be able to day. That would
Joey's). simultaneously update both the
Training Coordinator: Well, good
well. Here is the plan: point! Tbats the availability problem! But Ihave
1. If one of you
thought about that
receives the
person if he is updatecall (any updates to any
2. In case the available. schedule), ensure that you inform the other
other pe
Via email. It is a person is not available, ensure that you
3. When the must! inform him of all the updates to all schedules
all other person resumes duty,
schedules that he has the first thing he will do is update his fle with all che updates to
received via email.
Big DataAnalytics " 51
Wow!!! That is sure a Consistent and Available system!
Iooks like everything is in control. Wait a minute! There is a tiff that has taken place berween the ofice
dministrators. The two are pretty much available but are not talking to cach other which, in other words.
meansthat the updates are not flowing from one to the other. We bave to be partition toleran:!! As atrain-
ing coordinator, you instruct them sayingthat none of you are taking any calls requesting for schedules or
updates toschedules ill you patch up. This implies that the system is partition tolerant but not available at
that time.
rhree
In summary, one can at most decide to go with two of the
1. Consistent: The instructors or the training coordinator, once they have updated information with
vOu, will always get the most updated information when they call subsequenty.
2. Availability:The instructors or the training coordinators will always get the schedule if any or both of
che offce administrators have reported to work.
3. Partition Tolerance: Work will go on as usual even if there is communication loss between the office
administrators owing toa spat or a tiff!
When to choose consistency over availability and vice-versa...
1. Choose availability over consistency when your business requirements allow some Aexibility around
when the data in the system synchronizes.
2. Choose consistency over availability when your business requirements demand atomic reads and
writes.
Examples of databases that follow one of the possible three combinations:
1. Availability and Partition Tolerance (AP)
2. Consistency and Partition Tolerance (CP)
3. Consistency and Availability (CA)
Keter Figure 3.15toger aglimpse of databases that adhere to two of the hree characteristics of CAPtheorem.
A Is available/accessible/
operational at all times
AP Riak, Cassandra, CouchDB,
Traditional RDBMS CA
Dynamo like systems
PostgreSQL,
etc.
MySQL, Pick any Two!!
P
CP
System responds incorrectly
Commits are atomic HBase only when there is a total
across the entire MongoDB network failure
distributed systems Redis
MemcacheDB
BigTable like systems
CAP.
Figure 3.15 Databases and
lited
alymoes.
opertions,
3).
fast
scaling vetilalBeswer
deveks) a uoe Cadaing
Scalalkliy
ges yamni staustues. data
allaui
ngBcema,
d No -
fleribliy Schema
ales SeL No
manipulalion data
guehis
fi SQL tablesschemas,
fid use not databaes
do S&L No
(RSBM
s), SaLthaditional Ualike.
tuctd, velumei Vou large
e to
Gitgcy
a S No
SQL) ONLy CNDT NeSeL
.asandaa, Columus Apahe Ex'analia.
r
Key-vaaeD&
lata ato
Store Coumn- ).
optinizad
fait Riat. fot
as data Stee Aloes KeyVa
value -
B,Mongo docemeta
Ex
e D8JCouch
S. ON leiLle in data in data Sfe
OAocemnt-
asel BatabUnieted
Databasee: SaL No
toleene pantiton avadabitiy
an
Cempromiti
ng
n
dheem CAP sppot hoy
a)
Butation.
phopenties AcID Augpot
o (6Ne
ae bigod nolels. Selatinal
datalbales types
Suppati ).es
ne al and Stoage dit Scale
foR Ideal ata- Bata Big 4Suppo
hode
), Graph Riabases- Stae data asnatodky
and elatshps, dea o Setial
and eComuenlation System
DB!
bohy NoSeL? deignad to
Ne SeL databases are
moden dti challen that
hanll itugle wta.
tiaditienl SeL datebases
Realos'.
Ne SQL datal asel
ote BeAres cistead.
&tale out by adds machihe.
DB Casandaa.
Handling large
lange and Aata
anl Biveseunitnctured,
Beituictidand stuctud dta.
Ng SQL allows dynemic changes t
lst atailes wtudit heedig nigatiod.
Fast elwte opeyations
Real time applications
media elonmendalion engis and
4inenial thansaions.
neSaL datalaies ae
dutsitil euionmey
thenwll- ut _
waking
based applic aion.
DB (Aws), Fnusta
Synamo DE
bra_h Batab ases (Neali)T Jeal o
lika boala
Aelatiship heany
nitoks and Qstction.
Value Stais (Reda) Poft t
actng and Behion manageiet.
+Botumt stas (MbngoD6), Best
hested olata
Alvata No SQL Rotabases.
*Can eaiy Stale up and dourn.
+ chentbute.
Relies the diti
data Contitenly Aequineet
pe-defmal Auema.
* Sat ann be epliatid ti mltpl
wodes and Can be patilinel!
. alire dalument
Real-inne based:Rocunet *
Neic. eBayTuillen,
[Link] o
feals seA wl
-dl.
Loalmat menlilian, SelomCros
bases haaph
Nk
data wser wcb
Link shepping anals pais -Vale Key *
Gat,
nSaL No Useo
agplicatiaks
S&L. lasy
omgoDB, But
s the haweCasendra
slandar. have
swp ot not oes *
aofa SQL Cuppott No *
AcIb to
Ne
ohat
we
* Relational
dafined
Pre-odel
* *
mphas
kieariical
S data. Not votically Table Relational
tabae SEL
SaL Falebk
aa endaSQL Ne
best based vasus
en ft Acalalble &hena
ACID daabases
to Ne
phgs S&L
Casanda pradut
theeem. Hoafontaly *odel-s*Buynane
Follous nttuted
datoppoadh.
aNon-elationl, Ne
L.
. atabaselastldd SQ
olenetb
ryaluepaid.
ita'sts. Netllx,
Phaoshap.
Adobe [Link] ,
Brewals
se,: B, kialalble awtan,
CAP eBay
DLTPOLAP
se *Propetus
ctihitd
puinA Relatonml
fonatBata SatShenna
edela (omparisem
ok tadtiondl peanle
Tshansaclion
Praletig
maita
&tandPrline
l ot S&[Link]
data Nes
Tiis Batabase
to modelSaL.
Vesti yes
neo
Cale S yes database.
yes SeL mainlani
al antt
that
SlNe
Baling SaL, modem
Suppoti has
[Link] No ses
N
No RDBMS the
SaL SaL
kada Same
and
elational
yes 'Scale
et. [Link] as s
be Nes Balall.
be Theia allad
SaL
lolossespioeti. asos
distabitas
data tand'
hem tolerait ttat
* itin
phase Ma
* * D.
ule RedulMaep parviales * [Link] Conplens.
dita lnables
and
7iciently
ei- CHaloop
falt patallel distlted
maline!
tple them cato uking
o Dastibitid
[Link] latassts
Splanlattai
[Link]. tolerame VSmdller
aoss Hadof
p'iolaesig PenBoa
slag a
mlgha Fle
slalakle.,
an by
clnks and
glicaig pploah. hanlo to
fiamesik
Syilem)
hodes. falt
.Cot
Efelive;
4) Toleranle
2Fault Simltansu useurles
. [Link]
stuttrd, Scalatlat,by 3.4ARN
nstutid
data.
s Cae, adding
Gmmon mabtiple
Vabus Compared
Costa abess
tiu Met
that Seni- [Link] Can
deledules
and
hermolule Auo nodes the non
ttatl ata Sala handle apacitzo
doop Kso Typei'l Conodily to Haloop
otains pRevent patak Kesaa
nt hanl
Can
wskl oed st.
and hards Replicatd
ihiarie data cata at Ngolal
los.
)
YARN.
Base.A foanewdk
tell Negbata te
mp ltoHDFS
le to A " * Vesions
anto Hadosp.
20. HDFS Bata
are - heu data
3t (YARN) Atola u
ge
el paralleltaika stoey
data
fles
alled
epoatestage HDFS shena
20.
h
wtadap' has fanastk,
les.
Yet
esourle
Hodag
an stoes
datfleas ben ammes
is usten
[Link] ystaint.
Ecos addel.
Buppated managt Contmg
tk. be
ts Redu
Map " x
RedaleuniMapo Bat.
[Link] phootingan
ase Reseule fnion
t+
anl anak
lo)r SZokepel
data Aesuicaauctonealy taniia
"towit AN
ales. data OMaheut: st SI
-acles
btöesn
&tes Pag S&L.
toi data angag
and Thi
sr distibitid
uch Coweilelinto [Link]
S&L be
Hadoep ninigdata umplis
wl as a
&n Binilas
Aelitional a. Scalalle
macine
ApaiHadny
applicate,d. undetand to
[Link] a
allatndata RedueMap Hadoop hat t
[Link] anytne
clutai. Stadad datahti
Hadonp clustad.
Hadoyp ditoibtens
Coe aspeti o Hadop.
O. Hadonp Common.
2). HDEs.
Hadop YARN
stal dittib tion
Apache Haclosps a )
Clonlaay distibitp
inding Apache tadanp
Hadoop EMe Goeanpleg HD)
IBM Sntosphare
MapR Ms tliten)
Sata Seltien)
Hadorp VA SAL
tadop. SaL.
Scale out SCale P
Key Value Paas Rlstod, tlle
Functionl
JEMe
5tefate Hadop ata S
IBM dnfphee
YHP BigBita slation
clond. Based tHadoop Solitis
Aaonweb evia
cland based stitiak