Cloud Programming and Software
Environments
6.1 Features of Cloud and Grid Platforms
C lou d a n d G rid sy ste m s off e r v a riou s fe a tu re s to su p p o rt m o d e rn
co m p u tin g n e ed s. T h es e f ea tu re s f a ll in to m u ltip le c a te g o rie s:
ca p a b ilitie s, tra d itio n a l f ea tu re s, d a ta fe a tu res , and
p rog ra m m in g /ru n tim e su p p o rt. U n d e rsta n d in g th e se h e lp s in u s in g
clou d p la tfo rm s ef fe c tiv e ly .
6.1.1 Cloud Capabilities and Platform Features
• E la stic U tility C om p u tin g : C lo u d s c a n a u to m a tic a lly sc a le res ou rc es
(u p o r d ow n ) d e p e n d in g o n w o rklo a d d e m a n d s, of fe rin g f le x ib ility
a n d co st sa v in g s.
• P la tfo rm a s a S e rv ic e (P a a S ): P la tfo rm s like A z u re a n d A m a z on W eb
Se rv ice s (A W S) p rov id e m o re th a n in fra s tru c tu re:
• A z u re : O ffe rs T a b le s, Q u e u e s, B lo b s , S Q L D B , W e b & W o rke r
roles .
• A W S : O f fe rs S im p leD B , N o tif ic a tio n , M o n ito rin g , C D N , R D S
(R e la tio n a l D B ) , a n d H a d o op (M a p R e d u c e ).
• G oog le A p p E n g in e ( G A E ): P ro v id e s a stron g e n v iro n m e n t
fo r b u ild in g w e b a p p lic a tio n s.
Infrastructure Features (Table 6.2): Emphasis Features (Table 6.4):
• Includes core computing resources like VMs, storage, and • Newer features, such as data parallel file systems (e.g., Sector),
networking. are being widely adopted.
• These are typically managed by the cloud provider and • Academic clouds (e.g., Eucalyptus, Nimbus) may lack some
form the base for other services. advanced features in commercial clouds.
Programming Environments (Table 6.3):
• Includes support for traditional parallel and distributed programming models.
• Developers can use these for building scalable cloud applications.
• Examples include MPI, MapReduce, and service-oriented architecture (SOA).
6.1.2 Traditional Features in Grids and Clouds
T h e se a re es se n tia l f ea tu re s th a t s u p p ort w o rkf lo w m a n a g em e n t, d a ta m ov e m e n t, se c u rity , a n d sy ste m re lia b ility.
[Link] Workflow
• W o rkflow D e fin ition : A w orkf lo w c on n e c ts m u ltip le ta sks o r se rv ice s in to a c om p le te p ro ce ss .
• P o p u la r S y ste m s:
• P e g a su s, T a v e rn a , a n d K e p le r: O p e n -sou rc e to ols f or sc ie n tif ic w o rkflow s.
• T rid en t (b y M ic ros oft) : B u ilt on W in d o w s W ork flo w F o u n d a tion a n d in te g ra tes w e ll w ith A z u re a n d o th er p la tfo rm s.
• U s e C a se : W o rkf low s a llo w u se rs to lin k clou d a n d n o n -clou d s erv ic es e ffic ie n tly , e sp e c ia lly in s cie n tific o r d a ta -d riv e n a p p lica tion s .
[Link] Data Transport
• D a ta In g re ss /E g re ss C os ts: T ra n sf errin g la rg e d a ta se ts in to th e clou d is ofte n ex p en s iv e a n d slo w .
• A c a d em ic In te g ra tion : C lou d se rv ice s m a y so on in te g ra te b e tter w ith sy ste m s like T e ra G rid u sin g h ig h -sp e e d lin ks.
• D a ta S tru c tu res : C lo u d p la tform s o fte n s tore d a ta in b lo ck s ( A z u re b lob s) a n d ta b le s, w h ic h ca n b e op tim iz e d u sin g p a ra llel p roc es sin g .
• C u rre n t M eth od s : M os t d a ta is m o v e d u sin g H T T P p rotoc o ls, e sp ec ia lly in a c a d e m ic or p roto typ e e n v iro n m en ts.
[Link] Security, Privacy, and Availability
Th es e a re c ritica l f or tru st a n d d e p e n d a b ility in c lou d p la tfo rm s. K e y stra te g ie s in c lu d e:
• Virtu a l C lu s terin g : D y n a m ic a lly a llo ca te res ou rc e s to re d u ce ov e rh ea d a n d im p ro v e e ffic ie n c y.
• P e rsiste n t D a ta S tora g e: U s e relia b le stora g e sy ste m s th a t s u p p ort fa st q u e ry in g a n d d a ta re trie v a l.
• U se r A u th e n tic a tion A P Is: C lou d p la tf orm s p ro v id e se c u re A P Is fo r ta s ks like log g in g in a n d em a il se rv ic e s.
• Se c u re A c c es s P roto co ls : U se H T T P S , S S L , e tc ., to e n c ryp t d a ta in tra n sit a n d p ro tec t u se r p riv a c y.
• F in e -G ra in e d A cc e ss C on tro l: A s sig n s p e c ific p erm ission s to c on tro l w h o ca n a c ce s s o r m o d ify d a ta .
• P ro tec tion of S h a red D a ta : P re v en t u n a u th o riz e d m od if ic a tio n o r d ele tion of d a ta ; p rote ct in te lle ctu a l p rop e rty.
• D isa ste r R e c ov e ry: E n a b le V M m ig ra tion a n d oth e r m e c h a n ism s to m a in ta in a va ila b ility d u rin g fa ilu re s.
• R e p u ta tion -B a se d Se c u rity : T ru st m od els a re u se d to a llow on ly v erified c lie n ts, p re v e n tin g a c c es s from m a lic io u s u se rs o r p ira te s.
6.1.3: Data Features and Databases
[Link] Program Library
• A co llec tion of v irtu a l m a ch in e (V M ) im a g e s u se d in clou d c om p u tin g .
• T h e se V M im a g e lib ra ries sim p lif y th e p ro c es s o f d e p loy in g a n d m a n a g in g a p p lica tion s o n c lo u d p la tfo rm s.
• C lou d p la tfo rm s (e sp ec ia lly Ia a S) in clu d e too ls to ea sily d e p loy , c on fig u re , a n d m a n a g e V M im a g e s, m a k in g it u se r-frien d ly for b o th a c a d em ic a n d
c om m e rc ia l u s ers .
[Link] Blobs and Drives
• B a sic sto ra g e u n its in clou d p la tf orm s . E x a m p le s: A z u re B lob S tora g e, A m a z o n S 3 .
• T h e se a re a tta ch e d sto ra g e (like d is ks) fo r ru n n in g v irtu a l m a ch in e s— A z u re D riv e s a n d A m a z on E la stic B loc k S tore ( E B S ).
• A z u re u s es to g rou p b lo b s , s im ila r to fold e rs/ d ire cto rie s.
• C lou d sto ra g e is fa u lt-tolera n t b y d ef a u lt (a u to m a tic a lly b a c ke d u p a n d res ilie n t), w h ile tra d itio n a l sy ste m s like T e ra G rid re q u ire m a n u a l b a ck u p s .
• A P Is lik e th e S im p le Clou d F ile Sto ra g e A P I h elp in sta n d a rd iz in g b lob a n d f ile a c ce s s.
[Link] DPFS (Data-Parallel File Systems)
• S p e c ia l f ile sy ste m s b u ilt to h a n d le la rg e -s ca le d a ta w ith o p tim iz e d co m p u te -d a ta a ff in ity (i.e ., ke e p in g d a ta clos e to w h ere it’ s p roc e sse d ).
• G oo g le F ile S y ste m , H D F S (H a d oo p ), C o sm o s ( u se d in D rya d ).
• Id e a l f or d a ta -in te n s iv e a p p lic a tio n s.
• T h e se s ys tem s m a y n o t d ire c tly in teg ra te w ith b lob o r d riv e a rc h ite c tu re s, req u irin g d a ta m ov e m e n t b e tw e e n f orm a ts .
• C lou d s lik e A m a z on a n d A z u re d on ’ t y e t fu lly s u p p ort co m p u te -d a ta a ff in ity, b u t f ea tu re s like A z u re Af fin ity G rou p s a re a ste p f orw a rd .
[Link] SQL and Relational Databases
• B oth A m a z on a n d A z u re p ro v id e S Q L -b a se d rela tio n a l d a ta b a se s.
• E a sy to im p le m e n t in sm a lle r-s c a le e n v iro n m e n ts u n le ss e x trem e ly la rg e-s ca le d a ta is in v olv e d .
• O M O P p ro jec t u sin g O ra c le , S A S, a n d H a d oo p on F u tu re G rid .
• R a th e r th a n b u n d lin g d a ta b a se s in to a p p im a g e s, c lo u d s n o w d e p loy th e m a s se p a ra te se rv ic e s— ca lle d “ S Q L a s a S e rv ic e ” (like w ork er roles in
A z u re ).
• S im p lifies m a n a g e m en t— o n ly re q u ire s a s m a n y se rv ic e s a s f ea tu re s (N fe a tu re s → N s erv ice s) , in ste a d o f c o m p lex co m b in a tio n s (2 ^N im a g es ).
[Link] Table and NOSQL Nonrelational Databases
• NOSQL Overview: T h e se a re f le x ib le , sc a la b le d a ta b a se s w ith ou t rig id s ch e m a req u ire m e n ts.
• Examples: B ig Ta b le , S im p leD B , A z u re T a b le
• Characteristics:
• D e sig n e d fo r s ca la b ility a n d d istrib u tion
• S c h e m a -fre e (e a c h re c ord c a n h a v e d iffe re n t p rop ertie s)
• S o m e u se co lu m n fa m ilie s ( e.g ., B ig T a b le)
• Scientific Relevance:
• U s ed f or m a n a g in g la rg e sc ien tific d a ta s ets
• V O T a b le in a stro n om y a n d to ols like E x c e l sh o w ta b le u tility
• Academic Alternatives:
• A p a c h e H B a se ( lik e B ig Ta b le )
• C o u c h D B (d oc u m e n t sto re )
• M /D B (S im p leD B -like )
• Integration tools: S im p le C lo u d AP Is of fe r sta n d a rd a c ce ss m e th o d s fo r file , ta b le , a n d d oc u m e n t stora g e a c ros s c lou d p la tfo rm s.
[Link] Queuing Services
• Purpose: A llow d iffe re n t p a rts o f a n a p p lic a tion to c o m m u n ic a te v ia m e ss a g e p a ss in g .
• Features:
• Sh ort m es sa g e s (u n d e r 8 K B )
• R E S T fu l A P I ( w e b -b a se d in te rfa ce )
• “ D e liv e r a t le a st on c e ” relia b ility
• M es sa g e tim e ou ts c o n trol h ow lo n g a m e ss a g e is h e ld fo r p roc es sin g
• Examples: A m a z on S Q S , A z u re Q u e u es
• and a re op e n -s ou rc e s yste m s w ith sim ila r f u n ction a lity
• Q u eu e s e n a b le lo os e c ou p lin g b e tw e e n c o m p o n en ts a n d s u p p ort s c a la b le , relia b le a p p lic a tion d e sig n .
6.1.4 Programming and Runtime Support
T h is se ction fo c u se s on too ls a n d m od els th a t en a b le p a ra lle l c om p u tin g a n d
ru n tim e m a n a g em e n t in c lo u d en v iron m e n ts , e sp ec ia lly th ro u g h sy ste m s lik e
M ap Re duc e .
[Link] Worker and Web Roles
• W o rke r R o le s: T h es e a re b a c kg ro u n d p ro ce ss es in A z u re th a t a u tom a tica lly ru n
ta sk s. Th e y d o n ’ t n e e d m a n u a l sc h ed u lin g b ec a u s e th e c lou d h a n d le s it. T h e y
a re g re a t fo r p ro c es sin g ta s ks u sin g q u e u e s.
• W e b R oles : T h e se a re u s ed to b u ild a n d d e p loy w e b a p p lic a tion s. T h e y c a n
se rv e a s p o rta ls fo r u s er in te ra c tio n .
• Q u eu es : Q u e u e s a re u s ed to a ssig n ta sks in a f a u lt-tole ra n t a n d d istrib u te d
m a n n e r, h elp in g m a n a g e w orkloa d s ef ficien tly .
• G oo g le A p p E n g in e (G A E ) ta rg ets g e n e ra l w e b a p p s , w h e re a s sc ie n c e g a te w a y s
(u s ed in T e ra G rid ) c a te r to sc ien tific a p p lic a tio n s.
[Link] MapReduce
• M a p R ed u c e: A p o p u la r m od e l fo r p ro ce ss in g la rg e-s ca le d a ta in p a ra lle l. It
b re a ks ta sks in to tw o sta g e s: M a p ( p ro ce ss in d iv id u a l p ie c es o f d a ta ) a n d
R e d u c e ( co m b in e re su lts ).
• U n like old e r g rid sy ste m s, M a p R e d u c e s u p p orts : D yn a m ic e x e c u tio n - A d j u sts
ta sk s d u rin g ru n -tim e, H ig h fa u lt tolera n c e - R e co v ers from f a ilu res , a n d ea s e o f
u se .
• P op u la r Im p le m e n ta tio n s a re :
• H a d o op : O p en -s ou rc e, w id e ly a d op ted ( of fe red b y A m a z on ) .
• Dry a d : M ic ros oft’ s a lte rn a tiv e (e x p e c ted in A z u re ).
• Tw iste r: G o od fo r ite ra tiv e ta s ks like d a ta m in in g o r m a trix o p e ra tio n s.
• C lo u d e ra : A p la tf orm th a t of fe rs v a riou s H a d oo p d istrib u tion s , h e lp fu l f or
ru n n in g M a p R e d u ce o n d if fe ren t e n v iro n m e n ts, in c lu d in g Am a z on a n d Lin u x .
• R e a l-W orld U se : Sy ste m s like F u tu re G rid p la n to su p p o rt H a d oop , D rya d , a n d
T w iste r f or res ea rc h a n d tea ch in g .
[Link] Cloud Programming Models
• T h e se m od els d e fin e h ow p ro g ra m s a re w ritte n a n d ru n in c lou d e n v iro n m e n ts.
• E x a m p le s:
• G oo g le A p p E n g in e (G A E ) a n d A n e ka (b y M a n jra so ft): P rov id e fra m e w o rks a n d to ols f or b u ild in g c lo u d a p p lica tion s .
• Ite ra tiv e M a p R e d u c e : A n e x ten sio n o f M a p R e d u ce th a t su p p o rts rep e a ted d a ta p roc e ssin g (u se f u l in m a ch in e lea rn in g , c lu ste rin g ).
• P o rta b ility : T h e se m o d e ls h e lp m o v e a p p lic a tio n s a cro ss c lou d , h ig h -p e rfo rm a n c e co m p u tin g (H P C ), a n d c lu s ter sy ste m s.
[Link] SaaS (Software as a Service)
• S a a S C o n ce p t: U s ers a c c es s sof tw a re h os ted o n th e c lou d in ste a d of in sta llin g it loc a lly .
• B e n e fits:
• E a sie r d e p loy m en t (e .g ., S Q L a s a se rv ice ).
• S c a la b le a n d m o re m a n a g e a b le.
• T ec h n ica l T o ols N ee d e d :
• C lo u d se rv ic e s lik e M a p R ed u c e, B ig Ta b le, E C2 , S 3 , H a d o op , G A E , a n d W e b S p h ere .
• P ro te ction F ea tu res : Im p orta n t fo r b u ild in g tru st a n d relia b ility in S a a S :
• S c a la b ility
• S e cu rity
• P riv a c y
• A v a ila b ility
6.2 PARALLEL AND DISTRIBUTED PROGRAMMING
PARADIGMS
• P a ra lle l a n d d is t rib u t e d p ro g ra m m in g in v o lv e s ru n n i n g a p a ra lle l p ro g ra m o n a d i st rib u t e d c o m p u t in g s y st e m , c o m b in i n g t h e
s tre n g t h s o f b o t h p a ra d i g m s .
• A d i st rib u t e d c o m p u t in g s y st e m c o n s is ts o f m u lt ip le in d e p e n d e n t c om p u t e rs c o n n e c t e d v ia a n e t w o rk, w o rkin g to g e t h e r to
a c h ie v e a c om m o n g o a l, s u c h a s ru n n in g a n a p p lic a t io n .
• P a ra lle l c o m p u tin g m e a n s p e rfo rm in g m a n y o p e ra t io n s s im u lt a n e o u s ly u sin g m u lt ip le p ro c e s so rs , e it h e r w it h in a s in g le m a c h in e
o r a c ro s s m u lt ip le s ys t e m s .
• W h e n p a ra l le l p ro g ra m s ru n o n d is t ri b u t e d sy s t e m s , t h e re s u lt is fa s t e r e x e c u ti o n a n d b e tt e r p e rfo rm a n c e t h ro u g h sh a re d
c o m p u t a t io n a l e f fo rt .
• T h is se tu p o ffe rs b e n e fit s lik e re d u c e d re s p o n se t im e fo r u se rs a n d h ig h e r t h ro u g h p u t a n d re so u rc e u ti liz a tio n fo r s y st e m s.
• H o w e v e r, ru n n i n g p a ra lle l p ro g ra m s o n d is t ri b u t e d s y st e m s is c o m p le x , re q u irin g e ffi c ie n t c o o rd in a t io n , s y n c h ro n iz a t io n , a n d
d a t a c o m m u n ic a t io n a m o n g v a rio u s n o d e s .
• T y p ic a lly , t a s k s a re d iv id e d , p ro c e s s e d in p a ra ll e l a c ro s s d iffe re n t c o m p u te rs , a n d t h e n t h e re su lt s a re c o m b i n e d , w h i c h d e m a n d s
c a re fu l d e sig n a n d m a n a g e m e n t .
• D e sp it e t h e c o m p l e x it y , t h i s p a ra d ig m is w id e ly u s e d in h ig h -p e rfo rm a n c e c o m p u t in g , b ig d a t a , a n d re a l-tim e a p p lic a t io n s d u e to
it s e ffic ie n c y a n d s c a la b il ity .
6.2.1 Parallel Computing and Programming Paradigms
• In parallel computing within a distributed system, the first step is partitioning, which includes dividing both the computation and the
data.
• Computation partitioning breaks a program into smaller tasks that can run simultaneously.
• Data partitioning splits input or intermediate data into chunks for separate processing.
• mapping, where these tasks or data pieces are assigned to different networked nodes or workers. This helps ensure that all available
resources are used efficiently.
• Synchronization is crucial because different workers might depend on shared resources or on the output of others. It ensures that tasks
are coordinated and that issues like race conditions and data dependency are managed properly.
• Communication between workers is needed when one task depends on the result of another. This happens when intermediate data
must be transferred among workers during processing.
• Finally, scheduling handles the order in which tasks or data pieces are assigned, especially when there are more tasks than available
workers. A scheduler selects which task goes next, while a resource allocator maps them to specific workers. Scheduling also applies
when multiple programs are competing for limited system resources.
[Link] Motivation for Programming Paradigms,
• H a n d lin g t h e e n t ire d a t a flo w in p a ra lle l a n d d is t rib u t e d p ro g ra m m in g is c o m p le x , t im e -c o n s u m i n g , a n d re q u i re s d e e p t e c h n ic a l
k n o w le d g e , w h i c h c a n re d u c e p ro g ra m m e r p ro d u c t iv i ty a n d d e l a y p ro d u c t d e li v e ry.
• T o a d d re ss t h i s, p ro g ra m m in g p a ra d ig m s o r m o d e ls a re in t ro d u c e d to a b st ra c t th e lo w e r-l e v e l im p l e m e n ta t io n d e t a ils ,
a llo w i n g p ro g ra m m e rs to fo c u s o n t h e c o re lo g ic ra t h e r t h a n in fra s t ru c t u re .
• T h e s e m o d e ls m a k e it e a sie r t o w rit e p a ra lle l p ro g ra m s b y h id in g c o m p le x o p e ra t io n s , im p ro v i n g s im p lic it y a n d e ffic ie n c y .
• T h e m a in m o t iv a t io n s b e h in d u s in g s u c h p a ra d ig m s in c lu d e im p ro v in g p ro g ra m m e r p ro d u c ti v it y , re d u c in g t im e t o m a rk e t ,
u sin g re so u rc e s m o re e ff ic ie n tl y, in c re a sin g s ys t e m t h ro u g h p u t , a n d e n a b lin g h ig h e r a b s t ra c t io n le v e l s.
• P o p u la r m o d e rn m o d e l s lik e M a p R e d u c e , H a d o o p , a n d D ry a d w e re o rig in a lly d e sig n e d fo r in fo rm a t io n re t ri e v a l b u t a re n o w
w id e ly u se d in m a n y a p p lic a t io n s .
• T h e s e m o d e ls a lso o ffe r a d v a n ta g e s li ke lo o s e c o u p lin g , m a ki n g th e m id e a l f o r v irt u a l m a c h in e (V M ) e n v iro n m e n ts a n d
p ro v id in g b e t te r f a u lt t o le ra n c e a n d s c a la b ilit y c o m p a re d t o t ra d it io n a l m o d e ls l ike M P I.
6.2.2 MapReduce, Twister, and Iterative MapReduce
• MapReduce is a so f tw a re fra m e w o rk t h a t s im p lif ie s p a ra ll e l a n d d is t ri b u t e d c o m p u t in g o n la rg e d a t a s e t s . It h id e s t h e c o m p le x
d a t a flo w b y p ro v id in g u s e rs w i th t w o m a in f u n c t io n s : Map a n d Reduce.
• U se rs d e fin e h o w d a t a s h o u ld b e p ro c e s se d b y c u s t om iz in g t h e s e t w o fu n c t io n s, w h i c h a llo w s th e m t o c o n t ro l th e
c o m p u t a t io n lo g ic w it h o u t m a n a g in g t h e u n d e rly in g i n fra st ru c tu re .
• In t h is fra m e w o rk, d a t a is p ro c e s se d a s (key, value) p a irs . T h e "value" is t h e a c tu a l d a ta b e in g h a n d le d , w h ile t h e "key" h e lp s
t h e M a p R e d u c e s y st e m m a n a g e h o w d a t a is g ro u p e d a n d p a ss e d b e t w e e n t h e M a p a n d R e d u c e s t a g e s .
• T h e M a p fu n c t io n p ro c e s se s in p u t d a ta in t o in t e rm e d i a t e (k e y , v a lu e ) p a irs , w h ic h a re th e n g ro u p e d b y k e y a n d p a s s e d t o th e
R e d u c e fu n c t io n to g e n e ra t e th e fin a l o u t p u t .
• T h is a b st ra c t io n a ll ow s u s e rs t o w rit e sc a la b le p a ra lle l p ro g ra m s w it h o u t w o rry in g a b o u t d e t a il s lik e d a t a d is trib u t io n , fa u lt
t o le ra n c e , o r c o m m u n ic a t io n b e t w e e n n o d e s .
[Link] – Formal Definition of MapReduce
• T h e M a p R e d u c e fra m e w ork p rov id e s a h ig h -le v e l a b stra c tio n th a t h id e s th e c om p le x ste p s in v olv ed in p a ra llel a n d d istrib u te d d a ta p roc e ss in g , su c h
a s p a rtitio n in g , m a p p in g , sy n ch ro n iz a tion , c om m u n ic a tio n , a n d sc h e d u lin g .
• It ex p ose s tw o m a in in te rfa c e s — th e M a p a n d R ed u c e fu n c tion s — w h ic h u se rs ca n o v errid e to d e fin e c u s tom p ro ce ss in g log ic su ited to th e ir
a p p lic a tio n s.
• T o ex e c u te a job , th e u se r f irst d e fin e s th es e M a p a n d R e d u c e f u n ction s , th en c a lls a c en tra l fu n c tio n n a m e d M a p R e d u c e (S p e c , & R e su lts) to sta rt th e
p ro c es s.
• T h e S p e c o b je c t is in itia liz e d b y th e u s er a n d filled w ith e ss en tia l in f orm a tion su c h a s th e n a m e s o f in p u t a n d o u tp u t file s, tu n in g p a ra m e te rs, a n d th e
fu n ction n a m e s f or M a p a n d R e d u c e . T h is te lls th e sy ste m h o w to h a n d le th e job .
• T h e u se r’ s p rog ra m typ ic a lly in c lu d e s th re e p a rts: th e M a p fu n c tio n , th e R e d u c e f u n c tio n , a n d th e M a in fu n c tio n . T h e M a p a n d R e d u ce su b rou tin e s
a re in v ok ed d u rin g e x e cu tion to p e rfo rm th e re q u ire d ta sks , co n trolle d b y th e M a p R e d u ce en g in e .
[Link] MapReduce Logical Data Flow
• In th e M a p R e d u ce fra m e w ork, b oth in p u t a n d o u tp u t d a ta a re h a n d le d in th e fo rm o f
(k ey , v a lu e ) p a irs. F o r th e M a p f u n c tio n , th e in p u t c ou ld b e s om e th in g like a f ile lin e
of fse t a s th e ke y a n d th e lin e c o n ten t a s th e v a lu e.
• T h e M a p f u n c tio n p roc e sse s e a ch in p u t p a ir a n d g e n e ra tes in te rm ed ia te (ke y ,
v a lu e) p a irs. T h is p ro c es sin g h a p p en s in p a ra llel f or a ll in p u t p a irs.
• T h e sys tem th e n so rts a n d g ro u p s a ll in term ed ia te p a irs b a se d on k ey s s o th a t
v a lu es w ith th e s a m e ke y a re g rou p e d in th e fo rm ( ke y, [list of v a lu es ]). T h is
sim p lif ie s th e n e x t p ro c es sin g ste p .
• T h e R e d u ce fu n c tio n ta ke s ea c h g ro u p e d ke y a n d its a ss oc ia te d v a lu e s a n d
p ro c es se s th e m to p rod u c e fin a l ou tp u t p a irs. F or ex a m p le , su m m in g a ll v a lu e s f or
a sp e c if ic w o rd to c ou n t its o c cu rre n c es .
• A c la ssic e x a m p le is th e w o rd co u n t p rob le m . F o r tw o lin e s of te xt, th e M a p
fu n ction ou tp u ts ( w o rd , 1 ) f or e a c h w o rd , a n d th e R e d u c e fu n c tion g rou p s id en tic a l
w ord s a n d su m s th e c ou n ts , res u ltin g in ou tp u ts lik e (p eo p le , 2 ).
[Link] Formal Notation of MapReduce Data Flow
• T h e M a p f u n c tion is a p p lie d in p a ra lle l to e v e ry in p u t (k ey , v a lu e) p a ir, a n d p rod u ce s
a
n e w se t o f In term e d ia te ( ke y, v a lu e) p a irs a s fo llow s:
( 6 .1 )
• T h e n th e M a p R e d u c e lib ra ry c o llec ts a ll th e p rod u c ed in te rm e d ia te ( ke y, v a lu e ) p a irs
fro m a ll in p u t (k ey , v a lu e ) p a irs, a n d sorts th e m b a se d o n th e “ ke y ” p a rt. It th en
g ro u p s th e v a lu e s of a ll o c cu rre n c e s of th e sa m e ke y . F in a lly , th e R e d u c e f u n c tio n is
a p p lie d in p a ra lle l to ea c h g ro u p , p ro d u c in g th e co llec tion of v a lu e s a s ou tp u t, a s
illu stra te d h ere :
[Link] – Strategy to Solve MapReduce Problems
• Identifying Unique Keys:
T h e first ste p in solv in g a M a p R ed u c e p ro b lem is to d ete rm in e th e unique key. A f ter th e M a p fu n ction e m its in te rm e d ia te (k ey , v a lu e ) p a irs, th e
M a p R ed u c e fra m e w ork a u to m a tic a lly so rts a n d g ro u p s th e m b y th e se ke y s. T h is m a ke s ke y id en tifica tion th e f ou n d a tio n o f p rob le m -so lv in g in
M a p R ed u c e.
• Word Count Problem:
T o co u n t h ow m a n y tim e s ea c h w ord a p p e a rs in a d oc u m e n t c olle c tion , e a c h w o rd b e c om e s a key, a n d th e in te rm ed ia te value is 1 (fo r e a c h
oc c u rre n ce ). T h e R e d u c e f u n ction th e n s u m s th e v a lu e s fo r e a c h w ord .
• Words with Same Length:
If th e g o a l is to c ou n t h ow m a n y w ord s h a v e th e sa m e n u m b e r o f le tte rs, th e ke y is e a c h w o rd , a n d th e v a lu e is its length . T h e R ed u c e fu n c tio n
a g g re g a te s th es e le n g th s for a n a lysis.
• Anagram Count Problem:
F o r c ou n tin g a n a g ra m s ( w ord s w ith th e s a m e le tte rs in d iffe re n t o rd e r, like “ liste n ” and “ sile n t” ) , ea c h w o rd is tra n s form e d in to a sorted
sequence of letters ( e.g ., “ e iln st” ). Th is b e c om e s th e ke y, a n d th e v a lu e is th e c ou n t of su c h a n a g ra m s.
[Link] MapReduce Actual Data and Control Flow
• Data Partitioning: T h e M a p R e d u c e lib ra ry d ivid e s th e in p u t d a ta (sto re d in G o og le F ile Sy ste m – G F S ) in to m u ltip le ch u n ks. Th es e ch u n ks
c orre sp on d to th e n u m b e r of m a p ta s ks, s o ea c h m a p ta sk p ro ce s se s o n e p ie ce of d a ta .
• Computation Partitioning: Us e rs w rite co d e u sin g o n ly Map and Reduce functions . T h e M a p R e d u c e sy ste m cre a te s a n d ru n s co p ies o f th is u se r
p ro g ra m a c ross d if fe ren t m a c h in e s (u s in g a m eth o d like th e fork sy ste m ca ll) .
• Master and Workers: O n e p ro g ra m co p y b e co m es th e m a ste r, m a n a g in g ta sk a ss ig n m en ts , w h ile th e oth e rs b e c om e w o rke rs. W ork ers a re
a ss ig n ed e ith er m a p or red u c e ta sk s b y th e m a s ter.
• Reading Input Data: E a ch m a p w ork er rea d s its a ss ig n ed d a ta c h u n k ( sp lit) a n d p a sse s it to its M a p fu n ction . U su a lly , o n e d a ta s p lit is a ss ig n ed
p e r m a p w o rke r.
• Map Function Execution: Th e M a p fu n c tion p roc e sse s its in p u t (ke y , v a lu e ) p a irs a n d p ro d u c es in te rm ed ia te ( ke y, v a lu e ) p a irs a s o u tp u t for th e
n e x t s tep s.
• Combiner Function (Optional): A Combiner ca n ru n lo ca lly o n m a p w o rke rs to red u ce th e a m ou n t o f in te rm ed ia te d a ta th a t n e e d s to b e s e n t
a c ros s th e n etw o rk. It w o rks like a m in i R ed u c e fu n c tion to op tim iz e p e rform a n ce .
•
• Partitioning Function: In te rm ed ia te (ke y , v a lu e ) p a irs fro m m a p w o rke rs a re d iv id e d
in to R p a rtition s (e q u a l to th e n u m b e r of red u c e ta sk s). Th is e n su re s th a t a ll id en tic a l
ke y s g o to th e sa m e re d u ce w o rke r. U su a lly , a h a sh f u n c tion (like H a s h (k ey ) m od R )
is u s ed .
• Synchronization: R e d u c e w orke rs’ w a it u n til a ll m a p ta sks fin is h b e fo re sta rtin g
th e ir w ork. Th is e n su re s s m oo th d a ta tra n s ition fro m th e m a p to th e red u ce p h a s e.
• Communication: E a ch re d u c e w o rke r co llec ts its a ssig n e d p a rtition e d d a ta from a ll
m a p w o rke rs u s in g re m ote p ro ce d u re c a lls . T h is a ll-to -a ll co m m u n ic a tion c a n c a u s e
n e tw ork c on g es tion , w h ich is a kn o w n p erf orm a n c e issu e .
• Sorting and Grouping: O n ce th e red u c ed w orke r h a s rec e iv e d th e d a ta , it sorts a n d
g ro u p s th e in term e d ia te (k ey , v a lu e ) p a irs b y k ey . T h is p re p a res th e d a ta fo r fin a l
re d u c tio n .
• Reduce Function Execution: E a c h g rou p o f ( ke y, [v a lu es ]) is p a ss ed to th e R ed u c e
fu n ction , w h ic h p roc e sse s a n d w rite s th e f in a l o u tp u t to th e sp e c if ie d ou tp u t file s in
th e u s er's p rog ra m .
[Link] Compute-Data Affinity
• Origin and Implementation: M a p R e d u c e w a s o rig in a lly d e v e lo p e d b y G o o g le a n d it s f irst v e rs io n w a s im p le m e n t e d in C
la n g u a g e . It w a s d e si g n e d t o w o rk c lo se l y w it h G o o g le F ile S y st e m (G F S ) .
• GFS and Block Storage: G F S s t o re s file s b y b re a kin g th e m in to fix e d -s iz e b lo c k s (c a lle d c h u n k s ), w h ic h a re t h e n d is trib u t e d
a c ro ss d iff e re n t n o d e s in a c lu st e r.
• MapReduce and GFS Integration: M a p R e d u c e sp li ts in p u t d a t a in to b lo c k s a s w e ll, ju s t lik e G F S . S o , it e ff ic ie n tly u s e s G F S
b y s im p ly a s sig n in g M a p f u n c t io n s t o t h e n o d e s t h a t a lre a d y c o n t a i n t h e d a t a b lo c ks .
• Data-Compute Affinity Concept: In st e a d o f m o v in g d a t a a c ro s s t h e n e t w o rk t o w h e re c o m p u t a t io n h a p p e n s , M a p R e d u c e
s e n d s t h e c o m p u ta t io n ( M a p f u n c t io n ) t o w h e re th e d a t a a lre a d y e x is t s. T h is m in im iz e s d a ta t ra n sf e r a n d im p ro v e s
p e rfo rm a n c e .
• Block Size Matching: T h e d e fa u lt b lo c k s iz e i n b o t h G F S a n d M a p R e d u c e i s 6 4 M B , m a k in g t h e ir i n t e g ra t io n sm o o th a n d
[Link] Twister
e f fic ie n t . and Iterative MapReduce
• Performance Comparison – MPI vs. MapReduce:
It ’ s im p o rta n t t o c o m p a re t h e p e rfo rm a n c e o f M P I a n d M a p R e d u c e .
B o t h s u ffe r fro m p a ra lle l o v e rh e a d s , m a in l y d u e t o lo a d im b a la n c e a n d
c o m m u n ic a t io n d e la ys (w h ic h a re lik e s y n c h ro n iz a t io n d e la y s in
m u lt i-t h re a d e d s y st e m s) .
• Communication Overhead in MapReduce:
M a p R e d u c e h a s h ig h c o m m u n ic a t io n c o st s b e c a u s e it re a d s a n d w rit e s
d a t a t h ro u g h fi le s, w h ile M P I t ra n s fe rs o n ly t h e re q u ire d d a t a d ire c t ly
b e t w e e n n o d e s u s in g n e t w o rk c om m u n ic a t io n . T h is m e a n s :
• M P I u se s δ (d e l ta ) f lo w – o n l y th e n e e d e d u p d a t e s a re t ra n sf e rre d .
• M a p R e d u c e u se s f u ll d a t a f lo w – a ll d a t a is re a d a n d w rit t e n ,
in c re a s in g ov e rh e a d .
• Iterative Applications and Communication: M a n y p a ra lle l a p p lic a ti o n s
ha ve re p e a t e d c o m p u te -t h e n -c o m m u n ic a t e st e p s . To im p ro v e
p e rf o rm a n c e in s u c h c a s e s , t w o m a in s t ra t e g ie s c a n b e a p p lie d :
1 . S t re a m d a t a b e tw e e n s t e p s in s te a d o f w rit in g t o d isk .
2 . U s e l o n g -ru n n in g t h re a d s o r p ro c e s s o rs t h a t h a n d l e o n ly d e lt a
c h a n g e s b e t w e e n i te ra t io n s .
• W h i le t h e s e c h a n g e s s ig n ific a n t ly im p ro v e sp e e d , t h e y m a y re d u c e
fa u lt t o le ra n c e a n d m a k e it h a rd e r t o h a n d le
d y n a m ic c h a n g e s (li ke c h a n g e s in t h e n u m b e r o f n o d e s ).
• Twister – An Improved Framework: T w i st e r is a M a p R e d u c e + +
fra m e w o rk t h a t u s e s t h e a b o v e s t ra t e g ie s . It k e e p s
s t a t ic d a t a in m e m o ry a nd o n ly se n d s d y n a m ic δ fl o w b e tw e e n
ite ra t io n s . It ru n s M a p a n d R e d u c e f u n c t io n s in lo n g -ru n n in g th re a d s ,
a v o id in g fre q u e n t d is k I/ O .
• T w is t e r is m u c h fa st e r t h a n t ra d it io n a l M a p R e d u c e in t a s k s lik e
K -m e a n s c lu s t e rin g , a s sh o w n in p e rfo rm a n c e c o m p a ris o n s ( F ig u re
6 . 8 ).
• F ig u re 6 .9 c o m p a re s t h e in t e rn a l s t ru c t u re o f fo u r f ra m e w o rk s –
H a d o o p , D ry a d , T w is t e r, a n d M P I. N o t a b ly , D ry a d a lso a v o i d s h e a v y
d i sk u sa g e b y u s in g p ip e s fo r d a t a tra n s fe r.
6.2.3 Hadoop Library from Apache
• H a d o op is a n op en -s ou rc e im p lem e n ta tion of th e M a p R ed u c e p rog ra m m in g m o d e l. It is c od ed in J a v a b y A p a c h e a n d u se s H D F S ( H a d o op
D istrib u te d F ile S ys tem ) in ste a d of G F S ( G o og le F ile Sy ste m ) a s its sto ra g e la y er.
• H a d o op ’ s c o re is d iv id e d in to tw o m a in la y ers : th e M a p R e d u c e e n g in e a n d H D F S . Th e M a p R ed u c e e n g in e p e rform s th e c om p u ta tio n s, w h ile
H D F S m a n a g es d a ta s tora g e .
HDFS (Hadoop Distributed File System)
H D F S is a d istrib u te d f ile sy ste m m o d e led a f ter G F S. It sto res a n d o rg a n iz e s la rg e file s o v er a d is trib u ted c om p u tin g sys tem , e n a b lin g relia b le a n d
ef fic ien t
d a ta a c c e ss.
HDFS Architecture
• H DFS h as a m a ste r/s la v e a rc h ite ctu re . It in c lu d e s on e N a m e N od e ( th e m a s ter) and m u ltip le D a ta N od es (th e sla v e s/ w orke rs ).
W h e n a file is store d , H D F S sp lits it in to fixe d -s iz e b lo c ks (ty p ica lly 6 4 M B ) . T h e se b lo ck s a re d istrib u te d a n d store d a c ros s th e D a ta N od es .
• Th e N a m e N od e m a n a g e s th e sy ste m ’ s m e ta d a ta a n d n a m e sp a c e. It ke e p s tra ck o f w h ich D a ta N od e h o ld s w h ich file b lo ck s. It is th e ce n tra l
m a n a g e r o f b loc k m a p p in g a n d m e ta d a ta .
• E a c h D a ta N o d e u su a lly ru n s on a s ep a ra te n od e in th e c lu s ter. It h a n d le s th e sto ra g e a n d re trie v a l o f file b lo ck s a n d co m m u n ica te s re g u la rly w ith
th e N a m e N od e.
• Th e m e ta d a ta ref ers to th e f ile m a n a g em e n t in fo rm a tio n , su c h a s f ile lo c a tio n s, w h ile th e n a m e sp a c e is th e a re a u se d to store th is m e ta d a ta .
HDFS Features
• H D F S is a sp ec ia liz ed d is trib u te d file sys te m th a t su p p orts b ig d a ta a p p lic a tio n s, fo cu sin g m a in ly on la rg e-s ca le sto ra g e a n d b a tc h p roc e ssin g .
• U n like g e n e ra l-p u rp o se file sy ste m s, H D F S d o es n ot req u ire a ll s ta n d a rd fe a tu re s like h ig h -lev e l co n c u rren c y c on tro l o r se cu rity , a s it is n o t
d e sig n e d for in te ra c tiv e u se .
• Se c u rity is n ot su p p o rte d in H D F S b y d ef a u lt, w h ic h se ts it a p a rt f rom m a n y tra d itio n a l file sys tem s.
• Tw o k ey f ea tu res th a t m a ke H D F S d iffe re n t a re : 1 . F a u lt to le ra n c e – th e a b ility to c on tin u e o p e ra tin g d e sp ite h a rd w a re fa ilu re s.
2 . H ig h -th ro u g h p u t a c ce ss – o p tim iz e d f or f a st a c ce ss to la rg e file s, p rioritiz in g d a ta s trea m in g ov e r
la te n c y.
HDFS Fault Tolerance
• F a u lt tolera n ce is a k ey f ea tu re of H D F S , a s H a d o op ru n s o n low -co st h a rd w a re w h e re h a rd w a re fa ilu re s a re e x p e cte d .
• B loc k rep lica tion e n su re s relia b ility . H D F S d iv id e s file s in to b lo ck s a n d rep lica te s e a ch b loc k a c ros s th e clu ste r. B y d e fa u lt, ea ch b loc k is
s tore d in th re e d iffe re n t loc a tion s.
• R e p lic a p la c e m en t is c a re fu lly m a n a g e d . T o b a la n c e re lia b ility a n d c om m u n ic a tio n c os t, H DF S p la c es :
▶ O n e re p lic a on th e sa m e n od e a s th e orig in a l d a ta ,
▶ O n e o n a d iffe re n t n o d e in th e sa m e ra ck ,
▶ O n e o n a n od e in a d iffe re n t ra c k.
• H ea rtb e a ts a n d B loc kre p o rts a re se n t re g u la rly from e a ch D a ta N o d e to th e N a m eN o d e .
▶ H e a rtb ea ts in d ica te th a t th e D a ta N od e is a liv e a n d w o rk in g .
▶ B lo c kre p orts lis t a ll th e b loc ks sto red o n th e D a ta N od e.
▶ T h e N a m eN o d e u se s th e s e m es sa g es to m a n a g e a n d m o n ito r b loc k re p lic a s a cro ss th e s ys tem .
HDFS High-Throughput Access to Large Data Sets
• H D F S is d es ig n ed f or b a tc h p roc e ssin g , n o t f or re a l-tim e o r in tera c tiv e u se . S o, it f oc u se s m ore o n h ig h d a ta a c ce s s th ro u g h p u t ra th er th a n
lo w la ten cy .
• S in c e a p p lica tion s u s in g H D F S u su a lly d ea l w ith la rg e f iles , H D F S b re a k s th e m in to la rg e b lo ck s (e .g ., 6 4 M B or m ore ).
• U s in g la rg e b lo c k s iz e s h elp s in tw o m a in w a y s:
▶ It red u c es m e ta d a ta sto ra g e b y m in im iz in g th e n u m b e r o f b lo ck s p e r f ile.
▶ It sp e e d s u p d a ta re a d in g b y a llo w in g la rg e a m ou n ts of d a ta to b e re a d se q u e n tia lly a n d co n tin u o u sly, im p ro v in g stre a m in g
p erf orm a n c e .
HDFS Operation: Reading and Writing Files
• Reading a File:
▶ When a user wants to read a file, they first send an “open” request to the NameNode. The NameNode responds with the addresses of DataNodes
that hold the replicas of the file blocks. The number of addresses depends on the replication factor (default is 3).
▶ The user then uses the read function to connect to the nearest DataNode with the first block. After the first block is read, the connection closes,
and the user repeats the process for the next blocks until the whole file is read.
• Writing a File:
▶ To write a new file, the user sends a “create” request to the NameNode. If the file doesn't already exist, the NameNode approves and the user starts
writing using the write function.
▶ The data is first sent to a data queue, and a data streamer monitors this. The streamer then asks the NameNode for suitable DataNodes to store
the replicas of each block.
▶ The block is written to the first DataNode, which then forwards it to the second, and so on, until all replicas are stored. This process repeats for
each block of the file until the full file is written and replicated across the HDFS system.
[Link] Architecture of MapReduce in Hadoop
• T h e M a p R e d u ce e n g in e is th e top la y er of H a d o op , a n d it h a n d le s th e d a ta flow a n d c on tro l f low f or M a p R e d u c e job s ru n n in g on a d is trib u ted
sy ste m . It w o rks tog e th er w ith H D F S fo r d a ta s tora g e.
• J u st like H D F S, th e M a p R e d u c e e n g in e u s es a m a ste r/s la v e ( w orke r) a rc h ite c tu re. It in clu d e s on e J ob T ra c ke r (m a ste r) a n d m a n y T a sk T ra ck ers
(s la v e s) .
• T h e J o b T ra c ke r m a n a g e s a n d m on itors th e e n tire M a p R e d u c e jo b , a s sig n in g ta sks to th e T a sk T ra ck ers sp re a d a cro ss th e c lu s ter.
• E a ch T a skT ra c ke r h a n d les th e a c tu a l e x e cu tion of e ith e r m a p or red u ce ta sks o n its loc a l n o d e .
• E v e ry T a s kT ra ck e r h a s a se t n u m b e r of e x ec u tion slots, w h ich a re like th re a d s u s e d to ru n ta s ks in p a ra lle l. Th e n u m b e r of s lots d e p e n d s o n
th e C P U c ore s a n d th re a d s a v a ila b le o n th a t n od e (e .g ., N C P U s × M th rea d s = M × N s lo ts).
• O n e M a p T a sk p e r B loc k: E a c h d a ta b lo c k in H D F S is p roc es se d b y o n e m a p ta sk , a n d th a t ta sk ru n s in on e slot of a T a skT ra c ke r. T h is m ea n s
th e re is a on e -to-o n e c on n e c tion b etw e en a m a p ta sk a n d a d a ta b lo c k in H DF S.
[Link] Running a Job in Hadoop
• R u n n in g a H a d o op jo b in v o lv e s th ree m a in p a rts: th e u se r n o d e , th e J o b T ra c ke r, a n d m u ltip le T a skT ra c ke rs.
• S ta rtin g th e J ob : T h e u s er p rog ra m c a lls a f u n ction c a lle d ru n J ob (c on f) on th e u s er n od e, w h ere c on f h old s co n fig u ra tion s ettin g s fo r th e
M a p R e d u c e p roc es s a n d H D F S . T h is sta rts th e j ob e x ec u tion in H a d o op .
• J o b Su b m ission P ro ce ss :
• Th e u se r n o d e re q u es ts a n e w jo b ID fro m th e J ob Tra c ke r.
• It sp lits th e in p u t d a ta in to c h u n ks c a lled in p u t sp lits.
• Th e u se r n o d e co p ies n e c es sa ry file s like th e job ’ s J A R file , co n fig u ra tion , a n d in p u t sp lits to th e J ob Tra c ke r’ s file sy ste m .
• F in a lly , th e jo b is su b m itte d to th e J ob T ra c ke r u sin g th e su b m itJ ob () fu n c tio n .
• T h e J ob T ra c ke r cre a te s o n e m a p ta sk p er in p u t sp lit a n d a ssig n s th es e ta sks to T a s kT ra c ke rs’ a v a ila b le ex e c u tio n slots . It trie s to a ss ig n m a p
ta s ks w h ere th e d a ta is loc a ted ( d a ta lo ca lity ) fo r e ff ic ie n c y.
• It a lso c re a tes re d u c e d ta sk s a n d a ssig n s th e m to T a sk T ra ck ers , b u t w ith ou t c on sid e rin g d a ta lo ca lity . T h e n u m b e r o f re d u ce ta sks is s et b y th e u se r.
• E a ch T a sk Tra c ke r rec e iv e s th e jo b ’ s J A R f ile a n d sto res it loc a lly. It ru n s th e a s sig n e d m a p or red u c e ta sk s in sid e a J a v a Virtu a l M a c h in e (J V M ),
f ollow in g th e in s tru c tion s in th e J A R .
• T a sk T ra ck ers se n d h e a rtb e a t m e ssa g es reg u la rly to th e J o b T ra c ke r. T h e se h e a rtb e a ts te ll th e J o b T ra c ke r th a t th e T a s kT ra ck e r is a liv e a n d w h e th er
it is re a d y to ta ke n e w ta s ks or is still b u sy .
6.2.4 Dryad and DryadLINQ from Microsoft
• D rya d off ers m o re f le x ib ility th a n M a p R e d u c e b y a llo w in g u s ers to d e fin e th e
d a ta flow a s a n a rb itra ry d irec te d a cy c lic g ra p h (D A G ), w h e re v ertic es
re p re se n t c om p u ta tio n s a n d e d g es re p re se n t c o m m u n ic a tion ch a n n els.
• T h e D ry a d ru n tim e h id e s d eta ils su c h a s d a ta p a rtitio n in g , s ch ed u lin g ,
sy n c h ron iz a tio n , c o m m u n ic a tion , a n d fa u lt to le ra n c e to sim p lif y p ro g ra m m in g .
• T h e tw o m a in c o m p on en ts co n trollin g D ry a d a re th e job m a n a g er a n d th e
n a m e se rv er. T h e job m a n a g er b u ild s th e jo b ’ s D A G fro m th e u se r’ s
p ro g ra m , sc h e d u le s ta sks , a n d m on ito rs ex e c u tio n .
• T h e n a m e s erv e r p ro v id es re sou rc e a n d n e tw ork top olog y in form a tion to th e
job m a n a g e r, e n a b lin g e ff ic ie n t m a p p in g of ta sks c o n sid e rin g d a ta a n d
c om p u ta tio n loc a lity.
• E a c h c lu s te r n od e ru n s a lig h tw eig h t d a em o n to e x ec u te a ssig n e d ta s ks a n d
a c t a s a p rox y f or co m m u n ica tion a n d m on itorin g .
• D a ta tra n sf er b etw e en v e rtic e s u s es c h a n n e ls im p le m e n te d w ith sh a re d
m e m ory , TC P s oc ke ts, or d istrib u te d file s ys tem s , f orm in g a 2 D d is trib u te d
p ip e s tru ctu re .
• D rya d su p p orts d y n a m ic D A G m o d ifica tion s d u rin g ru n tim e , s u ch a s a d d in g
v ertic es o r e d g es , m erg in g g ra p h s , a n d h a n d lin g jo b in p u t/o u tp u t.
• F a u lt to le ra n c e is m a n a g ed b y re -e x ec u tin g fa ile d v e rte x ta s ks o n oth e r n od es
a n d re cre a tin g co m m u n ic a tion ch a n n els w h en e d g es fa il.
• D rya d is a g e n e ra l-p u rp o se f ra m ew o rk s u p p ortin g sc rip tin g , M a p R e d u c e -style
p ro g ra m m in g , a n d S Q L in teg ra tio n .
[Link] DryadLINQ from Microsoft
• DryadLINQ is b u ilt on top o f M ic ros oft’ s Dryad e x e c u tio n fra m ew o rk a n d
in teg ra tes it w ith .NET’ s LINQ (L a n g u a g e In teg ra te d Q u ery ), m a k in g
d istrib u te d c lu ste r co m p u tin g a c ce ss ib le to re g u la r p rog ra m m e rs.
• It en a b le s u s ers to w rite d a ta -p a ra lle l p rog ra m s in h ig h -le v e l la n g u a g es
like C # u s in g fa m ilia r LIN Q sy n ta x , w h ic h is co m p ile d in to a D rya d
e x ec u tion p la n .
• D ry a d L IN Q p ro g ra m s b eg in b y c re a tin g a L IN Q e x p re ssion ob je ct; d u e to
d e f erre d e xe c u tio n , th e c om p u ta tio n d oe sn 't sta rt u n til e x p lic itly trig g e re d .
• C a llin g To D rya d T a b le( ) in itia te s th e p roc es s— D rya d L IN Q co m p ile s th e
L IN Q e x p re ssion in to a Dryad execution plan, d iv id in g it in to
su b e x p re ssion s f or e x e cu tion on se p a ra te D rya d v e rtic e s.
• A c u sto m Dryad job manager m a n a g e s th e jo b ex e c u tio n , b u ild s th e D A G ,
a n d sc h e d u le s th e v ertice s b a se d on res ou rc e a v a ila b ility .
• E a ch v e rtex ru n s its sp e cific su b p ro g ra m , a n d th e ou tp u t is w ritte n to
re su lt ta b le s on c e th e jo b co m p le tes .
• T h e Dryad job manager te rm in a te s a fte r co m p letion a n d re tu rn s co n trol to
D ry a d L IN Q , w h ich th e n w ra p s th e res u lts in DryadTable ob je cts .
• F in a lly, th e u se r a p p lic a tio n reg a in s c on tro l a n d c a n a c c es s o u tp u t u s in g
sta n d a rd .N E T ite ra tors .
• N o t a ll a p p lica tion s fo llow a ll n in e e x e cu tion ste p s, b u t D ry a d L IN Q
p ro v id es a s ea m le ss in te rfa c e c om b in in g f a m ilia r p ro g ra m m in g w ith
p o w erf u l d istrib u te d ex e c u tio n .
6.2.5 Sawzall and Pig Latin High-Level Languages
• Sawzall is a high-level scripting language developed by Google, designed for parallel data processing on top of the
MapReduce framework.
• It was originally created by Rob Pike to process Google’s log files efficiently and transformed long batch jobs into
interactive sessions, enabling real-time insights.
• Sawzall performs distributed and fault-tolerant processing of very large datasets, including Internet-scale data.
• The execution model involves local data partitioning and filtering using on-site scripts, followed by aggregation for
final results.
• Users write analysis scripts in Sawzall, which are translated by the runtime engine into MapReduce jobs running
across multiple cluster nodes.
• Sawzall combines ease of scripting, cluster computing power, and reliability from redundancy and fault tolerance,
and it has been released as an open-source project.
[Link] Pig Latin
• P ig L a tin is a h ig h -le v e l d a ta f lo w la n g u a g e d e v elop e d b y Y a h o o!, im p le m e n ted on top
of H a d oo p a s p a rt of th e A p a ch e P ig p roj ec t.
• It p ro v id es a s crip tin g a p p ro a c h to p a ra lle l d a ta p roc es sin g , sim ila r to S a w z a ll a n d
D ry a d L IN Q , b u t s u p p orts m ore S Q L -lik e c on s tru cts su c h a s J oin , w h ic h S a w z a ll la c ks.
• W h ile D ry a d LIN Q is S Q L-b a se d , P ig L a tin a n d S a w z a ll f ollo w th e N oS Q L m od e l b u t
w ith p o w erf u l d a ta flow c a p a b ilitie s.
• P ig La tin a b stra c ts th e co m p lex ity of p a ra lle lism , le ttin g u s ers fo cu s o n e le m e n t-w ise
op e ra tio n s a n d s u p p orte d co llec tiv e fu n c tio n s, en su rin g s id e ef fe cts o cc u r o n ly in
c on tro lled o p e ra tio n s.
• U n like d e cla ra tiv e S Q L , P ig L a tin u se s a p roc ed u ra l d a ta flow p ip elin e , c le a rly
sp ec ify in g h o w d a ta is m a n ip u la te d ste p -b y -ste p .
• It s u p p o rts u se r-d e fin ed fu n ction s (U D F s ) a s first-c la s s e n titie s th a t c a n b e u se d w ith
op e ra tors like L oa d , S tore , G ro u p , F ilter, a n d F ore a c h , of fe rin g c u s tom p ro ce ss in g
fle xib ility.
• T h e A p a c h e P ig s ys tem tra n sla te s P ig L a tin sc rip ts in to M a p R ed u c e jo b s fo r
e x ec u tion on H a d oo p c lu s ters , en a b lin g s ca la b le a n d fa u lt-to le ra n t d a ta a n a ly sis
w ork flo w s.
6.2.6 Mapping Applications to Parallel and Distributed Systems
• A p p lic a tio n s c a n b e m a p p e d to d iffe re n t p a ra llel a n d d istrib u te d a rc h ite c tu res u s in g six c a te g orie s o f a p p lica tion m od els. T h e o rig in a l f iv e
c a teg orie s f oc u se d m a in ly o n sim u la tio n s, w h ile a s ix th ca teg ory a d d re sse s m o d e rn d a ta -in te n siv e c om p u tin g .
• C a te g o ry 1 in v olv e s S IM D ( Sin g le In stru c tion M u ltip le D a ta ) a rc h itec tu re s w h e re ta s ks e x e cu te in loc k-ste p . It w a s im p o rta n t h is to ric a lly b u t is
n o w la rg e ly o b so le te .
• C a te g o ry 2 is m o re sig n ific a n t tod a y a n d a lig n s w ith th e S P M D (S in g le P ro g ra m M u ltip le D a ta ) m o d e l o n M IM D ( M u ltip le In stru c tio n M u ltip le
D a ta ) sy ste m s. E a c h ta s k ru n s th e sa m e c o d e in d e p e n d en tly, id e a l fo r c om p le x , irre g u la r p rob le m s w ith co m p u te – c om m u n ica te p h a s es .
• C a te g o ry 3 fe a tu res a s yn c h ro n ou s ly in tera c tin g ob j ec ts, su ita b le f or e v e n t-d riv en sy ste m s like se a rc h a lg o rith m s o r O S -le v el th re a d m a n a g e m en t.
It relie s on sh a re d m e m ory f or fa st sy n c h ron iz a tio n , b u t it's n ot co m m on in la rg e -sc a le p a ra lle l sy ste m s.
• C a te g o ry 4 in clu d e s in d e p e n d e n t p a ra lle l c om p on e n ts , re q u irin g m in im a l c om m u n ic a tio n . T h is c a teg ory h a s g ro w n in re le v a n c e a n d is id e a l fo r
g rid a n d c lo u d co m p u tin g w ith p lea s in g ly p a ra llel w o rkloa d s.
• C a te g o ry 5 c ov e rs co a rse -g ra in ed w o rkflow s, w h e re d if fe ren t in d e p e n d e n t m od u le s a re lin k ed a t a h ig h er lev e l u sin g tw o-lev e l p ro g ra m m in g . It
su its g rid / clou d in fra stru c tu re s a n d su p p o rts w o rkflow -b a s ed c o m p u tin g .
• C a te g o ry 6 , c a lle d M a p R e d u ce + + , fo cu s e s o n d a ta -in ten s iv e a p p lic a tion s. It h a s th re e s u b ty p e s: ( 1 ) M a p -o n ly job s (like C a te g ory 4 ), (2 ) cla ss ic
M a p R ed u c e w ith p a ra lle l m a p a n d re d u c e s ta g e s, a n d ( 3 ) e x ten d e d M a p R e d u c e w ith m o re c om p le x d a ta p roc e ssin g . T h is c a teg ory o v erla p s
w ith 2 a n d 4 b u t em p h a siz e s d a ta I/ O a n d loos e ly sy n c h ron o u s stru c tu re s.
6.3 Programming Support of Google App Engine
• Language and Development Environment: G o og le A p p E n g in e ( G AE ) su p p orts d e v e lo p m e n t in b o th J a v a a n d P yth o n . J a v a d e v elop ers b e n e fit
f rom to ols su ch a s th e E c lip s e p lu g -in f or loc a l d e b u g g in g a n d th e G o og le W e b T o olkit (G W T ) f or b u ild in g w e b a p p lic a tio n s. D e v elop ers c a n
a lso u s e oth e r J VM -b a se d la n g u a g e s like J a v a S c rip t a n d R u b y th rou g h in te rp re te rs o r c om p ilers . F o r P yth o n , fra m e w orks like Dj a n g o a n d
C h erry P y a re co m m on ly u se d , a n d G o og le a lso p ro v id e s a b u ilt-in w e b a p p fra m e w ork to sim p lify d e v elop m en t.
• Datastore and Data Handling: G A E u se s a N oS Q L d a ta store to sto re e n titie s th a t ca n b e u p to 1 M B in s iz e . Th e s e e n tities h a v e sc h e m a -le ss
p rop e rtie s, e n a b lin g f le x ib le d a ta m od e lin g . Q u e ries c a n retrie v e e n tities o f a s p e c if ic k in d , filte red a n d so rte d b a se d on p rop e rty v a lu e s. J a v a
d e v elop ers c a n u se J a v a D a ta O b je c ts (J D O ) a n d J a v a P e rsiste n c e A P I (J P A ), b oth su p p o rted b y th e o p e n -so u rc e D a ta N u c le u s p la tfo rm .
P y th on d e v e lop e rs u se G Q L (G oo g le Q u e ry L a n g u a g e ), w h ic h is sim ila r to S Q L.
• Transactions and Consistency: G A E ’ s d a ta sto re is stro n g ly co n siste n t a n d u se s op tim istic c o n cu rre n c y co n trol. T ra n sa c tion s en su re a tom ic
o p e ra tion s a cro ss en titie s w ith in a n en tity g rou p . If m u ltip le p ro c es se s try to u p d a te th e sa m e e n tity , G A E re tries th e tra n sa c tio n a lim ited
n u m b e r o f tim e s. T h is tra n s a ction m od el e n su re s eith e r a ll o p e ra tio n s in a g ro u p su c ce e d or n o n e a t a ll. E n titie s in th e s a m e g ro u p a re sto red
to g e th e r to im p ro v e tra n sa c tio n e f fic ien c y .
• Blobstore and Caching Mechanisms: T o h a n d le la rg er d a ta o b je c ts, G AE p rov id e s a B lob store se rv ic e th a t a llow s s tora g e of f iles u p to 2 G B .
F or p erf orm a n c e op tim iz a tio n , G A E in c lu d e s a m e m c a c h e se rv ic e for in -m e m ory d a ta c a ch in g . T h is c a c h in g sy ste m w o rk s b oth in d e p en d e n tly
a n d a lo n g sid e th e d a ta s tore , h elp in g re d u c e a c c es s tim e s a n d s e rve r lo a d .
• External Communication Support: G A E a p p lic a tio n s c a n co m m u n ic a te w ith e x tern a l s yste m s v ia th e U R L F e tc h se rv ic e , w h ic h su p p o rts H T T P
a n d H T T P S p roto co ls . G o og le ’ s S ec u re D a ta C o n n ec tion ( SD C ) a llow s s ec u re tu n n e lin g f rom a p riv a te in tra n e t to th e G A E -h os ted a p p lic a tio n
o v e r th e p u b lic in te rn et, e n a b lin g h yb rid d ep lo ym e n ts.
• Integration with Google Services: G A E in te g ra tes w ith v a riou s G o og le s erv ice s su c h a s M a p s, C a le n d a r, Y o u T u b e , D oc s, a n d m o re th rou g h th e
G o og le D a ta A P I. T h is A P I en a b le s a p p lic a tion s to a c c e ss a n d m a n ip u la te res ou rc e s p rov id e d b y G o og le ’ s se rv ice s se a m le ss ly .
• User Authentication and Image Handling: G A E s u p p orts a u th e n tic a tion th rou g h G o og le A c c ou n ts . T h is m ea n s u s ers w ith e x istin g G o og le
c re d e n tia ls ( e.g ., G m a il) ca n log in to G A E a p p s w ith ou t n e e d in g a n e w a cc ou n t. A d d itio n a lly, th e Im a g e s s erv ic e a llo w s a p p lica tion s to p e rfo rm
b a sic im a g e o p e ra tio n s su c h a s re siz in g , rota tin g , flip p in g , c rop p in g , a n d e n h a n c in g .
• Background Processing and Task Scheduling: G A E s u p p orts b a c kg ro u n d ta s k e x ec u tion in tw o w a ys: sc h e d u led ta sks u sin g c ron jo b s a n d
d y n a m ic ta sk q u eu e s. C ro n jo b s a llow d e v e lo p e rs to ru n ta s ks a t sp ec if ic in terv a ls ( e.g ., h o u rly or d a ily ), w h ile ta s k q u eu e s le t a p p lic a tio n s
e n q u e u e ta sk s to b e p ro ce ss e d a syn ch ro n ou s ly a f ter u se r re q u es ts.
• Resource Quotas and Limits: G A E a p p lica tion s a re su b jec t to u sa g e q u o ta s th a t res tric t h o w m u c h C P U tim e , b a n d w id th , sto ra g e , a n d oth e r
re s ou rc es a n a p p c a n c on s u m e . Th es e q u ota s h e lp p re v e n t a n y on e a p p lic a tio n fro m ov e ru sin g s ys tem re so u rc es a n d en s u re co st c on tro l.
G A E a llo w s a ce rta in lev e l of f ree u sa g e , m a kin g it su ita b le fo r s m a ll-sc a le a p p s o r d e v e lo p m e n t/ te s tin g p h a se s.
6.3.2 Google File System (GFS)
• G oo g le F ile S ys tem ( G F S) w a s cre a te d to s tore th e h u g e a m ou n t o f d a ta n e e d e d b y G oo g le’ s s ea rc h e n g in e . It is a d istrib u ted file sy ste m
d es ig n ed to w ork relia b ly on c h e a p , u n re lia b le h a rd w a re. U n lik e tra d ition a l file sy ste m s, G F S w a s s p e cia lly m a d e to m e e t G o og le ’ s u n iq u e
n e e d s a n d is c lo se ly in te g ra te d w ith G o og le a p p lic a tion s .
• G F S w a s b u ilt w ith so m e ke y id ea s in m in d . It a ss u m es h a rd w a re f a ilu re s w ill h a p p en of ten b e ca u se it u se s in e xp e n sive c o m p o n en ts . It is
o p tim iz e d f or v e ry la rg e file s— u s u a lly o v er 1 0 0 M B a n d so m e tim e s s ev e ra l g ig a b y te s. T o m a n a g e th e se la rg e f iles e fficie n tly , G F S u se s b ig
b lo ck s iz e s of 6 4 M B , m u ch la rg e r th a n th e u s u a l 4 K B b lo ck s. T h e file s a re m os tly w ritten on c e a n d th en a p p e n d e d to, w ith v e ry fe w ra n d o m
w rite s. R e a d in g is m a in ly la rg e , co n tin u ou s stre a m s. B e c a u se o f th is , G F S foc u se s on h ig h d a ta th rou g h p u t ra th er th a n lo w d e la y.
• T o k ee p d a ta s a fe , G F S c op ie s e a ch d a ta b lo ck ( c a lled a c h u n k) o n a t le a st th ree d iffe re n t s erv e rs. O n e m a ste r se rv er m a n a g e s a cc e ss a n d
ke e p s m eta d a ta , m a kin g th e sy ste m sim p ler a n d re d u c in g c om p le x c oord in a tion . G F S d o e s n o t u se ca c h in g b ec a u s e its a c c es s p a ttern s
d on ’ t b en ef it f rom it. W h ile it w ork s like P O S IX sy ste m s, G F S a lso h a s sp e cia l A P Is th a t le t G oo g le’ s a p p s se e w h e re d a ta b loc ks a re s tore d
a n d u s e fe a tu re s lik e sn a p sh o ts a n d rec o rd a p p e n d s.
• W h e n w ritin g d a ta , th e clie n t firs t a s ks th e m a s ter w h ic h ch u n k s erv e r h a s co n trol o f th e c h u n k a n d w h ere its c op ie s a re. T h e m a ste r rep lies
w ith th e p rim a ry a n d se co n d a ry re p lica lo c a tio n s. T h e clie n t s e n d s d a ta to a ll rep lica s , w h ic h s tore it te m p ora rily in m e m ory . T h e n , th e c lie n t
te lls th e p rim a ry re p lica to a p p ly th e w rite in o rd e r. T h e p rim a ry a ss ig n s a s eq u e n ce n u m b er a n d p a sse s th e w rite to s ec on d a ry re p lic a s, w h ic h
d o th e sa m e in o rd e r. A f ter a ll re p lic a s c on f irm , th e p rim a ry in fo rm s th e clie n t. If so m e th in g f a ils, th e c lien t re trie s th e w rite .
• G F S o ffe rs a sp e c ia l “ rec o rd a p p e n d ” o p e ra tio n th a t lets m u ltip le clie n ts a d d d a ta to th e e n d o f a file sa f ely a n d a to m ic a lly . Th is is v e ry
u s ef u l fo r a p p lic a tio n s like w eb c ra w lin g , w h e re n e w d a ta is c o n tin u o u sly a d d e d . G F S ch oos e s th e of fse t a n d g u a ra n tee s th e a p p en d h a p p en s
a t le a st on c e .
• F a u lt to le ra n c e is v e ry im p o rta n t f or G F S . B oth th e m a ste r a n d c h u n k s e rv e rs ca n re sta rt q u ic kly to re d u ce d ow n tim e. B e c a u se c h u n ks a re
sto re d in th re e or m o re p la ce s, th e sy ste m c a n h a n d le a t le a st tw o fa ilu re s w ith ou t los in g d a ta . G F S a lso u se s a sh a d ow m a s ter to k ee p a
b a ck u p of m e ta d a ta a n d h e lp re c ov e r if th e m a in m a ste r f a ils. It c h ec ks d a ta in te g rity b y sto rin g c h e ck su m s fo r e v e ry 6 4 K B o f d a ta .
• G FS’ s a rch ite ctu re h a s on e m a ste r a n d m a n y c h u n k se rv e rs. Th e m a ste r co n trols th e f ile sy ste m ’ s n a m es p a c e , m eta d a ta , a n d loc kin g . It
ta lks re g u la rly to c h u n k se rv e rs to c h ec k th e ir h e a lth a n d a ssig n ta s ks like d a ta b a la n c in g a n d re c ov e ry. H a v in g a sin g le m a ste r m a ke s d es ig n
sim p le r b u t c ou ld ca u se a b o ttle n e ck . T o red u c e th is , G o og le u s es a sh a d ow m a ste r a n d c a c h es c on tro l m e ssa g es .
• O v e ra ll, G F S is d e s ig n ed for h ig h a v a ila b ility , g o od p e rfo rm a n ce , a n d sc a la b ility . It w o rks w e ll w ith la rg e f iles , su p p o rts stre a m in g rea d s a n d
w rite s, h a n d le s h a rd w a re f a ilu res g ra c e fu lly , a n d off ers ef ficien t a p p e n d op era tion s . Th is m a k es it id e a l f or p roc e ssin g b ig d a ta o n
in ex p e n siv e , c om m o d ity se rv ers .
6.3.3 BigTable, Google’ s NOSQL System
• B ig T a b le is G o og le ’ s N o SQ L sy ste m c re a ted to store a n d re trie v e m a ssiv e a m ou n ts of stru c tu re d a n d s em i-stru c tu red d a ta , s u ch a s w e b
p a g e s , u se r-sp ec if ic d a ta , a n d g eo g ra p h ic in form a tion . It h a n d le s d a ta lik e U R Ls w ith th e ir c on te n t a n d m e ta d a ta , p e r-u se r p re fe re n c es a n d
e m a ils , a s w e ll a s p h ys ic a l lo c a tio n d a ta u s ed in a p p lic a tio n s like G oo g le E a rth .
• T h e sc a le o f d a ta m a n a g e d b y B ig Ta b le is e n orm o u s, w ith b illion s of U R L s h a vin g m u ltip le v e rsio n s a v e ra g in g a b ou t 2 0 K B e a c h , h u n d red s
o f m illion s of u se rs g en e ra tin g th ou s a n d s o f q u e ries p e r s ec on d , a n d g eo g ra p h ic d a ta s ets e x c e ed in g 1 0 0 te ra b yte s. T ra d itio n a l
c o m m erc ia l d a ta b a s e s c a n n ot ef fic ien tly h a n d le th is sc a le a n d c om p le x ity , w h ich m otiv a ted G oog le to d e v e lo p B ig T a b le f or b e tte r
p e rform a n ce a n d sc a la b ility a t a low e r in cre m e n ta l c os t.
• B ig T a b le ’ s d e sig n foc u se s on a llow in g a sy n c h ron o u s u p d a tes b y m u ltip le p ro c es se s w h ile e n su rin g a cc e ss to th e m o st re c en t d a ta . It
s u p p orts v e ry h ig h rea d a n d w rite th rou g h p u t, p ote n tia lly m illion s of o p e ra tion s p er se c on d . It a ls o p rov id e s e ffic ie n t sc a n n in g ov e r e n tire
d a ta se ts or s u b s ets a n d su p p o rts c om p le x jo in s fo r la rg e d a ta se ts. A d d ition a lly , it tra c ks d a ta c h a n g e s ov e r tim e, su c h a s d iff ere n t
v e rsion s o f w e b p a g e s from m u ltip le c ra w ls.
• A rc h itec tu ra lly , B ig T a b le a cts a s a d is trib u ted m u lti-le ve l m a p , of fe rin g a f a u lt-to le ra n t a n d p e rsiste n t stora g e s erv ic e. It s ca le s to
th ou sa n d s of s erv e rs, m a n a g in g te ra b yte s of in -m e m o ry d a ta a n d p e ta b yte s of d isk -b a se d d a ta , w ith m illio n s of re a d s a n d w rite s p e r
s e co n d . Th e sy ste m is se lf -m a n a g in g , a llo w in g se rv e rs to b e a d d e d or re m ov e d d y n a m ica lly a n d a u tom a tic a lly b a la n c in g loa d a c ros s
m a ch in e s.
• B ig T a b le w a s d e sig n e d a n d im p le m e n te d sta rtin g in e a rly 2 0 0 4 a n d is u se d a cro ss v a rio u s G o og le s erv ic es , in c lu d in g G oo g le S e a rch , O rku t,
a n d G o og le M a p s. So m e B ig Ta b le d e p loy m e n ts m a n a g e d a ta o n th e ord er of h u n d red s of te ra b y tes sp re a d ov e r th o u sa n d s of se rv e rs .
• T h e sy ste m is b u ilt on G oo g le ’ s c lou d in fra s tru c tu re, le v era g in g ke y c om p on e n ts lik e th e G oo g le F ile S yste m (G F S ) f or p e rsiste n t sto ra g e ,
a sc h e d u le r fo r m a n a g in g job s, a loc k se rv ic e fo r m a s te r e le c tion a n d sy ste m b oo tstra p p in g , a n d M a p R e d u c e fo r p roc e ssin g a n d a c ce ss in g
B ig T a b le d a ta e ffic ie n tly .
[Link] Tablet Location Hierarchy (in BigTable)
• B ig Ta b le o rg a n iz es its d a ta in to u n its c a lle d ta b le ts, a n d loc a tin g th es e ta b le ts in v olv es a th ree -le v e l h iera rc h y . Th e top le v e l b eg in s w ith a sp e c ia l
file s tore d in C h u b b y, G oo g le ’ s d istrib u te d lo c k s erv ic e, w h ic h co n ta in s th e lo ca tion of th e ro ot ta b let.
• T h e roo t ta b let h o ld s in form a tion a b ou t a ll o th e r ta b lets in a sp e cia l in te rn a l ta b le c a lled th e M E T A D A TA ta b le . T h is root ta b le t its e lf is th e first
ta b le t o f th e M E T A D A T A ta b le b u t is tre a ted d if fe ren tly — it is n ev e r sp lit d u rin g ta b le t m a n a g e m en t to m a in ta in a m a x im u m of th re e h iera rc h ica l
le v e ls fo r loc a tin g a n y u se r ta b le t.
• E a ch M E T A D A TA ta b le t m a p s to a se t of u s er ta b le ts a n d store s th e ir loc a tio n s. T h e se m a p p in g s a re in d e x e d b y row ke y s, w h ich a re a
c om b in a tion of th e ta b le t’ s ta b le ID a n d its e n d row , a llo w in g q u ic k loo ku p a n d loc a tio n tra c kin g .
• F o r re lia b ility a n d e fficie n cy , B ig T a b le in clu d e s v a riou s o p tim iz a tion s . C h u b b y en su re s th e a v a ila b ility of th e roo t ta b let’ s loc a tio n file . T h e
B ig T a b le m a ste r c a n q u ic k ly d ete c t th e s ta tu s of ta b le ts b y sc a n n in g th e ta b le t se rv e rs.
• T o m a in ta in ef ficien c y a n d d a ta in te g rity , ta b le t se rv e rs p e rf orm co m p a ction op era tion s to c om p res s a n d o rg a n iz e sto re d d a ta . A lso , s h a red
lo g g in g is u se d so th a t m u ltip le ta b lets c a n w rite to a s in g le log , m in im iz in g sp a ce a n d im p rov in g s ys tem c on s iste n c y.
6.3.4 Chubby – Google’ s Distributed Lock Service
• C h u b b y is G oo g le’ s co a rse -g ra in e d d is trib u ted loc kin g a n d c oo rd in a tio n s erv ice u se d p rim a rily f or e lec tin g lea d e rs a n d m a n a g in g c on fig u ra tio n
a c ros s d is trib u ted s yste m s like G F S a n d B ig T a b le .
• C h u b b y a c ts lik e a lig h tw e ig h t file sy ste m w ith a h iera rc h ica l n a m e sp a c e in w h ich it s tore s s m a ll file s. T h e se file s co n ta in m eta d a ta o r c on trol
in form a tion ra th e r th a n la rg e d a ta se ts, w h ic h a re m a n a g e d b y G F S .
• T h e sy ste m is b u ilt u sin g th e P a x o s c on s en s u s p ro toc o l, w h ich m a k e s it h ig h ly f a u lt-tole ra n t a n d c a p a b le o f fu n c tio n in g e v en if so m e of its n od e s
fa il. T h is e n su re s co n siste n cy a n d a v a ila b ility of th e s yste m ’ s sta te .
• A C h u b b y c e ll c on s ists of fiv e se rv ers , e a c h m a in ta in in g th e sa m e n a m es p a c e . C lien ts c om m u n ic a te w ith C h u b b y s erv e rs via a C h u b b y lib ra ry ,
e n a b lin g th e m to p e rform file o p e ra tion s su c h a s re a d , w rite, loc k, a n d u n loc k.
• T h e P a x o s p roto c ol ru n n in g o n a ll se rv e rs en s u re s th a t a ll file op era tion s a re c on siste n t a n d re lia b le a c ros s th e sy ste m . Du e to its re lia b ility a n d
c on s iste n c y g u a ra n tee s, C h u b b y h a s b ec o m e G oog le ’ s p rim a ry in te rn a l n a m in g se rv ic e .
• S y ste m s like G F S a n d B ig T a b le d ep e n d o n C h u b b y to ele ct a p rim a ry se rv e r from a m on g se v era l rep lic a s, h e lp in g c oo rd in a te d is trib u ted a c tiv itie s
sa f ely a n d e ffic ie n tly .
6.4 Programming on Amazon AWS and Microsoft Azure
• In th is se c tio n , w e c on s id er th e p rog ra m m in g s u p p ort in th e A W S p la tf orm . It b eg in s w ith a re v ie w of A W S a n d its u p d a te d s e rv ic e off erin g s ,
f oc u sin g on h o w th e y su p p ort c lo u d -b a s e d a p p lic a tion d e v elop m e n t.
• T h e se c tio n ex a m in e s c ore A W S s erv ice s, in c lu d in g E C 2 (E la s tic C om p u te C lou d ), S3 (S im p le S tora g e S erv ice ), a n d S im p leD B , p rov id in g
p rog ra m m in g e x a m p le s to d em on stra te h ow th e se se rv ic e s a re u se d in p ra c tic e .
• A m a z o n , like A z u re , of fe rs a R ela tio n a l D a ta b a se S e rv ic e (R D S) , a lo n g w ith m es sa g in g in te rf a c es d is cu s se d ea rlie r in th e te x t. T h e se s u p p ort
stru ctu re d d a ta sto ra g e a n d co m m u n ic a tion b e tw ee n clou d c om p on e n ts .
• A m azo n’ s E la s tic M a p R e d u c e ( E M R ) p ro v id es H a d oo p -like b ig d a ta p roc e ssin g ca p a b ilitie s on E C 2 in sta n c e s. It e n a b le s sc a la b le a n d
d istrib u te d d a ta a n a lys is , th o u g h A m a z on d o e s n ot s u p p ort B ig T a b le d ire ctly .
• A W S s u p p orts N oS Q L d a ta b a s es th rou g h S im p le D B , w h ic h off e rs flex ib le d a ta sto ra g e w ith o u t a fixe d sc h em a . It w a s in trod u c ed a n d d isc u s se d
in e a rlie r se c tio n s a s p a rt o f A W S ’ s d a ta b a se se rv ice s.
• A m a z o n p ro v id es th e S im p le Q u e u e Se rv ice (S Q S ) a n d Sim p le N otif ic a tion S e rvic e ( S N S) , w h ic h a re c lo u d -b a se d im p le m e n ta tion s o f m e ss a g in g
sy ste m s th a t su p p ort d e c ou p lin g a n d n otific a tion d e liv e ry in d istrib u te d a p p lic a tion s .
• C lou d p la tfo rm s e ff ic ie n tly ru n b ro ke rin g sy ste m s, w h ic h a re u se fu l f or c on tro llin g se n so rs a n d su p p o rtin g m o b ile b a c ke n d op e ra tion s. T h e y a re
e sp ec ia lly ef fe ctiv e in m a n a g in g s erv ic es fo r s m a rtp h o n es a n d ta b le ts.
• A u to -sc a lin g a u tom a tic a lly a d ju s ts th e n u m b e r of E C 2 in sta n c e s b a se d on d e f in e d co n d itio n s. It e n su re s p e rfo rm a n ce d u rin g h ig h d em a n d a n d
c os t sa vin g s d u rin g low u sa g e b y d yn a m ic a lly sc a lin g re sou rc e s.
• E la s tic L o a d B a la n c in g d istrib u te s in c o m in g tra f fic a cro ss m u ltip le E C 2 in sta n ce s. It a v o id s f a iled n o d e s a n d b a la n c e s loa d s a c ross a c tiv e
in sta n c e s to m a in ta in p erf orm a n c e a n d re lia b ility.
• A W S C lo u d W a tc h e n a b le s m on itorin g o f clou d re so u rc es su ch a s E C 2 in sta n c es . It p rov id e s v is ib ility in to re so u rc e u sa g e, o p e ra tio n a l h e a lth , a n d
p erf orm a n c e m e tric s, in c lu d in g C P U, d is k a c tiv ity , a n d n etw o rk u s a g e .
6.4.1 Programming on Amazon EC2
• A m a z o n p io n e e re d v irt u a l m a c h i n e (V M ) h o s t in g fo r c lo u d a p p li c a t io n s, a llo w in g u s e rs t o re n t V M s in st e a d o f p h y sic a l
se rv e rs . T h is e n a b le s c u s t o m e rs t o ru n th e ir o w n s o ft w a re a n d m a n a g e a p p lic a t io n s fl e x ib ly .
• E C 2 of fe rs e l a s t ic c a p a b ilit ie s — u se rs c a n c re a t e , la u n c h , a n d t e rm i n a te s e rv e r in st a n c e s a s n e e d e d , p a yin g o n ly fo r t h e
ti m e t h e in s t a n c e s a re a c t iv e . T h is m a ke s E C 2 c o s t-e ff e c t iv e a n d s c a la b le .
• A m a z o n p ro v id e s p re in s ta lle d V M s c a ll e d A M Is (A m a z o n M a c h in e Im a g e s ), w h ic h a re te m p la t e s c o n fi g u re d w it h L i n u x o r
W i n d o w s op e ra t in g s y s te m s a n d o p ti on a l so f tw a re t o o ls f o r q u ic k d e p lo y m e n t .
• T h e re a re t h re e t y p e s o f A M Is— p u b lic , p riv a t e , a n d p a id . T a b le 6 . 1 2 o u t lin e s t h e ir c h a ra c te ris t ic s , s h o w in g op t io n s fo r
d iffe re n t le v e ls o f c u s t o m iz a t io n , s e c u rit y , a n d c o s t .
• T h e t y p ic a l E C 2 V M c re a t io n p ro c e s s in c lu d e s : c re a ti n g a n A M I, g e n e ra t in g a k e y p a ir f o r se c u re a c c e s s, c o n fi g u rin g a
fire w a ll fo r in s t a n c e p ro t e c t io n , a n d l a u n c h in g t h e in st a n c e .
• A M Is a re b a s e d o n A m a z o n ’ s v irt u a liz e d in fra s t ru c t u re , w h i c h in c lu d e s c o m p u ti n g , st o ra g e , a n d s e rv e r re s o u rc e s. T h i s
v irt u a l fo u n d a ti o n s u p p o rt s d yn a m ic in st a n c e d e p lo ym e n t a n d m a n a g e m e n t .
6.4.2: Amazon Simple Storage Service (S3
• Amazon S3 is a web service that allows users to store and retrieve any amount
of data at any time, from anywhere on the web. It provides object-oriented
storage accessible via web protocols.
• S3 stores data as objects, each placed in a bucket and accessed using a
unique key. Objects can contain values, metadata, and access control data,
making S3 essentially a key-value storage system.
• S3 supports both REST (Web 2.0) and SOAP interfaces. Users can read, write,
and delete objects ranging from 1 byte to 5 GB using these APIs, allowing
integration with various applications.
• Amazon S3 ensures high reliability with geographically distributed redundancy,
offering 99.999999999% durability and 99.99% availability. A lower-cost
Reduced Redundancy Storage (RRS) option is also available.
• S3 uses strong authentication mechanisms to prevent unauthorized access.
Objects can be made public or private, with access controlled using ACLs
(Access Control Lists) and per-object URLs.
• The default download protocol is HTTP, but S3 also supports BitTorrent,
reducing costs for high-volume data distribution.
• Storage costs range from $0.055 to $0.15 per GB per month, depending on the
volume. The first 1 GB/month of input/output is free, followed by tiered transfer
rates. There are no data transfer charges between EC2 and S3 in the same
region or between EC2 Northern Virginia and S3 U.S. Standard regions.
6.4.3 Amazon Elastic Block Store (EBS) and SimpleDB
• Amazon Elastic Block Store (EBS) offers persistent block-level storage that can save the state of EC2 instances,
which traditionally would be lost after shutdown. EBS allows users to store virtual machine images and mount them
to active EC2 instances for persistent use.
• Unlike Amazon S3, which functions as object storage, EBS acts more like a traditional file system, accessible through
operating system-level disk access. It offers a volume block interface, suitable for applications needing consistent,
low-latency storage.
• EBS volumes range from 1 GB to 1 TB and can be mounted as raw block devices to EC2 instances. Multiple volumes
can be attached to a single instance, and users can format these volumes with a file system or use them directly as
block storage.
• EBS supports incremental snapshots, enabling users to back up volumes efficiently and restore data quickly. This
feature enhances performance during data saving and recovery operations.
• EBS follows a pay-per-use model similar to EC2 and S3. Storage is billed at $0.10 per GB per month, while I/O
requests are charged at $0.10 per million requests (as of October 6, 2010), providing flexible, usage-based billing.
• An equivalent functionality to EBS is also available in open-source cloud platforms such as Nimbus, offering similar
persistent storage capabilities.
[Link] Amazon SimpleDB Service
• A m a z o n S im p le D B o f fe rs a s im p lif ie d v e rs io n of t h e re la t io n a l d a ta b a s e m o d e l, o rg a n i z in g d a t a i n t o d o m a in s (s im ila r t o t a b l e s) ,
it e m s (ro w s ), a n d a t t rib u t e s (c o lu m n s) . U n lik e t ra d it io n a l re l a t io n a l d a t a b a se s , S im p le D B a llo w s m u lt ip le v a lu e s f o r a s in g le
c e ll, g iv in g m o re fle x i b i lit y b u t w it h re la x e d c o n s ist e n c y .
• S im p le D B is d e s ig n e d fo r d e v e lo p e rs w h o n e e d q u ic k, sc a la b le d a t a s t o ra g e a n d re t rie v a l w ith o u t w o rry in g a b o u t rig id
s c h e m a s o r c o m p le x c o n s is t e n c y ru le s . T h is m a k e s i t i d e a l f o r lig h t w e ig h t a p p lic a t io n s a n d m e ta d a t a st o ra g e .
• S im p le D B c h a rg e s $ 0 .1 4 0 p e r M a c h i n e H o u r, w it h t h e fi rs t 2 5 h o u rs f re e e a c h m o n t h ( a s o f O c t ob e r 6 , 2 0 1 0 ). T h is m a ke s it
c o s t -e ff e c t iv e f or s m a ll-sc a le u s e c a s e s .
• S im p le D B c a n b e c o n s id e re d a “ L it t le T a b le ” fo r m a n a g in g s m a l l, s t ru c t u re d d a t a se t s , m u c h l ike A z u re T a b le . It c o n tra st s
w it h s ys t e m s lik e B ig T a b le , w h ic h a re d e s ig n e d fo r h a n d lin g la rg e -s c a le d a t a . S im p l e D B e v o lv e d fro m e a rlie r s y st e m s lik e
A m a z o n D y n a m o , w h i c h in s p ire d it s d is t rib u t e d a n d s c h e m a -le s s a rc h it e c t u re .
6.4.4: Microsoft Azure Programming Support
• M ic ro s o ft A z u re p ro v id e s a ric h p ro g ra m m in g m o d e l b u ilt o n v irt u a liz e d h a rd w a re a n d a s o p h ist ic a t e d c o n tro l e n v iron m e n t
k n o w n a s t h e A z u re fa b ric . T h is f a b ric s u p p o rt s d y n a m i c re s o u rc e a llo c a t io n , f a u lt t o le ra n c e , a n d d o m a in n a m e s y s te m ( D N S )
f u n c t io n s . S e rv ic e s c a n b e a u t o m a t ic a lly m a n a g e d a n d d e p lo y e d u s in g X M L t e m p la te s, w ith m u lt ip le in s t a n c e s la u n c h e d o n
d e m an d .
• O n c e s e rv ic e s a re ru n n in g , A z u re o ff e rs e x t e n si v e m o n it o rin g fe a t u re s , i n c lu d in g a c c e s s t o e v e n t lo g s, d e b u g g i n g t ra c e s ,
p e rf o rm a n c e c o u n t e rs , a n d c ra s h d u m p s . A lt h o u g h li v e d e b u g g in g o f c lo u d a p p li c a t io n s is n o t s u p p o rt e d , d e v e l o p e rs c a n d e b u g
u s in g t h e se re c o rd e d tra c e s a n d lo g s sa v e d in A z u re s t o ra g e
• A z u re d iv id e s it s fe a t u re s in t o s t o ra g e an d c o m p ute fu n c t io n s.
A p p lic a t io n s c o n n e c t t o t h e In t e rn e t v i a a " w e b ro le " — a c u st o m iz e d
v irt u a l m a c h in e fo r w e b h o s ti n g . T h e s e V M s a re c a lle d a p p lia n c e s a n d
s e rv e b a s ic w e b a p p l ic a t io n fu n c t io n s , p ro v id in g t h e fro n t e n d fo r
A z u re c lo u d a p p s.
• A lo n g s id e t h e w e b ro l e , A z u re in t ro d u c e s a " w o rke r ro le ," w h ic h i s u se d
fo r e x e c u t in g b a c k g ro u n d p ro c e s s e s a n d h a n d lin g t a sk s t h a t re q u ire
s c h e d u le d c o m p u t e re s o u rc e s . T h is st ru c t u re re fl e c t s t h e c lo u d ’ s
n eed fo r s c a la b le , d is t rib u t e d c o m p u t in g . W o rk e r ro le s s u p p o rt
H T T P (S ) a n d T C P p ro t o c o ls .
• E a c h rol e in A z u re h a s t h re e m a in life c yc l e m e t h o d s : O n S t a rt (), w h ic h
p e rfo rm s in it ia l iz a t io n a n d s ig n a ls re a d in e s s t o t h e lo a d b a la n c e r;
R u n (), w h ic h e x e c u te s t h e c o re a p p lic a t io n lo g ic ; a n d O n S t o p (), w h ic h
h a n d le s g ra c e fu l s h utd o w n. Th ese m e th o d s e n s u re st ru c tu re d
e x e c u t io n a n d m a n a g e m e n t of c lo u d se rv i c e s.
• T h e A z u re rol e -b a se d m o d e l o ff e rs fle x ib ilit y a n d m o d u la rit y , sim ila r t o
p ra c t ic e s in A W S a n d G o o g le A p p E n g in e . R o le s in A z u re c a n a ls o b e
lo a d -b a la n c e d , p ro v id in g s c a la b ili ty a nd a v a i la b ilit y fo r c lo u d
a p p l ic a t io n s, m a k in g A z u re ’ s p rog ra m m i n g m o d e l c o m p a ra b le t o
o t h e r m a jo r c lo u d p l a t fo rm s .
[Link] SQLAzure
• SQLAzure Service:
P ro v id e s S Q L S e rv e r a s a c lo u d -b a s e d s e rv ic e , e n a b l in g u se rs t o m a n a g e re la t io n a l d a t a b a s e s o n A z u re .
• Storage Access Interfaces:
A lm o s t a ll st o ra g e m o d a li tie s u s e R ES T in t e rfa c e s f o r a c c e s s , w h ic h a re a u t o m a t ic a lly lin ke d w it h U R L s .
E x c e p t io n : A z u re D riv e s (s im ila r t o A m a z o n E B S ) o f fe r a file s y st e m in t e rfa c e w it h d u ra b l e N T F S v o lu m e s b a c ke d b y b lo b
s to ra g e .
• Data Replication:
D a t a in S Q L A z u re st o ra g e is re p lic a t e d t h re e t im e s a c ro s s d iffe re n t p h y s ic a l lo c a tio n s t o e n s u re h ig h fa u lt t o le ra n c e a n d d a t a
d u ra b ilit y .
T h is re p lic a tio n g u a ra n t e e s c o n s is t e n t a c c e s s a n d a v a ila b i lit y.
• Blob Storage Structure:
T h e b a sic st o ra g e u n it is b l o b s , w h i c h a re a n a l og o u s t o A m a z o n S 3 o b je c t s .
B lo b s t o ra g e is o rg a n iz e d h ie ra rc h ic a lly in t h re e le v e ls :
• Account: T h e t o p -l e v e l ro o t .
• Containers: A n a lo g o u s t o fo ld e rs o r d ire c t o rie s in s id e a n a c c o u n t .
• Blobs: T h e a c t u a l d a t a fi le s, w h ic h a re e it h e r P a g e b lo b s o r B lo c k b l ob s .
• Block Blobs:
D e si g n e d p rim a rily fo r s t re a m in g la rg e a m o u n ts o f d a t a .
C o m p ris e d o f b lo c k s u p t o 4 M B e a c h .
E a c h b lo c k h a s a u n iq u e 6 4 -b y te ID .
M a x im u m siz e o f a b lo c k b lo b is 2 0 0 G B .
• Page Blobs:
D e si g n e d fo r ra n d o m re a d /w rit e o p e ra ti on s .
C o n s is t of a n a rra y o f p a g e s , a llo w in g e f fic ie n t u p d a t e s t o p o rt io n s o f th e b lo b .
M a x im u m siz e is 1 T B , s u ita b le f or v irtu a l h a rd d isk s a n d o t h e r la rg e file s re q u irin g f re q u e n t w rit e s .
• Metadata:
E a c h b lo b c a n h a v e m e t a d a t a a t t a c h e d a s n a m e -v a lu e p a irs ( ke y -v a lu e p a irs) .
[Link] Azure Tables (Detailed Bullet Points)
• Azure Table and Queue Storage:
D e si g n e d fo r s m a lle r-sc a le d a ta v o lu m e s a n d re lia b le m e ss a g e d e liv e ry b e t w e e n w e b a n d w o rk e r ro le s .
T h e s e s to ra g e m o d e s c o m p le m e n t A z u re ’ s c o m p u t e se rv ic e s b y su p p o rt in g a s y n c h ro n o u s c o m m u n i c a t io n a n d lig h t w e ig h t
d a t a s to ra g e .
• Queues:
Q u e u e s a re u se d f o r w o rk s p oo l in g , m e a n in g t h e y h e lp m a n a g e a n d d ist rib u te ta s ks b e t w e e n d iffe re n t p a rt s o f a n
a p p l ic a t io n .
Q u e u e s su p p o rt u n li m it e d m e s s a g e s w it h a m a x im u m si z e o f 8 K B p e r m e s s a g e .
S u p p o rte d op e ra t io n s in c lu d e : P U T (a d d m e s s a g e ), G E T ( re trie v e m e s s a g e ), D E L E T E (re m o v e m e s sa g e ), C R E A T E ( c re a t e
q u e u e ), a n d D E L E T E (d e le t e q u e u e ) .
T h is g u a ra n te e s re li a b le m e s sa g e d e liv e ry a t le a s t o n c e .
• Tables:
T a b le s a re N o S Q L k e y -v a lu e s t o re s d e s ig n e d t o h o ld la rg e n u m b e rs o f e n t it ie s d ist rib u te d a c ro s s m a n y se rv e rs.
N o lim it o n t h e n u m b e r o f e n t it ie s p e r t a b le , a llo w in g g re a t s c a la b ilit y .
E a c h e n t it y (ro w ) c o n t a in s p ro p e rti e s ( c o lu m n s ), u p t o 2 5 5 p e r e n t it y .
P ro p e rt ie s a re s t o re d a s t rip le s : n a m e , t y p e , a n d v a lu e .
• Mandatory Properties:
E a c h e n t it y m u s t h a v e t w o s p e c ia l p ro p e rt ie s :
• PartitionKey: G ro u p s re la te d e n t it ie s fo r e f fic ie n t st o ra g e a n d q u e ry in g . E n ti ti e s w it h t h e s a m e P a rt iti o n K e y a re s to re d
t o g e t h e r, w h ic h o p t im i z e s q u e ry p e rfo rm a n c e .
• RowKey: U n iq u e ly id e n t ifie s e a c h e n t it y w it h in a p a rt it io n .
• Entity Size Limit:
E a c h e n t it y c a n h o ld u p t o 1 M B o f d a t a .
F or l a rg e r d a t a , it is re c o m m e n d e d t o s to re t h e d a t a in b lo b st o ra g e a n d ke e p o n ly a re fe re n c e (lin k ) in th e t a b le .
• Querying Support:
A z u re T a b le s c a n b e q u e rie d u s in g f a m i lia r M ic ro so ft t e c h n o lo g i e s s u c h a s A D O . N E T a n d L IN Q , m a k in g it e a s y fo r
d e v e lo p e rs t o in te g ra t e .
6.5 Emerging Cloud Software Environments
• T h is se c ti o n d is c u ss e s v a rio u s c lo u d o p e ra t in g s y s te m s a n d so f tw a re e n v iron m e n t s th a t a re e m e rg in g o r p o p u la r in th e c lo u d
c o m p u t in g s p a c e .
• It b u ild s o n e a rlie r v i rt u a li z a t io n c o n c e p t s (in tro d u c e d in C h a p t e r 3 ) a n d n o w f o c u s e s o n p ro g ra m m in g re q u ire m e n t s a n d
m a n a g e m e n t fe a t u re s .
• C o v e rs se v e ra l o p e n -s o u rc e c lo u d p la t fo rm s :
• E u c a ly p t u s
• N im b u s
• O p e n N e b u la
• S e c to r/ S p h e re
• O pe nS ta c k
• A lso in t ro d u c e s A n e ka , a c l o u d p ro g ra m m in g to o lk it d e v e lo p e d b y t h e U n i v e rsi ty o f M e lb o u rn e , a im e d a t s u p p o rt in g c lo u d
a p p l ic a t io n d e v e lo p m e n t .
6.5.1 Open Source Eucalyptus and Nimbus
Eucalyptus Overview
• E u c a ly p t u s o rig in a t e d a s a re s e a rc h p ro je c t a t U n iv e rs it y o f C a li fo rn ia , S a n t a B a rb a ra .
• G oa l: t o e n a b le c lo u d c o m p u t in g c a p a b ilit ie s fo r a c a d e m ic s u p e rc o m p u t e rs a n d c lu s t e rs.
• O ffe rs a w e b s e rv ic e in t e rfa c e th a t is c o m p a t ib le w it h A m a z o n W e b S e rv i c e s (A W S ) E C 2 , m e a n in g u s e rs fa m ili a r w it h A W S
c a n in t e ra c t s im il a rly .
• P ro v id e s a st o ra g e se rv ic e c a lle d W a lru s, c om p a ra b le t o A m a z o n S 3 , f o r m a n a g i n g V M im a g e s a n d p e rsis t e n t s t o ra g e .
• In c lu d e s a u se r-frie n d ly in t e rfa c e fo r m a n a g in g c lo u d u s e rs a n d v irt u a l m a c h i n e im a g e s .
[Link] Eucalyptus Architecture
• E u c a ly p t u s is a n o p e n -s o u rc e s o ft w a re p l a t fo rm d e s ig n e d t o su p p o rt c lo u d c o m p u t in g .
• A rc h it e c t u re fo c u s e s o n m a n a g i n g v irt u a l m a c h in e (V M ) im a g e s e ffic ie n t ly .
• S u p p o rts b o t h c o m p u t e c lo u d (p ro c e s s in g re s o u rc e s ) a n d st o ra g e c lo u d ( d a t a s t o ra g e ) fu n c t io n a lit ie s .
• T h e a rc h it e c t u re is d e ta ile d in E u c a ly p t u s w h i te p a p e rs a n d w a s e a rli e r in t ro d u c e d w i th a v irt u a l c lu st e rin g p e rsp e c t iv e .
• D e si g n e d t o a llo w d e v e lo p e rs a n d s ys t e m a d m in is tra t o rs t o m a n a g e V M im a g e s fo r d e p lo y m e n t a n d sc a lin g .
6.5.2 OpenNebula, Sector/Sphere, and OpenStack
OpenNebula
• O p e n N e b u la is a n o p e n -s o u rc e t oo l kit d e s ig n e d t o tra n s fo rm e x is ti n g IT in fra s t ru c t u re in t o a n In fra s t ru c t u re a s a S e rv ic e ( Ia a S )
c lo u d w ith c lo u d -li ke i n t e rfa c e s.
• T h e a rc h it e c t u re (s e e F ig u re 6 .2 8 ) is fle x ib l e a n d m o d u la r, a llo w in g it t o in t e g ra t e w it h v a rio u s s t o ra g e s y s te m s, n e t w o rk s e t u p s ,
a n d h yp e rv iso rs .
• C o re c o m p o n e n t : A c e n tra liz e d m a n a g e r t h a t c o n t ro ls t h e e n t ire V M lif e c y c le , in c lu d in g :
• D yn a m ic n e t w o rk s e t u p fo r g ro u p s o f V M s .
• M a n a g in g st o ra g e re q u ire m e n t s like V M d is k im a g e d e p lo y m e n t a n d d y n a m ic s o ft w a re e n v iro n m e n t c re a t io n .
• C a p a c it y m a n a g e r (sc h e d u le r): C o n t ro ls h o w V M s a re a llo c a t e d re s o u rc e s . T h e d e fa u lt s c h e d u le r u se s a re q u ire m e n t / ra n k
m a t c h m a k in g p o l ic y .
• S u p p o rts m ore c o m p le x s c h e d u lin g u sin g le a se m o d e ls a n d re s e rv a t io n s.
• A c c e s s d riv e rs : A b s tra c t u n d e rly in g in fra st ru c tu re d e t a il s, e x p o s in g u n ifo rm f u n c t io n a lit y fo r:
• M o n it o rin g
• S t o ra g e
• V irt u a liz a t io n se rv i c e s
• T h is a b s t ra c t io n m e a n s O p e n N e b u la is p la t fo rm -in d e p e n d e n t , a b le t o w o rk w it h v a rio u s v irt u a liz a ti on t e c h n o lo g ie s .
• P ro v id e s m a n a g e m e n t in t e rf a c e s fo r in t e g ra t io n w ith d a ta -c e n t e r t o o ls lik e :
• A c c o u n t in g s y s te m s
• M o n it o rin g fra m e w o rk s
• Im p le m e n t s th e lib v irt A P I (a n o p e n V M m a n a g e m e n t in t e rfa c e ) a n d o f fe rs a c o m m a n d -lin e in t e rfa c e (C L I).
• E x p o s e s a s u b s e t o f f u n c t io n a lit y v ia a c l o u d in t e rfa c e fo r e x t e rn a l u se rs .
• S u p p o rts d y n a m ic e n v iro n m e n t s:
• C a n h a n d le th e a d d it io n o r fa ilu re o f p h ys ic a l re s o u rc e s .
• F e a t u re s lik e li v e m ig ra t io n a n d V M s n a p sh o t s h e lp m a in t a in fle x ib ilit y a n d fa u lt t o le ra n c e .
• S u p p o rts h y b rid c lo u d m o d e ls:
• C a n c o n n e c t lo c a l in fra s t ru c t u re w it h e x t e rn a l c lo u d s w h e n lo c a l re so u rc e s a re in su ffic ie n t .
• U se s c lo u d d riv e rs t o in t e rfa c e w it h p u b lic c l o u d s .
• S u p p o rts p e a k d e m a n d s c a li n g o r h ig h a v a ila b ili ty ( H A ) st ra t e g ie s .
• D ri v e rs c u rre n t ly a v a ila b le in c lu d e :
• A m a zo n E C 2
• E u c a ly p t u s
• E la st ic H o s t s
• S t o ra g e m a n a g e m e n t:
• H a s a n Im a g e R e p o si to ry fo r m a n a g in g V M d i sk im a g e s e a s ily .
• S im p lifie s im a g e s h a rin g a n d m u lt iu s e r e n v iro n m e n t s w it h im a g e a c c e s s c o n t ro l.
• U se rs c a n u se p re d e fin e d im a g e s o r u p lo a d a n d m a n a g e th e ir o w n .
[Link] Sector/Sphere
• Sector/Sphere is a software platform designed for large-scale distributed data storage and simplified distributed
data processing on clusters of commodity hardware.
• It supports deployment both within single data centers and across multiple data centers connected by high-speed
networks.
Sector (Distributed File System)
• Sector is a distributed file system (DFS) that manages large data sets and supports access from any location with a
high-speed network.
• Implements fault tolerance by replicating data and managing these replicas across the network.
• The system is network topology aware, placing replicas to improve reliability, availability, and data access
performance.
• Uses UDP for message passing (faster than TCP due to no connection setup) and UDT (UDP-based data transfer)
for reliable, high-speed data transmission over wide-area networks.
• Provides a client programming API, various tools, and a FUSE (Filesystem in Userspace) module for easy file system
access.
Sphere (Parallel Data Processing Framework)
• Sphere is tightly coupled with Sector, serving as a parallel data processing engine working directly on data stored in
Sector.
• Supports User-Defined Functions (UDFs) to process data segments in parallel, ideally where data is stored (data
locality), reducing data movement.
• Provides fault tolerance by restarting failed data segment processing on other nodes.
• Both inputs and outputs in Sphere applications are Sector files, enabling chaining of multiple Sphere jobs for
complex processing workflows.
• Data sharing and exchange happen through the Sector file system.
Sector/Sphere Architecture (Figure 6.29)
• C o m p o s e d o f fo u r m a in c o m p o n e n t s:
• Security Server: A u t h e n t ic a te s m a s t e r se rv e rs , s la v e n o d e s , a n d u se rs .
• Master Servers: C o re o f th e in f ra s t ru c t u re , m a n a g in g f ile sy s t e m m e t a d a t a , s c h e d u lin g jo b s , a n d h a n d lin g u s e r re qu e s t s.
S u p p o rts m u lt ip le a c t iv e m a s t e rs d y n a m ic a ll y jo in in g or l e a v i n g .
• Slave Nodes: S to re d a t a a n d e x e c u t e p ro c e ss in g t a s k s. C a n b e d is trib u t e d a c ros s m u lt ip le d a t a c e n te rs c o n n e c t e d b y
h ig h -sp e e d n e t w o rks .
•
Client Component: P ro v id e s A P Is a n d to o ls fo r d a t a a c c e s s a n d p ro c e ss in g b y u s e rs.
Space Component
• A n e w ly d e v e lo p e d fra m e w o rk w it h i n t h e p la t fo rm th a t s u p p ort s c o lu m n -b a s e d d is t rib u t e d d a t a t a b le s.
• T a b le s a re s t o re d b y c o lu m n s a n d s e g m e n t e d a c ro s s m u lti p le s la v e n o d e s .
• T a b le s a re i n d e p e n d e n t ; n o re la ti on s h ip s (jo in s) a re s u p p o rt e d .
• S u p p o rts a lim i te d se t o f S Q L -lik e o p e ra tio n s su c h a s :
• T a b le c re a t io n a n d m o d ific a t io n
• K e y -v a l u e u p d a t e s a n d lo o k u p s
• S e l e c t U D F o p e ra t io n s
[Link] OpenStack – Detailed Points
• Project Overview and Origin:
• O p e n S t a c k w a s jo i n t ly in it ia t e d b y R a c k s p a c e H o s t in g a n d N A S A i n J u ly 2 0 1 0 .
• It a im s t o d e v e lo p a n o p e n -s ou rc e c lo u d c o m p u t in g p la t fo rm fo r p u b lic a n d p riv a t e c lo u d s.
• B u ilt o n t h e p rin c ip le s o f o p e n n e s s a n d c o lla b o ra t io n , i t fo s t e rs a c o m m u n it y-d riv e n d e v e lo p m e n t m o d e l in v o lv in g
d e v e lo p e rs , re s e a rc h e rs , c lo u d p ro v id e rs, a n d e n t e rp rise s.
• Philosophy and Openness:
• T h e p ro je c t a d h e re s s t ric t ly t o a n o p e n so u rc e p h il os o p h y, a v o id in g p ro p rie ta ry c o d e .
• It p ro m o t e s t h e u s e o f op e n A P Is su c h a s t h o s e fro m A m a z o n W e b S e rv ic e s (A W S ), e n s u rin g b ro a d i n t e ro p e ra b ilit y a n d
v e n d o r n e u t ra lit y .
• T h is o p e n n e s s f a c ili ta t e s c o m m u n i ty in n o v a t io n a n d a v o id s v e n d o r lo c k -i n .
• Modular Components:
• O p e n S t a c k d e v e lo p m e n t is fo c u se d p rim a rily o n t w o c o re in f ra s tru c t u re s e rv ic e s :
1. OpenStack Compute (Nova):
• P ro v id e s t h e c o m p u ti n g fa b ric fo r c lo u d s e rv ic e s .
• R e s p on s ib le fo r p ro v is io n in g a n d m a n a g in g v irt u a l m a c h in e s (V M s ) a c ros s c lu st e rs o f p h ys ic a l s e rv e rs.
• M a n a g e s V M l ife c y c le o p e ra t io n s lik e b o o t , su sp e n d , re s u m e , t e rm in a t e , a n d liv e m ig ra t io n .
• S u p p o rts m u lt ip le h y p e rv is o rs s u c h a s K V M , X e n , H y p e r-V , a n d V M w a re E S X i.
• H a n d le s s c h e d u lin g o f V M in s t a n c e s , n e t w o rk c o n fig u ra t io n , a n d se c u rit y p o lic ie s .
2. OpenStack Object Storage (Swift):
• P ro v id e s h ig h ly a v a ila b le , d i st rib u t e d , a n d s c a la b le o b je c t s t o ra g e .
• Id e a l f o r s t o rin g u n st ru c t u re d d a t a lik e im a g e s , v id e o s , b a c k u p s , a n d v irt u a l m a c h in e sn a p s h o t s .
• U se s c o m m o d it y h a rd w a re in c lu s t e rs to e n s u re c o s t e ff ic ie n c y.
• S u p p o rts re p lic a t io n , d a t a d u ra b ilit y , a n d fa u lt t o le ra n c e th ro u g h re d u n d a n t st o ra g e a c ro s s d if fe re n t n o d e s a n d re g io n s .
• D e si g n e d t o h a n d le p e t a b y t e -sc a le d a ta w it h m u lt i-t e n a n t s u p p o rt .
• Image Repository Development:
• A n e w f e a t u re u n d e r d e v e lo p m e n t is t h e Im a g e R e p o sit o ry , m e a n t t o s u p p o rt im a g e m a n a g e m e n t fo r c o m p u t e s e rv ic e s .
• T h e im a g e re p o s it o ry h a s t w o k e y c o m p o n e n t s :
• Image Registration and Discovery Service:
• A llo w s u s e rs t o re g is t e r, c a t a lo g , a n d re trie v e V M d i sk im a g e s .
• S u p p o rts m e ta d a t a t a g g in g fo r e a sie r s e a rc h a n d f ilt e rin g .
• Image Delivery Service:
• E n s u re s e f fic ie n t d is trib u t io n o f V M i m a g e s fro m t h e s t o ra g e s y st e m t o th e c o m p u t e n o d e s.
• H a n d le s im a g e fe t c h in g , c a c h in g , a n d p ro v is io n in g .
• T h e s e s e rv ic e s w o rk t o g e t h e r t o e n a b le th e a u t o m a t e d d e p lo y m e n t o f v irt u a l m a c h in e s w it h sp e c i fic im a g e c o n fig u ra t io n s .
• Integration and Scalability Vision:
• O p e n S t a c k is m o d u la r a n d e x t e n si b le , a llo w in g c o m p o n e n t s t o o p e ra t e in d e p e n d e n t ly o r in c o n ju n c t io n .
• N e w s e rv ic e s l ike im a g e m a n a g e m e n t in d ic a te a b ro a d e r g o a l o f s e rv ic e in t e g ra t io n , in c lu d in g n e t w o rk in g (N e u tro n ),
d a s h b o a rd (H o riz o n ), id e n ti ty ( K e y st o n e ), a n d m o re .
• T h e o v e ra ll p ro je c t a im s t o b u ild a s c a la b le c lo u d O S c a p a b le o f h a n d l in g b o t h in fra s t ru c t u re -le v e l v irt u a liz a t io n a n d
s e rv ic e -le v e l o rc h e s tra t io n .