scieee Science in your language
[en] (orig)

Adaptive Scheduling Across a Distributed Computation Platform

Abstract

A programmable Java distributed system, which adapts to available resources, has been developed to minimise the overall processing time of computationally intensive problems. The system exploits the free resources of a heterogeneous set of computers linked together by a network, communicating using SUN Microsystems' Remote Method Invocation and Java sockets. It uses a multi-tiered distributed system model, which in principal allows for a system of unbounded size. The system consists of an n-ary tree of nodes where the internal nodes perform the scheduling and the leaves do the processing. The scheduler nodes communicate in a peer-to-peer manner and the processing nodes operate in a strictly client-server manner with their respective scheduler. The independent schedulers on each tier of the tree dynamically allocate resources between problems based on the constantly changing characteristics of the underlying network. The system has been evaluated over a network of 86 PCs with a bioinformatics application and the travelling salesman optimisation problem.

Read accessible full text

Adaptive Scheduling Across a Distributed Computation Platform

Author: Page, Andrew,Keane, Thomas,Naughton, Thomas J.
Publisher: IEEE Computer Society
Year: 2004
Source: https://mural.maynoothuniversity.ie/id/eprint/182/1/APageISPDC2004.pdf
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=PxN
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= xN
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.