Adap i e Scheduling Ac oss a Dis ibu ed
Compu a ion Pla o m
And ew Page, Thomas Keane and Thomas J. Naugh on
Depa men o Compu e Science,
Na ional Uni e si y o I eland,
Maynoo h, I eland.
Email: and [email protected], [email p o ec ed].ie, om.naugh [email protected]
Abs ac — A p og ammable Ja a dis ibu ed sys em, which
adap s o a ailable esou ces, has been de eloped o minimise he
o e all p ocessing ime o compu a ionally in ensi e p oblems.
The sys em exploi s he ee esou ces o a he e ogeneous se o
compu e s linked oge he by a ne wo k, communica ing using
SUN Mic osys ems’ Remo e Me hod In oca ion and Ja a socke s.
I uses a mul i- ie ed dis ibu ed sys em model, which in p incipal
allows o a sys em o unbounded size. The sys em consis s o
an n-a y ee o nodes whe e he in e nal nodes pe o m he
scheduling and he lea es do he p ocessing. The schedule nodes
communica e in a pee - o-pee manne and he p ocessing nodes
ope a e in a s ic ly clien -se e manne wi h hei espec i e
schedule . The independen schedule s on each ie o he ee
dynamically alloca e esou ces be ween p oblems based on he
cons an ly changing cha ac e is ics o he unde lying ne wo k.
The sys em has been e alua ed o e a ne wo k o 86 PCs
wi h a bioin o ma ics applica ion and he a elling salesman
op imisa ion p oblem.
I. INTRODUCTION
The compu a ional demands o mode n scien i ic esea ch
ha e been he d i ing o ce behind dis ibu ed compu ing [13],
p o iding la ge compu a ional esou ces cos e ec i ely. Cu -
en ly he e a e a ew no able dis ibu ed compu ing pla -
o ms such as SETI@home [3], Uni ed De ices [30], and
dis ibu ed.ne [h p://www.dis ibu ed.ne ], which ha e been
buil o y and sa is y he wo lds inc easing need o compu a-
ional powe , wi hou he adi ional high cos s associa ed wi h
dedica ed pa allel ha dwa e and clus e s. They wo k on he
p inciple o a use dona ing hei machine’s spa e clock cycles
o he sys em ac oss an in ane o he In e ne , so ha i s ee
esou ces can help o p ocess compu a ionally la ge p oblems.
The widesp ead success o he In e ne has mean ha hese
dis ibu ed sys ems ha e been able o ha ness la ge amoun s
o compu a ional esou ces om dono s’ machines, which
would o he wise ha e no been u ilised o hei ull po en ial.
These sys ems a e gene ally e e ed o as In e ne compu ing
sys ems, whe e hei esou ces a e massi ely dis ibu ed ac oss
he In e ne , and he p oblems hey a emp a e gene ally
i ially pa allelisable.
The e a e howe e p oblems wi h exis ing dis ibu ed sys-
ems which can be add essed h ough he use o Ja a. Many
sys ems [2], [3], [21], [29], [30] a e based on languages ha
a e no pla o m independen , such as C and Fo an, esul ing
in he equi emen o ha e mul iple e sions o he so wa e,
hus incu ing highe main enance and de elopmen cos s. Ja a
is a pla o m independen language, which allows by e-code
o be gene a ed which will un on a la ge a ie y o di e en
pla o ms. The pe o mance o Ja a is also compa able o
languages which use machine na i e code [8]. Many exis ing
sys ems a e limi ed o i ially pa allelisable p oblems, and a e
ha d coded o pe o m only a single ask [3], [21], [26]. Ja a
allows new classes o be loaded o upda ed whils a p og am
is unning, allowing o a p og ammable sys em o be easily
c ea ed.
Secu i y mus be conside ed in dis ibu ed sys ems due o
he use o insecu e ne wo ks, such as he In e ne , because
dono s pu hei machines unde he con ol o o he s. Secu i y
is essen ial o p o ec he in eg i y o he esul s ob ained om
dis ibu ed sys ems. Many exis ing sys ems use languages ha
make he implemen a ion o secu i y di icul whe eas Ja a has
mul iple secu i y o ien a ed API’s a i s co e, such as he Ja a
C yp og aphic A chi ec u e, a ailable o he p og amme o
implemen a secu e sys em wi hou he need o unde s and he
unde lying wo kings o he c yp og aphy. Digi ally signed JAR
iles, secu i y policies, and p og am execu ion in a sandbox
p omo e con idence in any p og ams unning, p o ec esul s
om in e e ence, and p o ec he dono ’s machine om
ha m ul damage.
Ano he p oblem wi h many exis ing dis ibu ed sys ems [2],
[3], [21], [30] is ha hey a e no ex ensible enough, and ha e
undamen al limi s on hei scalabili y. They use models such
as he single- ie clien se e model [11], which undamen-
ally limi s he size o he sys em. In p ac ice his model has
se ed SETI@home e y well, wi h up o ou million clien
machines as pa o he sys em [3]. Bu since he e is only
one se e (single machine o clus e ) o all o he clien s
he e is hus a ini e limi on he numbe o clien s he sys em
can handle a any one ime, wi h his limi depending on he
ne wo k esou ces and compu a ional esou ces a he se e . A
common solu ion is o inc ease he bandwid h o he se e ’s
In e ne connec ion and o upg ade he powe o he se e ,
bu his can be expensi e. Ano he solu ion, and one adop ed
by SETI@home, is o pa allelise he compu a ion a such a
coa se le el ha clien s ( ela i ely) in equen ly e u n o he
se e o mo e da a uni s. This ac ic, howe e , is only sui able
o pa icula p oblems and migh no be applied success ully
o he a bi a y p oblems o a gene al-pu pose p og ammable
sys em.
The mos commonly used communica ion echnologies o
pa allel compu ing, such as PVM and MPI, make he ask
o using mo e complex and unbounded models ha de o
implemen due o he equi emen o mo e low-le el de elop-
men . Ja a howe e allows o he de elopmen o dis ibu ed
sys ems a a high le el, h ough he use o RMI and Ja a
socke s, allowing o he c ea ion o mo e ex ensible sys ems
wi hou a co esponding inc eased in he complexi y o he
use s ask.
Many sys ems ha e been de eloped o a emp o add ess he
limi a ions o p e ious dis ibu ed sys ems. The Be keley Open
In as uc u e Ne wo king Compu ing (BOINC) [2] sys em is
a p og ammable successo o SETI@home, and a emp s o
make a mo e gene alisable sys em. Al hough i is p og amm-
able, only i ially pa allelisable p oblems a e conside ed. Also
BOINC only conside s p oblems which will be appealing
enough o ge la ge numbe s o use s ac oss he In e ne o
dona e hei ee esou ces o he p ojec . I s ex ensibili y
is also limi ed because i e ains a clien -se e a chi ec u e,
and implemen s a one-s ep p ocessing s age. I a compu a ion
equi es u he p ocessing o in e media e esul s, sepa a e
dedica ed machines mus be used. Uni ed De ices [30] p o ide
hei p og ammable dis ibu ed sys em on a comme cial basis,
wi h appealing In e ne compu ing p ojec s p ima ily being
used o p omo e hei comme cial dis ibu ed sys em so wa e.
Thei sys em is also limi ed by he use o pla o m dependen
na i e by e-code, and hei use o he clien -se e model limi s
he ex ensibili y o hei sys em. The p og ammable dis ibu ed
sys ems om K iege and V iend [21], and Sil es e e al. [26]
ha e simila aims o ou s, bu su e om he same scalabili y
p oblems as ou lined o p e ious clien -se e sys ems. In
addi ion, he na i e code used in hei sys em leads o pla o m-
dependence and se ious secu i y conce ns, bo h o which a e
alle ia ed in ou sys em h ough he use o Ja a.
A common limi a ion wi h he gene alisable sys ems we
e iewed [1], [2], [5], [21], [22], [29], [30] is he ac ha i is
no possible o un any compu a ion on any clien . I he co ec
clien (i.e. ope a ing sys em and a chi ec u e) is no a ailable
o a pa icula compu a ion, hen he compu a ion will ne e
ge un on hese sys ems. Almos e e y ope a ing sys em and
ha dwa e a chi ec u e suppo s a JVM, om desk op PCs o
mobile phones, and i is o ally pla o m independen .
A numbe o o he dis ibu ed sys ems implemen ed in Ja a
exis [1], [22], [26], [28]. Wo k by Ai-Ja oodi e al. [1],
using Ja a o c ea e a dis ibu ed sys em, is simila in many
espec s o he sys em desc ibed in his pape , al hough hey
ely on a dis ibu ed sha ed memo y model, which causes
signi ican o e heads in e ms o synch onisa ion, and con-
sis ency o da a. Ou use o RMI and Ja a socke s educes he
o e heads incu ed, and allows o a mo e scalable sys em.
Sil es e e al. [26] ha e c ea ed a dis ibu ed sys em in
Ja a o p ocess a bioin o ma ics p oblem, bu hei sys em
is no p og ammable. Su deanu and Moldo an [28] c ea ed a
dis ibu ed JVM, p o iding a pla o m o sequen ial mul i-
h eaded p og ams o be p ocessed by mul iple p ocesso s.
Thei ocus is on speeding up exis ing p og ams, a he han
using compu a ional esou ces e icien ly.
Ou aim is o design a p og ammable dis ibu ed compu ing
pla o m ha is un es ic ed in e ms o he ype o s uc u e o
compu a ions ha can be pe o med. A schedule , which can
adap o he e e changing esou ces a ailable o he sys em in
a he e ogeneous compu ing en i onmen , is equi ed o allow
mul iple di e en p oblems o be p ocessed simul aneously.
We e ain aspec s o he clien -se e model, bu in oduce
pee - o-pee communica ion wi hin a ee o scheduling nodes
ha se es o o e come he scalabili y and ex ensibili y limi-
a ions o employing a single se e . Ja a is used o ensu e
pla o m independence o bo h schedule nodes and p o-
cessing nodes. We ha e applied ou dis ibu ed compu ing
pla o m o p oblems om he ield o bioin o ma ics. The
so wa e we ha e de eloped is open sou ce and a ailable unde
he GNU GPL license ee o cha ge, and is no limi ed o
i ially pa allelisible p oblems, due o i s abili y o handle
message passing be ween di e en pa s o compu a ions and
he implemen a ion o a pipeline p ocesso .
The es o he pape is o ganised as ollows. In Sec . II, we
in oduce he mul i- ie ed model. In Sec . III, he designs o he
main componen s o he sys em a e p esen ed. Implemen a ion
and pe o mance e alua ion a e discussed in Sec . IV, and we
conclude in Sec . V.
II. OVERVIEW OF THE SYSTEM
The ounda ions o he mul i- ie ed dis ibu ed compu a ion
sys em we e laid in he Ja a Dis ibu ed Compu a ion Lib a y
(JDCL) [15] and i s ex ensions [20], which p o ided MIMD
capabili ies h ough emula ed pipeline p ocesso s. The JDCL
p o ided a simple clien -se e based de elopmen pla o m
o de elope s who wished o quickly implemen a dis ibu ed
compu a ion sys em. I a ose ou o he need o a pla o m-
independen dis ibu ed sys em ha was simple o c ea e,
adap ed o sys em changes, and was simple o deploy. Sys ems
such as SETI@home did no add ess hese issues e y well
and we e designed o be pla o m dependan and o a single
pu pose only. The JDCL does, howe e , su e om simila
scalabili y p oblems o hose o SETI@home in ha i has
one se e (single machine o clus e ). The design o he
cu en mul i- ie ed sys em aims o add ess his undamen al
limi a ion. An adap i e schedule has also been de eloped o
a emp o minimise he o e all compu a ion ime o p oblems
p ocessed by he sys em, whils also ma ching p oblems o
dono machines, and allowing he use o p io i ise p oblems.
A. Mul i- ie ed Model
The mul i- ie ed dis ibu ed compu ing sys em was c ea ed
wi h he in en ion o u ilising he maximum compu a ional
esou ces o he machines unde i s con ol, while no placing
any limi a ions on he maximum size o he sys em. The
clien -se e model alone was no su icien , so a hyb id
model was c ea ed ha combines he ad an ages o pee -
o-pee and clien -se e a chi ec u es wi hin he one model.
The sys em consis s o an n-a y ee o nodes whe e he
in e nal nodes pe o m he scheduling and he lea nodes
do he p ocessing, as shown in Fig. 1. The schedule nodes
communica e in a pee - o-pee ashion, which pe mi s op-
down econ igu abili y, ex ensibili y, and ee balancing. The
p ocessing nodes ope a e in a s ic ly clien -se e ashion
wi h hei espec i e schedule , which p omo es dono s’ us
in he sys em (anonymi y and secu i y) and admi s a simple
and obus design, h ough he use o Ja a RMI and he Ja a
C yp og aphy A chi ec u e API’s.
Scheduling node
P ocessing node
Remo e in e ace
Fig. 1. Example opology o mul i- ie ed dis ibu ed compu ing sys em.
The po en ial size o he dis ibu ed sys em, in wid h and
dep h, is unlimi ed in p inciple due o he abili y o he sys em
o dynamically add ano he scheduling node, ie o scheduling
nodes, o p ocessing node. Also, since he scheduling nodes
a e dis ibu ed, he po en ial o bo lenecks is educed, hus
imp o ing he pe o mance o he sys em. The dis ibu ed
na u e o he scheduling nodes means ha i one (o he han
he oo ) we e o ail, he es o he sys em could always
compensa e o he loss o ha b anch. By ha ing mul iple
scheduling nodes i means ha he dis ibu ed sys em is a
MIMD a chi ec u e. Al hough no an inc ease in capabili y,
going om a pu e MISD clien -se e amewo k o a MIMD
amewo k does inc ease he sophis ica ion o he algo i hms
ha can be dis ibu ed o e he sys em, and he e o e uns he
isk o inc easing he use ’s ask in p og amming a dis ibu ed
compu a ion. We ha e a emp ed o s uc u e he p og amme s’
in e ace as much as possible in his ega d, o ind a balance
be ween exp essi eness and simplici y. Fo example, he p o-
g amme is equi ed only o ex end wo classes o ully speci y
a mul i- ie ed dis ibu ed compu a ion, as explained la e in
Sec . III.
III. DESIGN
The e a e h ee dis inc pa s o he sys em. These a e
he clien (p ocessing node), he se e (scheduling node),
and he emo e in e ace. The use is equi ed o ex end wo
classes o c ea e a p oblem o un on he sys em, wi h he
use o Ja a enabling high le el abs ac design compa ed o
o he mo e low le el echnologies such as C and Fo an [8].
The Da aManage class (in he schedule ) speci ies how
he p oblem is o be pa i ioned in o uni s o wo k and he
in e media e esul s pu oge he , acili a ing he compu a ion
o mo e gene alisable p oblems, a he han being limi ed
o i ially pa allelisable p oblems [2], [3], [21], [26], [30].
The Algo i hm class (in he clien ) speci ies he ac ual
compu a ion. Figu e 2 depic s how a p oblem is ea ed wi hin
he sys em.
Addi ional Classes
P oblem
P oblem Da a
Da aManage
Algo i hm
Da aManage Da aManage …..
Queue Schedule
Communica ions Logs Secu i y
Se e
Cache
Clien
P oblem P ocesso
Secu i y
Communica ions
Resul s
Remo e In e ace
Communica ions
Secu i y
Ac ions
Resul s
Resul s
Use
Resul s
Fig. 2. How a p oblem is ea ed wi hin he sys em.
The use s o he sys em do no need any knowledge o he
opology o wo kings o he sys em in o de o submi p ob-
lems and ge hei p ocessed esul s back. They jus p o ide a
Da aManage , an Algo i hm, addi ional equi ed classes,
da a o be p ocessed (i equi ed), and an in ege p io i y
weigh ing, o c ea e a sel con ained P oblem objec , which
hen ge s p opaga ed o all se e s in he sys em. A P oblem
objec p o ides a p ede ined in e ace o he unc ionali y ha
he use has p o ided in hei Da aManage . They can also
associa e minimum CPU and memo y limi s, o enable he
schedule o ma ch p oblems o clien machines wi h su icien
ee esou ces (see Sec . III-E.2).
Communica ion wi hin he sys em is based on a combi-
na ion o Ja a RMI [27] and o dina y Ja a socke s. RMI is
a buil -in acili y in Ja a ha allows one o in e ac wi h
objec s ha a e ac ually unning in Ja a Vi ual Machines
(JVM’s) on emo e hos s on a ne wo k, and a oids he need o
he designe o wo y abou low le el communica ion issues.
By also using Ja a socke s o he ansmission o p oblem
da a iles, la ge amoun s o da a can be s eamed di ec ly
om he ha d disk wi hou b inging he la ge iles comple ely
in o memo y, hus educing he memo y equi emen , and
imp o ing he scalabili y and pe o mance o he sys em.
A. Clien
The main pu pose o he clien (p ocessing node) so wa e
is o con inually eques and p ocess da a uni s om i s se e
(one o he scheduling nodes). All communica ion is ini ia ed
by he clien , p o iding anonymi y and imp o ed secu i y.
Once he clien is s a ed, he clien a emp s o connec o he
se e , and eques s a da a uni o p ocess. Once i has ecei ed
he da a uni , i checks i i al eady has he equi ed algo i hm
in i s Algo i hm cache, and downloads i om he se e
i necessa y. The clien hen b eaks i s connec ion wi h he
se e and will no con ac i again un il he uni is p ocessed,
he maximum p ocessing ime is exceeded o he uni and i
needs o eques an ex ension, o an excep ion occu s. Hence
he se e is always kep in o med as o he s a us o he uni
being p ocessed by he clien . I an ex ension eques is no
ecei ed by he se e , i is assumed ha he clien has been
e mina ed (dono has swi ched o his/he machine) and he
se e edis ibu es he da a uni o ano he subo dina e node.
Finally, i a any ime he clien canno con ac he se e ,
i will go in o a ‘sleep mode’ and a emp o connec o he
se e pe iodically.
B. Remo e In e ace
The emo e in e ace is a s and-alone applica ion ha com-
munica es o e TCP/IP, and allows he adminis a o o ully
con ol all o he p oblems and scheduling nodes in he sys em.
We use Ja a RMI communica ions echnology o message
passing be ween he emo e in e ace and he se e and Ja a
Socke s o la ge da a iles, which can w i e da a di ec ly o
disk. P oblem objec s can be added, emo ed, con igu ed,
and hei cu en s a e iewed. The emo e in e ace can
connec and disconnec om se e s, o allow hem o be
upda ed, que ied, paused, and shu down wi hou a ec ing he
dis ibu ed sys em. The emo e in e ace also allows he owne
o a pa icula p oblem o iew i s cu en s a e and download
he esul s and log iles, which a e comp essed in o a single
a chi e using Ja a’s in buil GZip API o educe he bandwid h
equi ed.
C. Se e
The se e (scheduling node) is he engine o he en i e dis-
ibu ed sys em, con olling subo dina e se e s (subo dina e
scheduling nodes) and clien s (p ocessing nodes) and is la gely
desc ibed in [20]. We ha e modi ied his sys em in o de o
allow i o be a mul i- ie ed dis ibu ed sys em. A se e can
be added o he dis ibu ed sys em while he sys em is unning,
and likewise can be emo ed om he sys em wi hou a ec ing
he s abili y o he sys em o causing unning p oblems o
become co up ed ( he oo se e begin he excep ion). A new
se e con ac s he se e abo e i and synch onises i sel wi h
he es o he sys em. I a se e is emo ed, o ails, he es
o he sys em compensa es o his loss. The los p oblems a e
ealloca ed by he se e , abo e he ailing one, o o he s a he
same le el as he ailing one. This esilience o ailu e ensu es
he sys em can pe o m in unp edic able ope a ing condi ions.
Uni s o wo k which a e handed ou by a Schedule con-
ained wi hin he Se e , o subo dina e nodes, a e cached
on he ha d disk using a DiskHashTable we de eloped.
The buil -in hash able in Ja a s o es e e y hing in memo y,
educing scalabili y. The a ailabili y o cheap la ge ha d disk
s o age space allows o a much mo e scalable se e and
sys em as a whole, by only s o ing he uni s being used in
memo y.
D. Schedule
The e is a schedule in each scheduling node (se e ). In
he mul i- ie ed dis ibu ed sys em he schedule was c ea ed
o allow a limi less numbe (memo y limi s aside) o la ge
scale dis ibu ed compu a ions o un in pa allel on he sys em
and o allow he sys em o adap o he changing condi ions o
he ne wo k and compu a ional esou ces a ailable. Signi ican
op imisa ions ha e been achie ed by using adap i e scheduling
in dis ibu ed sys ems [23], compa ed o s a ic non-adap i e
scheduling algo i hms. The in o ma ion used by he schedule
o adap i s in e nal scheduling mechanism is collec ed pas-
si ely, so he schedule can only ga he in o ma ion when a
clien p esen s i . This is opposed o o he sys ems which adap
by p oac i ely sensing sys em condi ions [5], [7], [18], [32] o
by unning pla o m dependan hi d pa y p og ams [5], [32]
which cause addi ional o e heads o be incu ed.
When a schedule ecei es a new P oblem objec , i adds
i o he end o i s queue o p ocessing p oblems. Subo dina e
nodes eques uni s o wo k om he schedule ( ia he se e ),
which a e gene a ed by ins ances o he use -de ined Da a-
Manage [20] encapsula ed wi hin he P oblem objec s in
he queue. Subo dina e clien s and se e s bo h ecei e he
same uni s when hey eques mo e wo k, al hough he clien
p ocesses he uni o wo k while he se e b eaks i up u he
in o mo e subuni s, c ea ing a new P oblem objec o each
subuni . The p oblem is ecu si ely di ided (see Sec . III-
E.3) in o P oblem objec s in his manne by ins ances o
he Da aManage . P ocessed esul s a e sen back o he
Da aManage in he se e abo e, and combined in a simila
manne . When a P oblem is inished p ocessing i is emo ed
om he queue o cu en P oblems. A he oo , he inal
esul s and log iles a e w i en o iles and comp essed o
collec ion by he use h ough he emo e in e ace.
I he p ocessed esul s o a uni sen by a se e a e no
e u ned wi hin a speci ic pe iod, he uni is said o expi e and
he se e esends i o he nex eques ing se e /clien . Each
uni also con ains i s own ime so ha i he uni is close o
expi ing he clien , o se e , eques s an ex ension om he
se e abo e.
E. Scheduling Algo i hm
A schedule has been designed which aims o adap o he
changing esou ces a ailable o a dis ibu ed sys em and o
minimise he p ocessing ime o compu a ions. The schedule
akes in o accoun use speci ied p io i ies o each p oblem,
a emp s o e icien ly alloca e and manage he a ailable sys-
em p ocessing powe o e ed by he se o dono machines,
and dynamically al e s he g anula i y o da a uni s dis ibu ed.
The e icien alloca ion o sys em p ocessing powe is
loosely based on he ‘ma chmaking’ ea u e in Condo [24]
(see Sec . III-E.2). The dynamic adap a ion o he g anula i y
o da a uni s ha a e issued by he indi idual p oblems was
inspi ed by he dynamic window size o he TCP p o ocol [25]
(see Sec . III-E.3). The o e all scheduling s a egy is b ough
oge he in he use -de ined ai ma chmaking scheduling
mechanism d awing om esea ch in [4], [16], [18], [19], [24]
(see Sec . III-E.1).
1) Use -de ined Fai Ma chmaking Schedule : Each
scheduling node in he sys em has i s own schedule , which is
independen o all o he schedule s. The e o e, each sub- ee
in he sys em has he po en ial o sel -op imisa ion and has
he po en ial o adap o he cons an ly changing condi ions
o i s own esou ces. The schedule decides dynamically
how he esou ces o he sys em a e o be di ided among
he p oblems in i s queue, employing s a egies based on a
use -de ined ai ma chmaking scheduling mechanism [4],
[16], [18], [19], [24]. The use se s an in ege weigh ing Pi
o each p oblem i∈ {1,2,...,N}in he sys em. The alue
˜
Px=PxN
P
i=1
Pi−1
deno es he no malized weigh ing o
p oblem xo e all Np oblems in he sys em.
As p ocessed da a uni s a e e u ned o he schedule o
p oblem x, he ime aken (in seconds) o p ocess he uni τis
no ed and inco po a ed in o he ypical uni p ocessing ime x
o ha p oblem. Ra he han a simple a e age, xis calcula ed
om x= 0
x+L(τ− 0
x), whe e 0
xis he p e ious alue o x.
The lea ning a e L∈[0,1] de ines he in luence o p e ious
alues, wi h he in luence o olde alues ending owa ds
ze o o e ime. This echnique is bo owed om machine
lea ning [4], and allows scheduling nodes o con inuously
adap o he cons an ly changing condi ions encoun e ed by
he sys em. The alue ˜
x= xN
P
i=1
i−1
is he no malized
ypical p ocessing ime o he da a uni s o p oblem x.
The schedule ies o ai ly alloca e he pa allel compu a ion
ime esou ces by adop ing i s own in e nal weigh ing Fx
p opo ional o Pxand in e sely p opo ional o x, gi en
by Fx=c˜
Px/˜
x, c ∈R, o each p oblem x.I is his
weigh ing Fx ha ul ima ely de ines he p io i y con e ed on
each p oblem xin he sys em in o de o ai ly alloca e he
pa allel compu a ion ime.
2) Alloca ion and Managemen o P ocessing Powe : As
he i le o his sec ion sugges s, he e a e wo main aims
he e. Fi s ly he sys em a emp s o balance he compu a ional
equi emen s o he p oblems in he sys em wi h he capaci y
o he indi idual dono machines. The second aim is o
e icien ly alloca e and manage he o al amoun o a ailable
sys em p ocessing powe be ween he se o p oblems in he
sys em. The e a e wo di e en se s o inpu s o he scheduling
algo i hm. The i s ype o inpu occu s when a clien makes
a eques o a da a uni (supplying i s CPU speed and amoun
o memo y a ailable on he dono machine). The o he ype
o inpu occu s when a clien e u ns a se o esul s. In his
case he inpu consis s o he p oblem ID and how long he
uni ook o comple e.
I a clien is making a eques o a da a uni , hen he
schedule algo i hm goes h ough i s queue o p oblems and
que ies each p oblem o see i he dono machine mee s i s
minimum equi emen s. When a sui able p oblem is ound ha
is eady o issue a da a uni ,such a da a uni is e u ned o he
clien . I no sui able p oblem is ound in he queue, hen he
clien is sen a message o sleep and e y la e .
I he inpu consis s o a esul s se , hen he schedule
eco ds how long he uni ook o be p ocessed. A e e e y
esul is ecei ed, he schedule calcula es wha we call a
“se icing alue” (de ined below) o each p oblem and eso s
he queue o p oblems by hei se icing alue in descending
o de . The se icing alue, Si∈[−1,1], is a a io o how much
sys em p ocessing ime each p oblem has ecei ed compa ed
o he p io i y o he p oblem. A se icing alue Si<0means
ha p oblem iis o e -se iced, Si>0means ha p oblem
iis unde -se iced, and Si= 0 means p oblem iis cu en ly
being se iced app op ia ely.
The se icing alue o each p oblem iis gi en by
Si=
Pi
N
P
j=1
Pj
−
Fi×Ui
N
P
j=1
(Fj×Uj)
, i ∈
{1,2,...,N},
N
P
i=1
Si= 0,
whe e Uiis he numbe o uni s handed ou so a o
p oblem i, and whe e all o he a iables we e explained in
Sec . III-E.1. P oblems ha ha e ecei ed he leas p ocessing
ime ela i e o hei p io i y will be p omo ed o he op o
he queue and will be alloca ed mo e sys em ime.
3) Recu si ely Spli ing up a P oblem: Any p oblem ha
is un on ou dis ibu ed sys em mus be pa allelised by he
use . The use speci ies how he p oblem can be pa i ioned
in o uni s o wo k in he Da aManage hey p o ide. Fo ex-
ample, his could be pai s o indices speci ying subsequences
o a genome o be analysed. They speci y wha he smalles
possible uni o wo k can be o he p oblem, also called
g anula i y, allowing he schedule o c ea e uni s o wo k
which a e mul iples o his a omic g ain size. This pa i ioning
happens a each ie in he sys em, wi h he p oblem being
b oken up in o smalle pa s a each ie . A ca e ully chosen
g anula i y can gi e signi ican pe o mance inc eases [10].
Gi en a p oblem, he maximum speedup achie able is
limi ed by he numbe o a omic uni s ha he p oblem can
be b oken in o, and he numbe o clien s in he sys em.
Techniques bo owed om machine lea ning [4] allow us o
dynamically adjus he pa i ioning o da a uni s acco ding o
he e e -changing compu a ional esou ces a he disposal o
he sys em.
E en hough use s speci y he minimum equi emen s
(memo y, p ocesso speed) o hei p oblem when hey en e
hei p oblem in o he sys em, i is no possible o ha e a p io i
knowledge o he capaci y o he ne wo k and he powe o he
dono machines. The e o e he app op ia e g anula i y o he
pa allel compu a ion will ha e o be asce ained dynamically.
Wi h oo ine a pa allelism, clien s will e u n esul s a e a
e y sho p ocessing ime and migh o e load he ne wo k.
Too coa se a pa allelism migh esul in some p ocesso s being
le idle and could cause la ge amoun s o p ocessing ime o
be was ed i a dono machine is unexpec edly swi ched o .
A e e e y xuni s a e ecei ed by a p oblem (de aul alue o
xis 50), ou sys em ies o modi y he p oblem’s g anula i y
so ha he a e age ime a o p ocess a uni will app oach some
a ge . I ais no su icien ly close o , hen an a emp is
made o al e he g anula i y o subsequen uni s. The ac ion
o a( ep esen ed by d) by which he g anula i y needs o be
al e ed is calcula ed om
d=( −a
a,|( −a
a)×100|>
0,o he wise
whe e is he a ge ime, a=1
x
N
P
i=1
Uiis he a e age p o-
cessing ime o he p e ious xuni imes Ui,Nis he numbe
o p oblems in he schedule , and is he pe cen age a iance
h eshold below which he p ocessing imes a e allowed o
luc ua e. The de aul alues o and a e 15% and one hou ,
espec i ely. A posi i e alue o dindica es ha he p oblem’s
g anula i y should be inc eased and a nega i e alue means
ha he p oblem’s g anula i y should be dec eased. The alue
o dis sen o he use ’s Da aManage and i is up o his
code o ake he app op ia e ac ions o al e he g anula i y o
subsequen uni s. Choosing lowe ( espec i ely, highe ) alues
o xand will allow he sys em o adap mo e ( espec i ely,
less) quickly o changes in he ne wo k, bu will cause i o be
mo e ( espec i ely, less) likely o o e eac in he p esence o
ansien luc ua ions in ne wo k capaci y and clien p ocesso
powe .
F. Sys em Secu i y
Secu i y should be a e y impo an aspec o any dis ibu ed
sys em. The essence o ou sys em is ha is akes an a bi a y
piece o code and sends i ou o a clien o be dynamically
loaded. This p inciple alone h ows up many implica ions o
he secu i y o he dono machine, and i is a majo p oblem
wi h all exis ing dis ibu ed sys ems ha do no use Ja a [1],
[2], [3], [7], [21], [29], [30]. The JVM allows us o include
a comp ehensi e secu i y manage in he clien so wa e. The
downloaded code is es ic ed in he ope a ions i can pe o m,
by a secu i y policy, which is a s anda d ea u e o Ja a. I a
secu i y excep ion is de ec ed, he code is immedia ely s opped
and he se e is no i ied o he b each in secu i y. In he
same manne , any secu i y excep ions on he se e will esul
in he use ’s p oblem code being ejec ed om he sys em.
One o he impo an secu i y mechanism in ou so wa e is
ha all o he sys em JAR iles a e digi ally signed using
he “MD5wi hRSA” algo i hm, which is included as s anda d
in Ja a. I ou so wa e is ampe ed wi h in any way, o an
unau ho ised clien is downloaded, i will no un.
IV. IMPLEMENTATION
The en i e sys em was designed using objec s a a high
le el, hus, when i came o implemen a ion, Ja a was chosen
because o i s objec o ien a ed capabili ies and also o
i s pla o m-independence and ease o implemen a ion. Ja a
p og ams a e compiled o an a chi ec u e neu al by e-code
o ma , he e o e a Ja a applica ion can un on any sys em, as
long as ha sys em implemen s he JVM [14].
The compu a ional esou ces a ailable o his sys em’s
deploymen we e based on machines wi h a ying ope a -
ing sys ems, such as Mic oso Windows 98/2000/NT, and
Linux dis ibu ions Redha 7.2/Debian/Fedo a Co e 1/Man-
d ake 9.2/Gen oo 1.4, and a ying ha dwa e a chi ec u es such
as hose o In el, SUN, and HP.
A. Applica ions
The mul i- ie ed dis ibu ed sys em has so a been used
o analyse ube culosis and E-Coli genomes sea ching o
duplica ed pa e ns, o b eak c yp og aphic schemes based
on he disc e e loga i hm p oblem, such as ElGamal [12],
using a dis ibu ed Polla d-Rho algo i hm, and has been used
o model ligh p opaga ion in issue using he Mon e Ca lo
me hod. All o hese p oblems a e i ially pa allelisable, so o
show he gene alisabili y o he sys em, he a elling salesman
op imisa ion p oblem was also p ocessed using he sys em.
B. Pe o mance E alua ion
We ca ied ou pe o mance es s o e alua e he capabili ies
o he mul i- ie ed dis ibu ed sys em. We se ou o show
ha adding mo e se e s would inc ease he capaci y o he
dis ibu ed sys em, and we also se ou o show ha he
scheduling algo i hm would a emp o op imally balance he
compu a ional esou ces a i s disposal. The capaci y o he
sys em is he maximum numbe o p ocessing nodes in he
sys em, whe e he addi ion mo e o p ocessing nodes would
educe he pe o mance o he sys em. These expe imen s we e
ca ied ou in a compu e labo a o y wi h a dedica ed ne wo k
o 86 PCs. Each had a 600 MHz Pen ium III p ocesso wi h
128 MB o RAM and 20 GB o ha d disk space and was
connec ed o a 10 Mb/s E he ne LAN.
Figu e 3 shows how balanced he schedule is du ing ope -
a ion. We eco ded he sum o he absolu e alues o all o he
se icing alues,
N
P
i=1
|Si|, whe e Nis he numbe o p oblems
in he sys em, e e y ime a p ocessed uni was e u ned o he
schedule o demons a e how balanced he se o p oblems
was a each poin in ime. A alue o ze o indica es a
balanced schedule and is he op imal alue, wi h highe alues
indica ing he deg ee o which he schedule is unbalanced.
Ini ially i e p oblems we e added o he sys em, wi h wo
addi ional p oblems being added on h ee sepa a e occasions,
co esponding o he peaks in he g aph. The p oblems used o
he es we e ins ances o he a elling salesman op imisa ion
p oblem, and each p oblem had an equal p io i y o 1. The
schedule quickly mo es o y and educe he imbalance in
he schedule , hus ending owa ds a balanced sys em.
Figu e 4 shows he speedup ha was achie ed. To show
how gene alisable he sys em is, by allowing message passing
be ween in e media e esul s in he scheduling nodes, we again
used he a elling salesman op imisa ion p oblem. Based on
he speedup da a, he a e age e iciency using 86 p ocesso s
is 95.6%. In Fig. 4 he speedup was calcula ed om he
equa ion S(n) = P1
Pn, whe e Pndeno es he p ocessing ime
o np ocesso s and P1deno es he p ocessing ime o one
0 1000 2000 3000 4000 5000 6000 7000 8000 9000
0
0.2
0.4
0.6
0.8
1
1.2
1.4
1.6
1.8
2
No. o Uni s P ocessed
Imbalance in he se icing alues
Fig. 3. Expe imen al esul s showing he schedule dynamically balancing
he sys em as new jobs a e added.
p ocesso , wi h each poin on he g aph being he a e age o
a minimum o 5 uns o he expe imen .
0 10 20 30 40 50 60 70 80 90
0
10
20
30
40
50
60
70
80
90
Speedup
Nume o P ocesso s
Linea speedup ( heo e ical uppe bound)
Measu ed speedup
Fig. 4. Speedup achie ed wi h a a elling salesman op imisa ion algo i hm.
Since a se e can handle in he o de o housands o
clien s, and such a numbe o machines was no a ailable o us,
we had o simula e conges ion o demons a e ha he mul i-
ie ed dis ibu ed sys em can inc ease he numbe o clien s
in a sys em compa ed o he adi ional clien -se e opol-
ogy. This was achie ed by using Shun a’s eewa e Nimbus
bandwid h h o ling so wa e [h p://www.shun a.com]. The
bandwid h o each se e was cons ic ed o only 14.4 kb/s,
and in addi ion, each uni e u ned om a se e o clien was
bloa ed wi h ex a da a o make he e u ned esul s la ge . The
size o his bloa ed da a was p opo ional o he size o he
uni . The p oblem posed o he dis ibu ed sys em was a pa e n
ma ching exe cise wi h he ube culosis genome ( ou million
nucleo ides in leng h) o ind all duplica ed s ings wi hin he
genome.
The p oblem was ini ially un wi h 1 se e and nclien s,
a one- ie adi ional clien -se e model o se e al alues o
n. Nex we an he p oblem using 5 se e s, a anged as 1
se e in he op ie and 4 se e s in he nex ie , wi h n
clien s. Figu e 5 shows he esul ing plo o se e al alues
o n. This plo shows ha a e a ce ain numbe o clien s,
ne wo k conges ion a he se e causes he p ocessing ime o
ac ually inc ease (app oxima ely 20 clien s, wi h conges ion,
in he case o adi ional 1-se e opology). As he numbe
o se e s inc eases, his c i ical numbe can be inc eased.
E en ually, he mul i- ie ed sys em oo eaches i s capaci y
bu i can handle many mo e clien s han he single se e
sys em. The op imum was calcula ed assuming linea speedup
om he iming wi h one clien .
0 20 40 60 80 100
0
1000
2000
3000
4000
5000
6000
7000
Numbe o p ocesso s
P ocessing ime (seconds)
1 se e
5 se e s
Op imal pe o mance
Fig. 5. P ocessing ime compa isons wi h 1- and 5-se e dis ibu ed
compu a ion sys ems, in he p esence o simula ed conges ion.
V. CONCLUSION
We ha e examined a a ie y o o he exis ing dis ibu ed
sys ems [1], [2], [3], [5], [7], [9], [15], [20], [21], [24], [26],
[28], [29], [30] and ha e p oduced a sys em ha combines
he ad an ages o e ed by each o hese exis ing sys ems and
o e comes many o he disad an ages o each sys em. Cen al
o ou success was he use o Ja a o implemen he dis ibu ed
sys em. Ou dis ibu ed sys em is capable o being deployed
in a ypical in e ne /in ane en i onmen . Some o he ea u es
o ou sys em include a mul i-p oblem adap i e schedule
ope a ing independen ly on mul iple se e s o ganised in a
hie a chical model, a emo e se e in e ace, emo e upda ing
o clien so wa e, dynamic changing o da a uni sizes, and
inbuil comp ession o da a. We also p esen ed an adap i e
schedule ha a emp s o minimise he p ocessing ime o
p oblems in he sys em, and balance he se o p oblems based
on use p io i y and p oblem complexi y.
Fu u e imp o emen s o he sys em will allow o he
dynamic ebalancing o he opology o he sys em o imp o e
pa allel e iciency, allow subo dina e se e s o ake o e
supe io se e s in he e en o ailu e hus p o iding a mo e
obus sys em, and enhancemen o he scheduling s a egy o
include a neu al ne wo k which employs online lea ning [6].
We will also inco po a e suppo o SQL da abases, o allow
mo e s anda dised and less complex message passing be ween
in e media e esul s gene a ed by p oblems, as well as s o age
o s a e in o ma ion abou p oblems o acili a e a mo e obus
sys em such as in [7], [9], [17], [31].
The so wa e is eely a ailable unde an open sou ce GNU
GPL licence om he sys em homepage loca ed a
h p://www.cs.may.ie/dis ibu ed/
VI. ACKNOWLEDGEMENT
Suppo is acknowledged om he I ish Resea ch Council
o Science, Enginee ing, and Technology, unded by he
Na ional De elopmen Plan.
REFERENCES
[1] J. AiJa oodi, N. Mohamed, H. Jiang, and D. Swanson. Middlewa e
in as uc u e o pa allel and dis ibu ed p og amming models in he -
e ogeneous sys ems. IEEE T ansac ions on Pa allel and Dis ibu ed
Sys ems, 14(11):1100–1111, No embe 2003.
[2] D. Ande son. Public compu ing: Reconnec ing people o science. In
Con e ence on Sha ed Knowledge and he Web, pages 17–19, Mad id,
Spain, No embe 2003.
[3] D. Ande son, J. Cobb, E. Ko pela, M. Lebo sky, and D. We hime .
Massi ely dis ibu ed compu ing o SETI. Compu ing in Science &
Enginee ing, 3(1):78–83, Feb 2001.
[4] C. G. A keson, A. W. Moo e, and S. Schaal. Locally weigh ed lea ning.
A i icial In elligence Re iew, 11(1-5):11–73, 1997.
[5] F. Be man, R. Wolski, H. Casano a, W. Ci ne, H. Dail, M. Fae man,
S. Figuei a, J. Hayes, G. Obe elli, J. Schop , G. Shao, S. Smallen,
N. Sp ing, A. Su, and D. Zago odno . Adap i e compu ing on he g id
using AppLeS. IEEE T ansac ions on Pa allel and Dis ibu ed Sys ems,
14(4):369–382, Ap il 2003.
[6] J. P. Bigus and J. Bigus. Cons uc ing In elligen agen s wi h Ja a.
Wiley Compu e Publishing, New Yo k,USA, 1998.
[7] K. Bi man, R. an Renesse, and W. Vogels. Na iga ing in he s o m:
using as olabe o dis ibu ed sel -con igu a ion, moni o ing and adap-
a ion. In Au onomic Compu ing Wo kshop, pages 4–13, Sea le, WA,
USA, June 2003.
[8] J. M. Bull, L. A. Smi h, L. Po age, and R. F eeman. Benchma king
ja a agains c and o an o scien i ic applica ions. In Ja a G ande,
pages 97–105. ACM P ess, 2001.
[9] G. Deen, T. Lehman, and J. Kau man. The Almaden Op imalG id
p ojec . In Au onomic Compu ing Wo kshop, pages 14–21, Sea le, WA,
USA, June 2003.
[10] M. D ozdowski and P. Wolniewicz. Ou -o -co e di isible load p o-
cessing. IEEE T ansac ions on Pa allel and Dis ibu ed Sys ems,
14(10):1048–1056, Oc obe 2003.
[11] H. Edels ein. Un a eling clien /se e a chi ec u e. DBMS, 7(5):34–41,
May 1994.
[12] T. Elgamal. A public key c yp osys em and a signa u e scheme based
on disc e e loga i hms. IEEE T ansac ions on In o ma ion Theo y,
31(4):469–472, July 1985.
[13] M. J. Fische and M. Me i . App aising wo decades o dis ibu ed
compu ing heo y esea ch. Dis ibu ed Compu ing, 16(2–3):239–247,
Sep embe 2003.
[14] D. Flanagan. Ja a in a Nu shell. O’Reilly and Associa es, UK, 4 h
edi ion, 2002.
[15] K. F i sche, J. Powe , and J. Wald on. A ja a dis ibu ed compu a ion
lib a y. In 2nd In e na ional Con e ence on Pa allel and Dis ibu ed
Compu ing Applica ions and Technologies, pages 236–243, Taipei, Tai-
wan, July 2001.
[16] S. Ha i i, H. Topcuoglu, and M. Wu. Pe o mance-e ec i e and
low-complexi y ask scheduling o he e ogeneous compu ing. IEEE
T ansac ions on Pa allel and Dis ibu ed Sys ems, 13(3):260–274, Ma ch
2002.
[17] M. Jelasi y, M. P euß, and B. Paech e . A scaleable and obus
amewo k o dis ibu ed applica ion. In P oceedings o he Cong ess on
E olu iona y Compu a ion, pages 1540–1545, Honolulu, Hawaii, USA,
May 2002.
[18] C. Jin, D. Wei, S. H. Low, G. Buh mas e , J. Bunn, D. H. Choe, R. L. A.
Co ell, J. C. Doyle, W. Feng, O. Ma in, H. Newman, F. Paganini,
S. Ra o , and S. Singh. FAST TCP: F om heo y o expe imen s.
submi ed o IEEE Communica ions magazine, Ap il 2003.
[19] S. Kanhe e, A. Pa ekh, and H. Se hu. Fai and e icien packe scheduling
using elas ic ound obin. IEEE T ansac ions on Pa allel and Dis ibu ed
Sys ems, 13(3):324–326, Ma ch 2002.
[20] T. Keane, R. Allen, T. J. Naugh on, J. McIne ney, and J. Wald on.
Dis ibu ed Ja a pla o m wi h p og ammable MIMD capabili ies. In
N. Guel i, E. As esiano, and G. Reggio, edi o s, Scien i ic Enginee -
ing o Dis ibu ed Ja a Applica ions, olume 2604, pages 122–131.
Sp inge Lec u e No es in Compu e Science, Feb ua y 2003.
[21] E. K iege and G. V iend. Models@Home: dis ibu ed compu ing in
bioin o ma ics using a sc eensa e based app oach. Bioin o ma ics,
18(2):315–318, Feb ua y 2002.
[22] M. Miglia di, V. Sunde am, A. Geis , and J. Donga a. Dynamic
econ igu a ion and i ual machine managemen in he Ha ness me a-
compu ing sys em. In Lec u e No es in Compu e Science, olume 1505,
pages 127–134. Sp inge Ve lag, 1998.
[23] F. Paganini, Z. Wang, J. C. Doyle, and S. H. Low. Conges ion con ol o
high pe o mance, s abili y and ai ness in gene al ne wo ks. submi ed
o publica ion, Ap il 2003.
[24] R. Raman, M. Li ny, and M. Solomon. Ma chmaking: Dis ibu ed
esou ce managemen o high h oughpu compu ing. In P oceedings
o he Se en h IEEE In e na ional Symposium on High Pe o mance
Dis ibu ed Compu ing, Chicago, IL, USA, 1998.
[25] J. Semke, J. Mahda i, and M. Ma his. Au oma ic cp bu e uning.
In P oceedings o ACM SIGCOMM ’98, pages 315–323. ACM P ess,
1998.
[26] T. Sil es e, E. Nugues, G. Pe i`e e, M. Gouy, and L. Du e . Phyloja a
: a gene ic clien -se e ool o phylogene ic ee econs uc ion -
applica ion o g id compu ing. In M.-F. Sago and H.-P. Lenho ,
edi o s, Eu opean Con e ence on Compu a ional Biology, Pa is, F ance,
Sep embe 2003.
[27] Sun Mic osys ems Inc. Ja a RMI - Dis ibu ed Compu ing o Ja a.
Whi e pape .
[28] M. Su deanu and D. Moldo an. Design and pe o mance analysis o
a dis ibu ed ja a i ual machine. IEEE T ansac ions on Pa allel and
Dis ibu ed Sys ems, 13(6):611–627, June 2002.
[29] D. Thain, T. Tannenbaum, and M. Li ny. Condo and he G id. In
F. Be man, A. Hey, and G. Fox, edi o s, G id Compu ing: Making The
Global In as uc u e a Reali y. John Wiley, 2003.
[30] Uni ed De ices. G id MP Pla o m A chi ec u e, 2003. Whi e Pape .
[31] R. an Renesse, K. Bi man, and W. Vogels. As olabe: A obus and
scalable echnology o dis ibu ed moni o ing, managemen and da a
mining. ACM ansac ions on Compu e Sys ems, 21(2):164–206, May
2003.
[32] R. Wolski, N. T. Sp ing, and J. Hayes. The ne wo k wea he se ice: a
dis ibu ed esou ce pe o mance o ecas ing se ice o me acompu ing.
Fu u e Gene a ion Compu e Sys ems, 15(5–6):757–768, 1999.