A Technique to Automatically Determine Ad-hoc Communication Patterns at Runtime
Abstract
Producción Científica
Full text
A Technique o Au oma ically De e mine Ad-hoc Communica ion
Pa e ns a Run ime
Ana Mo e on-Fe nandez, A u o Gonzalez-Esc ibano, and Diego R. Llanos
Depa amen o de In o m´a ica, Edi . Tecn. de la In o maci´on, Uni e sidad de Valladolid, Campus Miguel Delibes,
47011 Valladolid, Spain.
Abs ac
Cu en High Pe o mance Compu ing (HPC) sys ems a e ypically buil as in e connec ed clus-
e s o sha ed-memo y mul ico e compu e s. Se e al echniques o au oma ically gene a e pa -
allel p og ams om high-le el pa allel languages o sequen ial codes ha e been p oposed. To
p ope ly exploi he scalabili y o HPC clus e s, hese echniques should ake in o accoun he
combina ion o da a communica ion ac oss dis ibu ed memo y, and he exploi a ion o sha ed-
memo y models.
In his pape , we p esen a new communica ion calcula ion echnique o be applied ac oss
di e en SPMD (Single P og am Mul iple Da a) code blocks, con aining se e al uni o m da a
access exp essions. We ha e implemen ed his echnique in T asgo, a p og amming model and
compila ion amewo k ha ans o ms pa allel p og ams om a high-le el pa allel speci ica ion
ha deals wi h pa allelism in a uni ied, abs ac , and po able way.
The p oposed echnique compu es a un ime exac coa se-g ained communica ions o dis-
ibu ed message-passing p ocesses. Applying his echnique a un ime has he ad an age o
being independen o compile- ime decisions, such as he ile size chosen o each p ocess. Ou
app oach allows he au oma ic gene a ion o p e-compiled mul i-le el pa allel ou ines, lib a ies,
o p og ams ha can adap hei communica ion, synch oniza ion, and op imiza ion s uc u es o
he a ge sys em, e en when compu ing nodes ha e di e en capabili ies.
Ou expe imen al esul s show ha , despi e ou un ime calcula ion, ou app oach can au o-
ma ically p oduce e icien p og ams compa ed wi h MPI e e ence codes, and wi h codes gen-
e a ed wi h au o-pa allelizing compile s.
Keywo ds: SPMD models, Dis ibu ed communica ions, Pa allel p og amming, T asgo
1. In oduc ion
Pa allel machines a e becoming mo e he e ogeneous, mixing de ices wi h di e en capabil-
i ies in he con ex o hyb id clus e s, wi h hie a chical sha ed- and dis ibu ed-memo y le els.
The ocus on pa allel applica ions is shi ing o mo e di e se and complex solu ions, exploi ing
I
Email add ess: {ana,a u o,diego}@in o .u a.es. (Ana Mo e on-Fe nandez, A u o
Gonzalez-Esc ibano, and Diego R. Llanos)
P ep in submi ed o Pa allel Compu ing June 26, 2017
se e al le els o pa allelism, wi h di e en s a egies o pa alleliza ion. P og amming in his kind
o en i onmen is challenging. Many app oaches ha e been p oposed o e he las ew yea s ol-
lowing wo di e en pa hs: P og amming using pa allel p og amming models, o au oma ically
gene a ing pa allel code om sequen ial p og ams.
Using cu en pa allel p og amming models (e.g. MPI, OpenMP, In el TBBs, Cilk, and
PGAS languages such as Chapel, X10, o UPC), he applica ion p og amme s ill aces many
impo an decisions no ela ed wi h he pa allel algo i hms, bu wi h implemen a ion issues ha
a e key o ob aining e icien p og ams. Fo example, decisions abou pa i ion and locali y s.
synch oniza ion/communica ion cos s; g ain selec ion and iling; p ope pa alleliza ion s a egies
o each g ain le el; o mapping, layou , and scheduling de ails.
Au oma ic code gene a ion echniques can be used o au oma ically ans o m high-le el
pa allel exp essions o sequen ial codes o pa allel p og ams ha ake in o accoun many o
hese issues. Mos o hese echniques we e designed o be applied a compile ime [1, 2, 3,
4]. Fo example, he wo k p esen ed in [2] p oposes a echnique ha , om a sequen ial code,
gene a es a low-le el pa allel code o dis ibu ed-memo y sys ems using he Message Passing
In e ace (MPI) lib a y. This echnique imp o es p e ious schemes because he code i gene a es
is pa ame ic in he numbe o p ocesses and p oblem sizes, educing he communica ed olume
o da a. Howe e , i needs o ix a single ile size a compile ime, e en i he dis ibu ed sys em
has nodes wi h di e en capabili ies.
In his pape , we p esen a new communica ion calcula ion echnique o be applied ac oss di -
e en SPMD (Single P og am Mul iple Da a) blocks o code, ha con ain se e al di e en da a
accesses exp essions o he same da a s uc u e, whose indexes a e calcula e wi h uni o m a ine
exp essions in he indexes selec o s. We conside as uni o m a ine exp essions, hose exp es-
sions ha de i e in accesses o a mul i-dimensional pa alelo ope o he da a s uc u e domain.
Fo wo dimensions, his means ec angula shapes. The echnique suppo s codes wi h se e al
da a accesses o he same da a s uc u e. Thus, he esul ing domain accessed by a code block is
a compound o pa alelo ope shapes, ha can be non-con ex.
We ha e implemen ed his echnique in T asgo [5], a p og amming model and compila ion
amewo k o gene a e pa allel p og ams om a high-le el pa allel speci ica ion ha deals wi h
pa allelism in a uni ied, abs ac , and po able way. The p oposed echnique compu es a un ime
exac coa se-g ained communica ions be ween wo consecu i e pa allel blocks o dis ibu ed
message-passing p ocesses.
The main con ibu ions o his pape a e he ollowing:
•A echnique o de e mine a un ime communica ions ac oss SPMD blocks o dis ibu ed-
memo y sys ems. The communica ions calcula ed a e:
–Coa se-g ained in he sense ha communica ion calcula ion ac oss wo pa allel SPMD
blocks is done once o he whole index space mapped o a p ocess a un ime, inde-
penden ly o he numbe o sizes o iles gene a ed inside he p ocess. This enables
di e en ile sizes o be used in he same compu a ion a he same hie a chical le el,
an impo an ea u e in achie ing a good pe o mance on he e ogeneous sys ems ha
include machines wi h di e en a chi ec u es [6].
–Exac because hey a e op imal in e ms o communica ion olume. Ou un ime
calcula ion skips all he duplica ed da a elemen s in he communica ion. Thus, no
da a is communica ed wice, and no unneeded o con ol da a is communica ed ac oss
any wo dis ibu ed p ocesses.
2
** Illus a i e example
Inpu s:
a: 1s s encil pa ame e
b: 2nd s encil pa ame e
M[n][n]: Ma ix wi h ini ial alues
limi : Numbe o i e a ions
Ou pu :
M[n][n] : Ma ix wi h esul alues
Tempo al a iables:
M_ emp[n][n]: Auxilia ma ix.
** Time loop
Do i e = 1 o limi
** Fi s SPMD block: Upda e M_ emp
Do i = a o n-b
Do j = a o n-b
M_ emp[i][j] = M[i][j]
** Second SPMD block: Compu e s encil ope a ion
Do i = a o n-b
Do j = a o n-b
M[i][j] = (M_ emp[i-a][j] + M_ emp[i+b][j] + M_ emp[i][j-a] + M_ emp[i][j+b])/4;
Figu e 1: Sequen ial algo i hm o he illus a i e example.
•The implemen a ion o he p oposed communica ion calcula ion echnique in a pa allel
p og amming amewo k, whose inpu language allows he p og amme o eason in e ms
o logical p ocesses, wi hou acing decisions abou g anula i y, h ead managemen , o
in e p ocess communica ion.
•An expe imen al s udy o e alua e ou un ime p oposal compa ing i wi h a compile- ime
s a e-o - he-a ool ha gene a es communica ion codes, and wi h pu e-MPI e e ences
codes.
Ou echnique pos pones o un ime pa o he analysis and decisions needed o ans o m
p og am abs ac ions o ac ual p ocesses. Thus, he p og ams can adap hei beha io a un ime,
dealing wi h di e en pa i ions, g anula i y, da a-dis ibu ion, memo y hie a chies, ile sizes, o
synch oniza ion and communica ion s uc u es.
We s a ou discussion wi h an illus a i e example based on a s encil compu a ion in Sec . 2,
showing he ans o ma ion echniques p esen ed o he es o he pape . Sec ion 3 desc ibes
he T asgo model and i s ools. Sec ion 4 in oduces he new echniques applied in T asgo. To
show he applicabili y and e iciency o he app oach, we include se e al expe imen al s udies
in Sec . 5, compa ing pe o mance on dis ibu ed- and sha ed-memo y pla o ms wi h MPI e -
e ence codes, and codes gene a ed wi h au o-pa allelizing compile s. The esul s show ha ou
app oach can au oma ically p oduce e icien p og ams despi e he o e head o he calcula ion
pe o med a un ime. Sec ion 6 discusses some ela ed wo k. Sec ion 7 p esen s he conclusion,
and u u e wo k.
2. Illus a i e example and O e iew
This sec ion p esen s an illus a i e example ha se es as a quick o e iew o he echniques
p esen ed in his pape . The example is a modi ica ion o a Jacobi PDE sol e o Poisson’s
3
(1) SPMD:
Upda e
(2) SPMD:
Compu a ion
o (i=lowe _x; i<uppe _x; i++)
o (j=lowe _y; i<uppe _y; j++)
M[i][j]= 0.25* (M2[i-a][j] + M2[i+b][j]
M2[i][j-a] + M2[i][j+b] )
o (i=lowe _x; i<uppe _x; i++)
o (j=lowe _y; i<uppe _y; j++)
M2[i][j]= M[i][j]
Communica ion
s age
Dis ibu e comp.
among p ocesses
}
Communica ion
s age
o ( =0; <limi ; ++){
lowe _x = map_lowe (a, n-b, myRank, num_p ocesses, 0)
uppe _x = map_uppe (a, n-b, myRank, num_p ocesses, 0)
lowe _y = map_lowe (a, n-b, myRank, num_p ocesses, 1)
uppe _y = map_uppe (a, n-b, myRank, num_p ocesses, 1)
// De e mine communica ion pa e n 1
// Execu e communica ion pa e n 1
// De e mine communica ion pa e n 2
// Execu e communica ion pa e n 2
Figu e 2: Scheme o he pa allel algo i hm ollowing an SPMD model o he illus a i e example (le ), and code
exce p s o he main blocks ( igh ). The pa ame e ade e mines he halo sizes on he op and le sides, and he b
pa ame e de e mines he same on he bo om and igh sides.
equa ion o compu e hea ans e in a disc e ized wo-dimensional su ace. This is a simple
da a-pa allel example, ha aims o in oduce he eade s o he basis o he app oach. This ex-
ample can be used o show basic concep s o MPI synch oniza ion (see e.g. [7]). I con ains
clea compu a ion and communica ion s ages wi h uni o m da a access exp essions. The sequen-
ial algo i hm is shown in Fig. 1. In ou example we also in oduce wo in ege pa ame e s a,b
ha a e used in he access exp essions o selec a un ime he dis ance o he posi ions consid-
e ed neighbo s in e ms o he s encil ope a ion ca ied ou o ha pa icula in oca ion o he
compu a ion.
2.1. P og amming wi h an SPMD model
We i s p esen an o e iew o he ypical app oach used o p og am he illus a i e exam-
ple ollowing an SPMD model. In he sequen ial algo i hm ( ecall Fig. 1), we can dis inguish
wo di e en blocks o code inside he ime loop ha can be pa allelized independen ly wi hou
iola ing any da a dependence ( ans o ming hem in o SPMD blocks). Figu e 2 (le ) shows a
dis ibu ed pa allel p og amming scheme o he illus a i e example ollowing an SPMD model.
To p og am his algo i hm in pa allel, he only da a dependences ha mus be aken in o accoun
a e hose p oduced be ween he SPMD blocks. Fo his eason, a communica ion/synch oniza-
ion s age is inse ed be ween hem.
When we execu e his pa allel example algo i hm in sha ed-memo y sys ems, a synch oniza-
ion s age is enough o a oid da a dependences be ween he wo pa allel s uc u es. Howe e ,
in dis ibu ed memo y sys ems, he da a w i en by a p ocess in he i s SPMD block should be
sen o o he p ocesses ha need hese da a o execu e he second SPMD block. No ice ha a
synch oniza ion is implici in his communica ion.
P og amming o dis ibu ed-memo y sys ems has he challenge o he p og amme o deal-
ing wi h he da a and compu a ion dis ibu ion be ween di e en p ocesses. In ou example (see
Fig. 2, igh ), we ha e dis ibu ed he compu a ion using he mapping unc ions map lowe () and
4
WOM_ emp,1,p
WIM_ emp,2,p
U
Read
exp essions
uppe _ylowe _y
lowe _x
uppe _x
Exp : [i-a][j] Exp : [i+b][j] Exp : [i][j-a] Exp : [i][j+b]
U U
Exp : [i][j]
=
lowe _x-a
uppe _x-a
lowe _x+b
uppe _x+b
uppe _y-alowe _y-a uppe _y+blowe _y+b
W i e
exp essions
NORM
Figu e 3: Using he ead and w i e da a-access exp essions inside he pa allel s uc u e o he illus a i e example o
calcula e he wo king inpu and ou pu index se s (W2
I,W1
O) o M emp o a gene ic p ocess. This example assumes
posi i e pa ame e s a,b≥0. The No m ope a ion on he las s age educes he numbe o boxes used o ep esen he
union o domains. The d awings conside he pa icula case o a=2,b=3, o a domain wi h 8 ×8 con iguous indexes.
map uppe (). These mapping unc ions ecei e he i s and las i e a ion o he loop o be dis-
ibu ed, he iden i ie o he cu en p ocess, and he o al numbe o p ocesses. These unc ions
e u n he loop limi s co esponding o he chunk o i e a ions ha should be execu ed by he
cu en p ocess, a oiding o e lappings wi h o he p ocesses.
Each p ocess is eady o pe o m he compu a ion as soon as i has he co esponding da a
in i s local memo y. These da a include no only he da a in he posi ions ha ma ch he loop
i e a ions, bu also he da a ha co espond o hei halo, which may be owned by o he p ocesses.
A communica ion s age be o e each SPMD block should e ie e he co esponding halos, hus
ensu ing ha each p ocess has a local copy o all he needed da a. In ou pa icula example,
be o e he execu ion o he i s SPMD block, a communica ion s age is no necessa y as each
p ocess al eady has all he da a i needs. Howe e , in he second SPMD block, each p ocess
needs da a ha ha e been upda ed by o he p ocesses (da a in he halo). The communica ion
phase in his case depends on execu ion pa ame e s, such as he ma ix size, he ile size, he
numbe o p ocesses, he pa i ion policy ( he way in which he da a we e pa i ioned among he
p ocesses), and he alues o aand bpa ame e s which de ine he halo sizes.
The echnique p esen ed in his pape de e mines au oma ically a un ime he communi-
ca ion pa e ns needed be ween wo consecu i e pa allel s uc u es, aking in o accoun hese
pa ame e s, ega dless o hei a ailabili y a compile o execu ion ime, ega dless o he appli-
ca ion o o he sequen ial o iling op imiza ion echniques inside each p ocess.
2.2. O e iew o he communica ion de e mina ion echnique
We p esen he e an o e iew o he p oposed echnique o de e mine au oma ically he com-
munica ion pa e ns among di e en SPMD blocks. As we ha e seen in he p e ious desc ip ion,
we usually need a communica ion s age be ween wo consecu i e SPMD blocks o ensu e a co -
ec execu ion. In ou example, we need o communica e some da a w i en in he i s SPMD
block by each p ocess o o he p ocesses ha ead hese da a in he second SPMD block. Ou
echnique consis s o wo s eps:
1. In o de o de e mine he da a ead and w i en in each pa allel s uc u e, o each SPMD
block in he p og am, we gene a e a compile ime a pa ame ized unc ion ha , a un ime,
e u ns he se o da a being ead o w i en o a gi en p ocess iden i ie p.
Fo he illus a i e example, we show an example o bo h se s o indexes e u ned by hese
gene a ed unc ions in Fig. 3. We name he se s o indexes o he ma ix M ead and w i en
5
3
6 7
5
8
0 1 2
3
6 7
5
8
0 1 2
4 4
o (p=0; p<Numbe o p ocesses; p++)
i (p != 4) {
Calcula e in e sec ion o WOM_ emp,1,p and WIM_ emp,2,4
S o e in e sec ion in CR
}
Da a needed om P1:
Da a needed om P3:
Da a needed om P5:
Da a needed om P7:
CR
Communica ion Recei e (CR) pa e n o p ocess 4:
Figu e 4: Calcula ion o Communica ion Recei e (CR) pa e n be ween he wo pa allel s uc u es o he illus a i e
example. Example o a numbe o p ocesses P=9, a p ocess iden i ie 4, he p e iously gene a ed unc ions WM emp,2,∗
I
and WM emp,1,∗
O, a mapping/pa i ion unc ion ha e u ns i egula ec angula blocks, and pa icula alues o symbolic
pa ame e s a=2,b=3.
by he k- h SPMD block in he code, a a gi en p ocesso pas ollows: Inpu Wo king Se
Indexes WM,k,p
I, and Ou pu Wo king Se Indexes WM,k,p
O. These unc ions a e gene a ed
a compile ime using he da a-access exp essions ound in he inpu code. In Fig. 3, we
see how he se WM emp,2,p
Iis calcula ed by applying he uni o m access exp essions ound
in he code o he calcula ed loop limi s (lowe x, lowe y, uppe x, uppe y). The se o
indexes is no malized o be ep esen ed wi h a se o ec angula shapes.
2. In he second s ep, we apply an algo i hm a un ime o de e mine he communica ion
pa e ns and o s o e a compac desc ip ion o hem in an objec . To calcula e he da a ha
ha e o be ecei ed by a local p ocess om a emo e p ocess p, he algo i hm in e sec s
he se o indexes ead by he local p ocess in he second SPMD block ( ha is, WM,k+1,local
I)
wi h he se o indexes w i en by he p ocess pin he i s SPMD block (WM,k,p
O). Figu e 4
shows a isual ep esen a ion o ou un ime algo i hm o he illus a i e example. The
example shows he calcula ion o he Communica ion Recei e (CR) pa e n o p ocess 4.
On he le , we can see he M emp ma ix dis ibu ed among 9 p ocesses wi h an a bi a y
i egula pa i ion policy, and he da a se (do ed lines) o be ead by p ocess 4 in he
ollowing SPMD block. On he igh , we see he da a ha should be ecei ed by p ocess 4
om di e en emo e p ocesses. The pa e ns a e calcula ed using he p ope in e sec ions
o he da a owned by di e en p ocesses.
To calcula e he da a o be sen by a local p ocess o a p ocess p, he algo i hm pe o ms
he opposi e in e sec ion, be ween he se o indexes w i en by he local p ocess in he
i s SPMD block ( ha is, WM,k,local
O) and he se o indexes ead by he pp ocess in he
second SPMD block (WM,k+1,p
I). An emp y in e sec ion indica es ha no send (o ecei e)
ope a ion is needed o ha pa icula p ocess p.
6
F on −end ansla o
Na i e compile
High le el sou ce code
CMAPS
SPC−XML speci ica ion
XML
XML
SMP Code + HIT calls
C + un ime HITmap
Mapped p og am
Mul ile el code
C + un ime HITmap + OpenMP
Bina y execu able
P og am ep esen a ions T ans o ma ions
Code ans o ma ions
Exp ession builde
Back−end
Sha ed memo y
Polyhed al model: Plu o
HITmap lib a y
Xsl
Xsl
Figu e 5: S uc u e o he T asgo ans o ma ion amewo k.
A e applying his echnique o e e y p ocess and e e y a ay we can apply he de e -
mined send and ecei e pa e ns o pe o m he ac ual communica ions.
Ou echnique equi es ha any symbolic pa ame e mus ha e he same alue on e e y p o-
cess. Thus, he se o indexes accessed by a emo e p ocess pcan be calcula ed by any p ocess,
wi h no in e -p ocess communica ion. A deepe discussion abou he cons ain s, he ea u es
used o educe he complexi y, and a o mal de ini ion o ou echnique is p esen ed in Sec . 4.
3. The T asgo Model
The p e ious sec ion b ie ly desc ibes ou p oposed echnique. In his sec ion, we e iew he
T asgo pa allel p og amming and execu ion model, whe e ou echnique has been implemen ed.
The T asgo model [5] p oposes he use o an explici ly pa allel, bu high-le el and s uc-
u ed ep esen a ion o pa allel algo i hms. I uses es ic ed synch oniza ion a he highe le el,
gene a ing mo e e icien and less synch onized pa allel s uc u es a he low le el. The o iginal
model is based on he SP (Se ies-Pa allel) p ocess model [8] and da a-dis ibu ion algeb as, p o-
iding clea and well-de ined seman ics [9]. The model is ee o ace condi ions, unexpec ed
dead-locks, o s ochas ic beha io s. The high-le el code uses a global iew app oach in hie -
a chical decomposi ions. The seman ics p o ide clea synch oniza ion poin s and hie a chical
global s a es ha simpli y es ing and debugging.
3.1. O e iew o he code ans o ma ion amewo k
Figu e 5 shows he s uc u e o he T asgo ans o ma ion amewo k. The le column shows
he p og am ep esen a ions, and he igh columns he ans o ma ion laye s.
A on -end ansla es he inpu language (in ou case CMAPS [10]) o an XML in e nal
ep esen a ion. The main pa o he ans o ma ion laye ans o ms he global add ess space
in o a pa i ioned add ess space, building he unc ions o compu e communica ions ac oss i ual
p ocesses. The ans o med code is ew i en by a back-end ha gene a es C code wi h calls o
he Hi map un- ime lib a y (see Sec . 3.3). The esul ing compu a ion code, gene a ed o he
local dis ibu ed p ocess, can inally be op imized h ough polyhed al ools o gene a e op imized
pa allel code o he sha ed memo y le el using OpenMP. Plu o compile [11] has been shown
7
1/* Second Sequen ial unc ion */
2
oid
upda eCell( in
double
up, in
double
down, in
double
le , in
double
igh ,
3ou
double
esul ) {
4* esul = ( up + down + le + igh ) / 4 ;
5}
6
7/* Fi s Sequen ial unc ion */
8
oid
upda eDa a( in
double
Da a, ou
double
esul ) {
9* esul = Da a ;
10 }
11
12 /*
Pa allel
s encil code wi h pa ame ic dependences */
13 coo dina ion
oid
caseA( inou ile
double
M[][], in
in
limi ,
14 in
in
a, in
in
b, in layou pa i ion_policy) {
15
double
M_ emp[][];
16 /* Dis ibu e a ays*/
17
A ayMap
( M, pa i ion_policy );
18
A ayMap
( M_ emp, pa i ion_policy );
19
20 /* Time loop */
21 loop( in [1:limi ] ) {
22 /* Fi s SPMD block: Upda e M_ emp */
23
pa allel
( i,j in M) {
24
do
: upda eDa a( M[i ][ j ], M_ emp[ i ][ j ] );
25 }
26 /* Second SPMD block: Compu e s encil ope a ion */
27
pa allel
( i,j in M ) {
28
do
: upda eCell( M_ emp[ i-a ][ j ], M_ emp[ i+b ][ j ],
29 M_ emp[ i ][ j-a ], M_ emp[ i ][ j+b ],
30 M[ i ][ j ] );
31 }
32 }
33 }
Figu e 6: CMAPS code o he illus a i e example.
o be a good ool o au o pa allelize sequen ial codes o sha ed-memo y sys ems. As a p oo
o concep , we use Plu o on he codes ob ained a e applying he mapping policy. Thus, we
make i op imize he local compu a ion inside each MPI p ocess, in o de o e icien ly exploi
in e nally he mul i-co e p ocesso o each node. The me hodology used o in eg a e Plu o wi h
he Hi map oolchain was desc ibed in [12]. The inal code is compiled wi h a na i e C compile .
Figu e 6 shows he illus a i e example coded in CMAPS, he cu en T asgo inpu language.
Fi s , we de ine he sequen ial unc ions o apply o each elemen , speci ying he ou pu o inpu
ole o he pa ame e s (lines 2,8). The pa allel unc ion is de ined using he coo dina ion modi-
ie , and also speci ying he ole o he pa ame e s (line 13). In i s body, he unc ion A ayMap()
dis ibu i ely alloca es bo h a ays M emp and Min e ms o he esul s o a mapping/pa i ion
policy named pa i ion policy, which is also a pa ame e (line 17-18). A e he dis ibu ion, he
code upda es and compu es each elemen o M emp in pa allel using he sequen ial unc ions
p e iously de ined (lines 23, 28). The pa allel s uc u e in CMAPS maps he compu a ion indi-
ca ed in he do: code o each indexes pai in he domain speci ied in he clause inside he pa allel
b acke s.
3.2. No a ions and de ini ions
In his sec ion, we p esen de ini ions used in he es o he pape . In his wo k, we ocus on
a ays wi h egula dense and s ided domains.
8
Signa u e is a iple o in ege numbe s Shb,e,si(meaning begin,end, and s ide). The se o
indexes exp essed by a signa u e is Shb,e,si={b≤i≤e: (i−b) mod s=0}.
Domain is a subspace o Zn. Rec angula n-dimensional pa allelo ope domains, dense o s ided,
can be ep esen ed by a uple o n Signa u es. Le us conside an n-dimensional domain
Dnhs0,s1, ..., sn−1i ∈ Zn, whe e s0∈S∗, .., sn−1∈S∗a e he signa u es whose Ca esian
p oduc de ines he domain. This kind o s uc u es only ep esen ec angula shapes.
Wo king se is a gene ic se o indexes. We ep esen gene ic se s as unions o signa u e do-
mains, WM=∪q
i=0di:q∈θ(N),di∈Dn.
Tile is an objec ha associa es da a elemen s o a gi en ype o index elemen s o a domain.
The domain o a ile is deno ed as D(T), T:D→ ype.
Logical p ocess is a uple Ph ,WM
I,WM
Oi, whe e is a unc ion o subp og am, WM
Iis he
wo king se ha P ecei es as inpu ( he da a indexes o M ha a e ead), and WM
Ois he
se used as ou pu ( he da a indexes o M ha a e w i en). Logical p ocesses may be
composed in sequence, o in pa allel. A sequen ial composi ion P1.P2indica es ha 1
is execu ed be o e 2, and ha da a modi ica ions in oduced in he ou pu iles o P1a e
p opaga ed in he co esponding inpu iles o P2. Sequen ial composi ion is associa i e bu
no commu a i e. A pa allel composi ion P1◦P2indica es ha 1and 2can be execu ed
in pa allel. Pa allel composi ion is associa i e and commu a i e.
Wa e- on composi ion (P1•(W )P2), is a pa allel composi ion wi h explici ly added o de
dependences o a bi a y iles, ac oss p ocesses. I he e is o e lapping o he shapes o
W wi h he inpu wo king se o P2and he ou pu wo king se o P1, he unc ion 2canno
be s a ed un il 1has inished. The da a ep esen ed by he o e lapping shapes should be
p opaga ed. This does no allow gene ic da a- low composi ions o be exp essed, wi h
ansi i e cycles ha could lead o dead-locks. This composi ion ope a ion allows pa allel
s uc u es, such as wa e- on compu a ions, and mac opipelines o be exp essed when
used inside a loop.
Vi ual Topology V(N,R) is a g aph whe e he e ices N ep esen i ual p ocesses, associa ed
wi h compu a ional esou ces (g oups o p ocesso s), and he edges R ep esen neighbo
ela ions.
Layou L:D→Vis a unc ion ha maps domain subspaces (indexes o iles o logical p o-
cesses) o he i ual p ocesses in a i ual opology.
3.3. Hi map lib a y
Hi map [13] is a lib a y o he managemen and un- ime mapping o hie a chically iled
a ays used in T asgo. I is based on an SPMD model, and on he message-passing pa adigm.
Hi map has h ee main unc ionali y modules: (a) Domain and ile managemen ; (b) Mapping
modules; and (c) Communica ion pa e ns. Hi map de ines objec s o decla e and manipula e
index se s as mul idimensional pa allelo opes wi h op ional s ide, o as spa se se s. Hi map
de ines a plug-in sys em o include new mapping modules: Vi ual opology cons uc o s and
mapping unc ions named layou s. The modules gene a e objec s ha can be que ied a un ime
o ob ain in o ma ion abou he esul o he mapping o he local, o any emo e p ocess. Finally,
i con ains unc ionali ies o build eusable communica ion pa e ns o iles, o sub iles, ac oss
9
Table 1: Inpu da a sizes (N×N), ime loop i e a ions (T), and h eshold pa ame e , o he di e en benchma ks in he
expe imen al s udies conduc ed in He acles and CETA.
Machine He acles CETA
Benchma ks Sizes, i e a ions, h eshold Sizes, i e a ions, h eshold
Illus a i e example N =7500, T =200 N =7500, T =200
S encil-Op N =5000, Th eshold =0.001 N =5000, Th eshold =0.001
Cannon’s algo i hm N =7680 N =7680
Ma mul N =4000 N =4000
Jacobi-2d N =7000, T =1000 N =5000, T =800
Gauss-Seidel N =7000, T =1000 N =5000, T =800
Blu -Robe s N =13000 N =13000
Gem e N =8000 N =600
compu a ion cos o he communica ion calcula ion a un ime g ows linea ly wi h he numbe
o i ual p ocesses Pin he opology. In many applica ions, he calcula ion o he communi-
ca ion pa e ns can be mo ed ou o he loops and compu ed a ini ializa ion ime. Howe e , in
applica ions whe e he communica ion exp essions a e pa ame ized wi h loop indexes o o he
pa ame e s, he exp essions a e no in a ian , and he communica ion pa e ns should be com-
pu ed a e e y loop i e a ion.
The new T asgo p o o ype allows he addi ion o specialized ans o ma ion modules ha , by
inpu code inspec ion, can de ec pa allelism pa e ns, and subs i u e he gene ic communica ion
calcula ion by speci ic op imized unc ions ha do no a e se all he o he p ocesses o compu e
wo king-se index in e sec ions. The ime o compu e communica ions in hese cases does no
g ow wi h he numbe o p ocesses. Fo example, by checking he exp essions used in he ile
selec ions, i is possible o de ec s encil compu a ions, ha de i e in neighbo synch oniza ion
s uc u es. Simila ly, a ci cula shi pa e n is also de ec able in Cannon’s algo i hm o ma ix
mul iplica ion. Ou T asgo p o o ype includes modules o some simple s encil and shi pa e ns,
which subs i u e he gene ic communica ion calcula ion by code ha calcula es he in e sec ions
only wi h he needed neighbo s o bo h inpu and ou pu wo king se s.
These modules o de ec speci ic pa e ns educe he calcula ion communica ion imes. Ne -
e heless, hey canno be gene alized o any dependences pa e n o mapping policy chosen.
Mo e well-known applica ions o design pa e ns can be analyzed and implemen ed. I is an
in e es ing esea ch ques ion i any well-de ined pa e n can be de ec ed, and i s co esponding
communica ion code can be gene a ed o any mapping policy chosen a un ime.
5. Expe imen al s udy
We ha e conduc ed an expe imen al s udy o alida e ou app oach, and o e i y he e i-
ciency o he esul ing codes. We p esen ou di e en pe o mance s udies:
•One o ou main con ibu ions is he abili y o ou echnique o au oma ically calcula e
communica ions on a dis ibu ed-memo y p og amming model wi hou a ixed ile size a
compile ime. Fo his eason, he expe imen al s udy s a s wi h a pe o mance s udy o
show he pe o mance imp o emen achie ed when he ile size is independen ly uned a
un ime o each machine in ol ed in he compu a ion.
16
•Second, we e alua e he po en ial o e head ha ou un ime echnique can in oduce when
adding mo e compu ing elemen s. The simula ion s udy shows a compa ison o he un-
ime cos o he gene al communica ion de e mina ion using he desc ibed algo i hms
wi h espec o using communica ion pa e ns o speci ic applica ions al eady included
in T asgo.
•Thi d, we pe o m an end- o-end measu e, including c ea ion o da a s uc u es, da a ini-
ializa ion, and he es o T asgo po en ial o e heads o compa e he p og ams gene a ed
by ou T asgo p o o ype wi h MPI p og ams manually de eloped and op imized.
•The las s udy p esen s a compa ison in e ms o compu a ion and communica ion (de-
e mina ion plus execu ion) imes wi h a s a e-o - he-a polyhed al code gene a o o
dis ibu ed-memo y sys ems, he Plu o-MPI compile [2], using se e al benchma ks o he
PolyBench [14].
A mo e concise desc ip ion o each applica ion and benchma k used in he s udies can be
ound in he supplemen a y ma e ial o his pape .
5.1. Expe imen al pla o ms and se up
Two clus e s ha e been used in he di e en expe imen al s udies. The i s one is a homo-
geneous dis ibu ed-memo y sys em called CETA. I is a hyb id clus e ha belongs o CIEMAT
and he Spanish go e nmen . The clus e nodes a e connec ed by In iniband echnology, and hey
ha e wo In el Nehalem-based Xeon 5520 CPUs a 2.27 GHz, wi h 4 co es each. Using 8 nodes
o he clus e , we exploi up o 64 compu a ional uni s.
The o he clus e , A las, is composed by wo mul ico e machines (He acles and Zeus), ha
ac s as a dis ibu ed-memo y clus e . He acles is a Dell Powe Edge R815 se e , wi h 4 AMD
Op e on 6376 p ocesso s a 2.3 GHz, wi h 16 co es each, and 64 co es in o al. Zeus is a 6-co e
In el E5-2620 3 a 2.40GHz (up o 12 h eads wi h hype h eading). He acles is used in all
he expe imen s o es he scalabili y o he gene a ed message-passing codes in sha ed-memo y
pla o ms.
In he expe imen al s udies, all he codes, including he MPI e e ence codes and he p o-
g ams gene a ed by T asgo o Plu o, a e compiled wi h GCC 4.8.3 wi h -O3 lag. We use
mpich3 3.1.3 as MPI implemen a ion.
Table 1 shows he benchma ks, inpu sizes o he ma ices, h eshold, and numbe o i e a ions
o he expe imen al s udies using He acles alone and CETA, excep o he i s expe imen al
s udy, whe e we use he ull A las clus e wi h ma ices o 1 500x1 500 da a elemen s. These
p oblem sizes ha e been chosen o gene a e enough compu a ional load o ob ain signi ican
esul s o ou expe imen al pla o ms. In some cases, like Gem e in CETA, he size is he
maximum suppo ed in he pla o m by he dis ibu ed e sion o Plu o-MPI ha eplica es he
memo y oo p in o he whole global da a s uc u es o each p ocess.
As he ocus o ou wo k is he e icien au oma ic calcula ion and execu ion o communica-
ions among p ocesses, we ha e launched, in all he expe imen s, a dis ibu ed p ocess o each
compu a ional uni , wi hou exploi ing he sha ed memo y o he machines. All he esul s p e-
sen ed a e he minimum execu ion ime on en epe i ions o each expe imen , o elimina e he
ou lie s p oduced by s ochas ic delays in he communica ion sys ems.
17
Table 2: Compu a ion imes (seconds), o he ma ix mul iplica ion benchma k wi h a size o 1500x1500 on he clus e
A las wi h di e en unning o he ile size. TS-He acles, when applying he bes ile size o he He acles machine in
bo h machines. TS-Zeus, when applying he bes ile size o he Zeus machine in bo h machines. Bes TS ep esen s he
p og am ha chooses he bes ile size ound empi ically o each machine.
P ocesses TS-He acles TS-Zeus Bes TS
1+1 11.32 6.23 6.23
2+2 5.69 3.17 3.17
4+4 2.87 2.24 1.78
6+6 2.00 2.11 1.75
12+12 1.16 1.28 0.99
Table 3: Execu ion ime o he communica ion de e mina ion o Jacobi-2D sol e (seconds).
P ocesses Gene al Model Speci ic Model Hie a chical mapping policy
256 2.27 ×10−33.08 ×10−57.60 ×10−4
1 024 3.92 ×10−33.12 ×10−51.48 ×10−3
16 384 0.12 2.97 ×10−51.54 ×10−3
262 144 2.41 3.27 ×10−52.38 ×10−1
1 048 576 53.22 3.21 ×10−59.53 ×10−1
5.2. Imp o emen achie ed by uning he ile size o each p ocess
We ha e de eloped an expe imen al s udy o show he posi i e impac on he pe o mance o
uning he ile size a un ime based on he execu ion machine de ails [6].
We use as benchma k he classical algo i hm o ma ix mul iplica ion (in CMAPS). We ha e
execu ed he p og am wi h di e en ile sizes in he wo di e en machines o he A las clus e ,
o empi ically de e mine which is he bes ile size o each machine. Then, we execu e he
p og am dis ibu ing he p ocesses ac oss bo h machines a he same ime. Table 2 shows he
un ime execu ion imes when he ma ix mul iplica ion is execu ed (1) using he bes ile size
ound o He acles in bo h machines (TS-He acles), (2) using he bes ile size ound o Zeus in
bo h machines (TS-Zeus), and (3) using, on each machine, i s bes ile size (Bes TS),
We obse e ha he bes pe o mance is achie ed when he ile size is uned o each machine
independen ly. Unlike p e ious echniques, ou solu ion does no analyze he index domain and
no does i gene a e communica ion code a compile ime. This allows he ile size o be adjus ed
a un ime using di e en app oaches. Fo example, ixing a pa ame ic ile size in an al eady
iled inpu code o using wo ks, such as [16], ha p o ide di e en ways o choose he bes
ile size a un ime. Ou me hod allows he applica ion o his kind o op imiza ion echniques
wi hou changing he communica ion codes. This ea u e is no ound in o he echniques o he
ela ed wo k, whe e he iling is also used as a main ea u e o gene a e he communica ion code,
and hus he ile size canno be changed a un ime.
5.3. Gene al communica ions model s. pa e ns o speci ic applica ions
In his s udy, we e alua e he po en ial o e head ha ou un ime echnique can in oduce
when he numbe o p ocesses is eally high. Each p ocess compu es i s communica ion s uc u e
independen ly. Thus, we can isola e and un he code o compu e he communica ion s uc u es
on a single p ocess wi h he p ope pa ame e s o simula e he calcula ions ha would be done i
18
G oup: 0
Le el: 0 Le el: 1
GP: 0
GP: 1
GP: 4
GP: 5
GP: 8
GP: 9
GP: 12
GP: 13
GP: 2
GP: 3
GP: 6
GP: 7
Le el: 2
27 27 27 2826
19
35
G oup: 1
G oup: 2
G oup: 3
Figu e 10: Applica ion o he p oposed communica ion calcula ion echnique when using a hie a chical QuadT ee map-
ping policy o dis ibu e a ma ix on 64 p ocesses. Whi e ec angula shapes ep esen he WIo he 28- h p ocess (id=27)
o he Jacobi-2D benchma k. C osses poin ou he p ocesses o g oups o p ocesses whose mapped da a do no in e sec
wi h he locally accessed domain a a gi en le el, de ec ing ha no communica ion is needed. Thus, hey a e no checked
a lowe le els.
he applica ion we e o be launched wi h any gi en numbe o p ocesses. This allows us o ob ain
accu a e measu es o he cos o he communica ion de e mina ion alone o a huge numbe o
p ocesses. This s udy has been ca ied ou wi h one MPI p ocess in He acles.
Table 3 shows he execu ion imes ob ained o compu e he communica ions s uc u e o he
Jacobi-2D sol e . I compa es he ac ual ime o calcula e communica ions using he gene al
model, desc ibed in Sec . 4.2, wi h he op imized pa e n desc ibed in Sec . 4.3. Mo eo e ,
we also compa e he ime spen o calcula e communica ions using he gene al model when a
Q ee hie a chical mapping policy is used o dis ibu e he da a. This kind o mapping policies
de ine he dis ibu ion o da a in se e al le els (see Fig. 10). Ou echnique is pe o med a each
hie a chical le el ecu si ely. The calcula ion ope a ions (such as in e sec ions o sub ac ions
among domains) a e only pe o med in he nex lowe / ine le el o he pa s o he cu en le el
whose in e sec ion wi h he local accessed domain is no emp y. Thus, in he example shown in
he igu e wi h 64 p ocesses, ins ead o compa ing wi h all he p ocesses, he compa ison is only
pe o med wi h 4 domains a le el 0, 12 a le el 1 and 12 a he las le el. Using his kind o
hie a chical mapping echniques, we ob ain a calcula ion ime bounded by O(log(P)) o pa e ns
which imply communica ion wi h a cons an numbe o p ocesses.
The ime o he gene al model g ows linea ly wi h he numbe o p ocesso s as expec ed
(P). Fo he Jacobi-2D example, we can see on Tab. 3 ha i is less han wo seconds, e en o
hund eds o housands o p ocesso s. Bu i may be a hind ance o millions o p ocesso s, o i
he pa e n needs o be ecompu ed a di e en i e a ions o a loop. In hese cases, a hie a chical
decomposi ion and nes ed pa allel s uc u es can alle ia e he p oblem, as we can see in he
hi d column o he able. This simula ion e i ies ha he communica ion calcula ion ime when
using he gene al model on a hie a chical mapping policy, is educed o be p opo ional o log(P)
ins ead o P.
On he o he hand, he speci ic pa e n o he Jacobi-2D s encil always compu es neighbo
synch oniza ion pai s, independen ly o he da a sizes o he numbe o p ocesses. The ime in
his case is bounded and negligible.
19
Table 4: Pe o mance (in seconds) ob ained o he h ee benchma ks chosen, o MPI e e ence e sions, and o T asgo
gene a ed codes. Cannon’s algo i hm by design equi es a numbe o p ocesses wi h a pe ec squa e oo .
Illus . case S encil-Op Cannon’s
Machine MPI T asgo MPI T asgo MPI T asgo
He acles-4 14.48 18.04 163.47 196.86 174.31 186.01
He acles-8 17.51 20.28 116.76 118.68 - -
He acles-16 9.79 11.07 57.34 67.43 56.66 62.94
He acles-32 4.42 5.43 35.83 43.61 - -
He acles-64 4.00 5.08 31.62 35.89 16.10 17.95
CETA-4 25.22 27.69 147.79 166.72 173.46 173.71
CETA-8 22.68 23.20 108.94 114.01 - -
CETA-16 11.03 12.72 100.97 111.19 48.89 58.46
CETA-32 6.04 6.64 80.11 86.08 - -
CETA-64 3.37 3.63 57.54 60.20 18.54 21.37
The use o hie a chical mapping policies is an in e es ing ea u e o scale he use o he
p oposed echnique o eally huge amoun s o p ocesses, wi h he a ge o exascale compu ing
in mind. Howe e , o he small sizes o he machines used in he es o ou cu en expe imen al
s udies, we use he gene al model wi h a simple one le el mapping policy, because de e mining
he communica ion pa e ns o hese cases has an unno iceable impac on he pe o mance, o
any model o applica ion.
5.4. Compa ison wi h MPI e e ences
In his s udy, we compa e he p og ams gene a ed by ou T asgo p o o ype wi h MPI p o-
g ams manually de eloped and op imized. In o de o do ha , we pe o m an end- o-end mea-
su e, including c ea ion o da a s uc u es, da a ini ializa ion, and he es o T asgo o e heads.
Fo his s udy we use: he Illus a i e example, an implemen a ion o Cannon’s algo i hm o
ma ix mul iplica ion [17], which is specially de ised o dis ibu ed-memo y sys ems in o de
o minimize he memo y oo p in , and an op imized s encil compu a ion (S encil-Op ) whose
communica ion pa e n is da a-dependen and ede e mined on each ime i e a ion o he s encil
p og am. These benchma ks a e desc ibed in he supplemen a y ma e ial.
Table 4 shows he pe o mance ob ained by he MPI e e ence e sions, and he T asgo gen-
e a ed p og ams o he cases o s udy. We execu e he applica ions on bo h, sha ed-memo y
and dis ibu ed-memo y machines. We see ha T asgo p og ams scale qui e well, bu losing
some pe o mance in compa ison wi h he manually op imized MPI codes (less han 20% in he
wo s cases). This shows ha he implemen a ion o he p oposed echnique in T asgo p oduces
e icien pa allel p og ams a he communica ion le el.
5.5. Compa ison wi h a s a e-o - he-a ool
In his s udy, we ha e pe o med an in-dep h compa ison o he pe o mance be ween he
codes gene a ed by T asgo, and hose gene a ed o dis ibu ed-memo y by Plu o-MPI (dis -
mem), a polyhed al model compile ha includes s a e-o - he-a echniques o gene a ing com-
munica ion code [3]. We choose Plu o-MPI because: (1) Acco ding o he au ho s, i is he
i s wo k ha epo ed an end- o-end ully au oma ic dis ibu ed-memo y pa alleliza ion and
code gene a ion o inpu p og ams and ans o ma ion echniques; (2) I is a ee a ailable ool
easy o ins all, which suppo s all he benchma ks es ed in he pape ; (3) Many esea ch wo ks
ha e appea ed ha use Plu o as baseline, o bo h sha ed- and dis ibu ed-memo y sys ems, such
as [18, 19, 20]; (4) To he bes o ou knowledge, he me hods o Plu o o code gene a ion in
20
Table 5: Maximum a ia ion in he execu ion imes o each benchma k in He acles and CETA, when using T asgo and
Plu o. The a ia ion a io is de ined using he ollowing o mula: ((max ime −min ime)/min ime).
Machine He acles
Benchma ks T asgo Plu o
Ma mul 0.0231 0.0575
Jacobi-2d 0.0490 0.0707
Gauss-Seidel 0.0079 0.0160
Blu -Robe s 0.2670 0.2591
Gem e 0.1192 0.1965
Machine CETA
Benchma ks T asgo Plu o
Ma mul 0.0262 0.0260
Jacobi-2d 0.0196 0.0177
Gauss-Seidel 0.0119 0.0056
Blu -Robe s 0.1695 1.1790
Gem e 0.5304 0.4878
dis ibu ed-memo y sys ems, compa ing wi h o he s, a e he ones ha educe he communica ed
olume o da a, being he gene a ed code pa ame ic in he numbe o p ocesses and p oblem
sizes.
We ha e selec ed i e examples om Polybench [14], a collec ion o examples o be used
o es ing and de eloping polyhed al model compila ion echniques. The examples a e Jacobi-
2D, Gauss-Seidel, Blu -Robe s, he classical Ma ix Mul iplica ion algo i hm and he Gem e
algo i hm. We ha e sligh ly modi ied he s encil in he Jacobi-2D, and he Gauss-Seidel exam-
ples. Ins ead o compu ing a 5-poin s a s encil, we compu e he 4-poin s a s encil o Poisson’s
equa ion. I educes he compu a ion load pe p ocess, leading o a sligh ly bigge impac o
he communica ions, which is he ocus o his s udy. Gauss-Seidel has been selec ed because i
p esen s wa e- on dependences, de i ing in a mac opipeline compu a ion. Ma ix Mul iplica-
ion is a good case o simple linea algeb a p oblems wi h a nice compu a ion/communica ion
balance. On he o he hand, Gem e p esen s a communica ion/compu a ion balance ha is no
adequa e o dis ibu ed-memo y sys ems (communica ion ime ypically is highe han compu-
a ion). The Blu -Robe s il e is a kind o s encil wi h a single i e a ion. They we e chosen o
show he pe o mance o his kind o applica ions when he echniques discussed in his pape
a e applied.
We compile he codes wi h he Make iles p o ided by de aul wi h Plu o-MPI, which include
adis op op ion enabling he use o he FOP communica ions model [2], and a se o de aul lags
and ile size alues o each example. The Make iles o he s encil examples do no include he
lag -l2 iles o enable mul i-le el iling. The in e nal ools used in Plu o-MPI do no handle i in
an a o dable compila ion ime.
Fo ime measu ing, in his s udy, we ha e aken in o accoun only he main compu a ion
and communica ion imes, named Seq. and Comm. in he esul ables. In Plu o-MPI, he ull
ma ices a e alloca ed and ini ialized in all p ocesses, al hough only he pa s needed o he
local compu a ion a e used. A he end o a pa allelized a ine loop nes , Plu o needs o c ea e
again a common global s a e communica ing local esul s o each p ocess. We ha e skipped
all hese imes in ou Plu o-MPI measu es. Fo ai compa ison, we measu e in all he codes as
communica ion imes only: The cos o packing and unpacking he da a, he cos o communica -
ing con ol in o ma ion needed o he communica ions (only in Plu o-MPI), ne wo k la encies
and synch oniza ion wai s. In he compu a ion pa s, we ha e selec ed he same ile sizes on
each example o bo h, T asgo and Plu o codes. In addi ion, in o de o p o ide a measu e o
he s ochas ic delays, we also show in Tab. 5 he maximum a ia ion o he execu ion imes,
o he di e en benchma ks, o each di e en execu ion pla o m, and o each di e en ool
es ed, T asgo and Plu o. The maximum a ia ion is ep esen ed as a a io o he ime di e ence
be ween he minimum and he maximum execu ion imes. We obse e ha in he applica ions
Blu -Robe s and Gem e , whe e he communica ion ime is much highe han he compu a ion
21
Table 6: Main execu ion imes (in seconds) o he i e benchma ks chosen om he Polybench.
Jacobi-2d Gauss-Seidel ma mul Gem e Blu -Robe s
Machine T asgo Plu o T asgo Plu o T asgo Plu o T asgo Plu o T asgo Plu o
He acles-4 142.65 101.67 428.36 274.04 41.34 28.46 0.65 0.75 0.38 1.36
He acles-8 83.35 77.74 253.40 196.52 23.23 14.41 1.15 0.76 0.29 1.49
He acles-16 46.61 58.26 142.41 156.02 13.24 8.26 1.60 0.78 0.24 1.75
He acles-32 24.18 59.47 113.21 138.58 7.23 4.67 1.85 0.83 0.14 2.05
He acles-64 18.51 61.32 64.91 128.11 4.05 2.39 2.52 0.86 0.11 2.77
CETA-4 53.44 38.20 122.13 72.08 26.19 30.31 0.0025 0.0049 0.48 1.61
CETA-8 37.03 28.73 71.08 59.87 13.33 14.36 0.0999 0.0046 0.38 5.67
CETA-16 24.28 33.20 62.86 44.93 7.84 12.73 0.1291 0.0787 0.41 20.88
CETA-32 14.62 21.13 40.25 46.65 7.05 6.84 0.4696 0.1673 0.35 47.05
CETA-64 14.56 24.36 35.09 70.30 7.56 4.13 0.4641 0.1310 0.20 59.78
ime, he a ia ion o he execu ion imes is eally high because o he s ochas ic delays o he
ne wo k. In o de no o ake in o accoun his s ochas ic e ec s o he ne wo k in as uc u e, in
he es o ables, we show he minimum execu ion ime achie ed in he expe imen s.
In Table 6, we p esen he o al execu ion imes (Seq+Comm) conside ed o each benchma k.
In Table 7, we p esen independen ly he accumula ed ime expended in he compu a ion s ages,
and he accumula ed ime expended in he communica ion calcula ion, execu ion, and synch o-
niza ion wai s. Each esul is he measu e ob ained o he p ocess ha expended mo e ime in
he co esponding s ages. Thus, we can obse e in which cases he communica ion cos is highe
o lowe , independen ly o he compu a ion code.
The esul s o he Jacobi-2d, and he Gauss-Seidel examples indica e ha he Plu o code is
mo e e icien o a low numbe o p ocesses, al hough i does no scale as well as T asgo. The
ans o ma ions pe o med by Plu o, including skewing he ime loop o pa allelize, de i e in a
lo o e-u iliza ion and exploi a ion o memo y hie a chies inside he p ocesses. T asgo codes
exploi only spa ial pa allelism a he dis ibu ed-le el, as in classical manual message-passing
app oaches. This de i es in coa se-g ained communica ions, bu ewe oppo uni ies o exploi
compu a ion code op imiza ions inside he p ocesses due o ex a synch oniza ions. I is specially
no iceable o he mac o-pipeline s uc u e om which Gauss-Seidel de i es. Howe e , as he
numbe o p ocesses g ows, Plu o e eals i s mo e clumsy communica ion calcula ions, while he
g anula i y o he T asgo communica ions dec eases, and i s educed cos s o communica ions
become much mo e ele an . See he communica ion cos in Table 7 o hese examples.
In he case o Plu o codes, he ull ma ices a e alloca ed and ini ialized in all he p ocesses.
On he o he hand, T asgo use ac ually dis ibu ed a ays, wi h much lowe memo y oo p in .
Howe e , T asgo needs a communica ion s age be o e he execu ion o each SPMD block. In
he case o he ma ix mul iplica ion, he e is only one execu ion o an SPMD block. Thus,
he communica ion imes in Table 7 o T asgo only include he ime o edis ibu e he da a
needed o each p ocess o compu e i s local pa . On he o he hand, Plu o does no need a
communica ion s age beyond he global s a e consolida ion, which we do no conside in ou
s udy. The communica ion imes o he ma ix mul iplica ion in Plu o codes in Table 7 a e due
mainly o he synch oniza ion imes o exchange con ol in o ma ion in o de o de e mine ha
no communica ion is necessa y o any p ocess.
As expec ed, he Gem e example shows poo scalabili y o bo h T asgo and Plu o dis ibu ed-
22
memo y p og ams. The compu a ional load is eally low, wi h se e al high- olume communi-
ca ion s ages. The cos o execu ing he communica ions is highe han he compu a ion. The
sequen ial algo i hm in he Gem e benchma k is no a good candida e o dis ibu ed-memo y
p og amming in gene al. The pe o mance in his case can be imp o ed in bo h app oaches using
a di e en mapping policy o dis ibu e he compu a ional load among he p ocesses. Howe e ,
his s udy is beyond he scope o his pape .
The Blu -Robe s il e is a kind o s encil p og am wi h a single i e a ion wi h wo SPMD
blocks. The esul s show ha he ans o ma ion o SPMD blocks in o an a ine loop nes pe -
o med by Plu o some imes implies poo pe o mance, specially in dis ibu ed-memo y sys ems,
due o he need o mul iple communica ion s ages in he gene a ed pipeline (al hough locali y
is imp o ed). This e ec is highly no iceable o CETA, he dis ibu ed-memo y clus e (see
Comm. imes in Table 7). Using solu ions such as diamond- iling can alle ia e his p oblem in
Plu o. On he o he hand, T asgo issues a single communica ion s age o each SPMD block,
wi h he expec ed scalabili y.
We conclude ha T asgo codes scale e y well due o hei e y e icien communica ion
s uc u es, despi e he ac ha he compu a ion code can s ill be op imized u he .
6. Rela ed Wo k
The polyhed al model p o ides a o mal amewo k o de elop au oma ic ans o ma ion
echniques a he sou ce code le el [21]. I is applicable o codes based on sequen ial s a ic loops
wi h a ine exp essions (SCoP). All he polyhed al echniques p esen ed so a o dis ibu ed
memo y need o pa ame ize he i e a ion space polyhed a and analyze dependences a compile
ime. The e a e many simila i ies be ween ou wo k and speci ic polyhed al model echniques.
Fo dis ibu ed memo y, he bes communica ion calcula ion me hods so a (see e.g. [2, 22, 3])
compu e communica ions o sequences o a bi a y nes ed loops wi h egula (a ine) accesses,
also known as a ine loop nes s. The loops a e ans o med, iled and inally pa allelized. Com-
munica ions canno be calcula ed ac oss di e en sec ions o a ine loop nes s unless loop usion
can be done. These echniques analyze a compile ime he oo p in o ile da a used by o he
iles. This implies ha he ile size mus be ixed a compile ime and mus be he same o all
he machines in ol ed in he compu a ion. Mo eo e , using hese me hods, he e a e s ill cases
o duplica ed o unnecessa y da a communica ions [2]. The un- ime complexi y o he gen-
e a ed code is dependen on bo h he da a size (numbe o iles) and he numbe o p ocessing
elemen s [23]. Mul i-le el iling echniques can be used o alle ia e he p oblem. Howe e , his
in oduces mo e ile-size decisions a compile ime. Communica ion ac oss iles o di e en le -
els is ha de , and communica ing ac oss iles o bigge sizes inc eases he cases o edundan
communica ions.
The e a e ools which b ing oge he he ad an ages o he polyhed al model and he ask-
o ien ed p og amming p oposals [19]. These app oaches educe global ba ie s in sha ed-memo y
pa allel codes by launching asks and compu ing hei dependences and oo p in s. Howe e , his
app oach also pe o ms a da a pa i ion a e a iling echnique is applied wi h p ede ined ile sizes
a compile ime [11]. The bes ile size depends, among o he hings, on he a chi ec u e de ails
o he a ge machine whe e he p og am will be execu ed [6]. Choosing ile sizes a compile ime
p e en s au oma ic uning o di e en de ices in he e ogeneous en i onmen s. The e a e some
p oposals such as [24] ha gene a e loops ha i e a e o e ull ec angula iles, wi h unknown
pa ame ic ile size. Howe e , i has no been demons a ed ha hese echniques can be applied
wi h he cu en communica ion code gene a o s o dis ibu ed-memo y sys ems.
23
Table 7: Pe o mance (in seconds) o Polybench codes, gene a ed o dis ibu ed-memo y by T asgo, and by Plu o-MPI, b oken down in o compu a ion and communica ion
imes (including calcula ion and execu ion).
Jacobi-2d Gauss-Seidel Ma mul
T asgo Plu o T asgo Plu o T asgo Plu o
Machine Seq. Comm. Seq. Comm. Seq. Comm. Seq. Comm. Seq. Comm. Seq. Comm.
He acles-4 141.94 2.24 85.35 49.43 226.32 283.24 192.41 116.34 41.24 0.09 28.46 2.73
He acles-8 80.56 16.51 52.52 54.18 129.52 182.18 100.02 123.70 23.05 0.18 14.41 2.69
He acles-16 45.93 15.24 30.04 48.48 70.61 108.35 53.80 122.33 13.05 0.20 8.26 3.30
He acles-32 22.46 7.83 22.23 54.78 52.88 95.59 33.25 123.64 7.05 0.18 4.67 3.49
He acles-64 15.57 7.86 22.96 63.83 27.73 55.60 22.11 123.51 3.79 0.26 2.39 2.39
CETA-4 53.04 8.09 36.88 15.98 62.93 75.75 55.65 29.37 26.12 0.07 30.31 2.62
CETA-8 36.17 13.19 22.15 17.87 32.89 52.95 28.60 36.43 13.22 0.13 14.36 2.66
CETA-16 20.24 13.14 16.49 30.66 22.83 53.43 14.39 35.78 6.97 0.89 12.72 6.16
CETA-32 8.63 9.40 8.54 22.41 19.62 36.24 10.79 43.36 5.32 1.79 6.81 5.85
CETA-64 5.77 12.29 8.55 25.39 8.81 33.09 9.35 70.47 5.60 1.69 4.07 4.13
Gem e Blu -Robe s
T asgo Plu o T asgo Plu o
Machine Seq. Comm. Seq. Comm. Seq. Comm. Seq. Comm.
He acles-4 0.1816 0.4648 0.7514 0.7538 0.35 0.03 0.70 0.93
He acles-8 0.1607 0.9864 0.7553 0.7585 0.21 0.06 0.47 1.41
He acles-16 0.1417 1.4545 0.7790 0.7848 0.19 0.07 0.41 1.78
He acles-32 0.0504 1.7966 0.8158 0.8286 0.08 0.05 0.39 2.19
He acles-64 0.0383 2.4794 0.8442 0.8679 0.09 0.05 0.39 3.25
CETA-4 0.0010 0.0019 0.0026 0.0033 0.38 0.10 0.75 1.01
CETA-8 0.0005 0.0989 0.0031 0.0045 0.38 0.15 0.74 4.92
CETA-16 0.0003 0.1203 0.0037 0.0462 0.22 0.33 0.70 19.48
CETA-32 0.0003 0.3767 0.0036 0.1987 0.11 0.30 0.37 45.20
CETA-64 0.0002 0.4633 0.0035 0.1967 0.05 0.16 0.41 53.93
24
O he simila app oaches based on compile- ime in e sec ions o pa ame ic polyhed a ha e
been p oposed o educe da a ans e s in accele a o s, such as FPGAS [25], whe e communi-
ca ions a e calcula ed and op imized only be ween he hos and he accele a o . Dis ibu ed-
memo y p og ams in oduce he complexi y o dealing wi h da a pa i ion policies, and di e en
communica ion pa e ns ac oss a numbe o p ocesses only known a un ime.
The wo k in [26] p esen s a hyb id compile - un ime ansla o scheme, simila o ou ap-
p oach, ha calcula es he communica ion pa e n needed among he SPMD blocks. Howe e ,
hey only suppo egula and epe i i e applica ions whe e he communica ion pa e n is he
same in all he i e a ions o he ou e se ial loop ha encloses he SPMD blocks. This cons ain
is also ound in o he dis ibu ed-memo y app oaches ha in eg a e classical polyhed al ech-
niques o egula codes, wi h inspec o /execu o echniques [27] o suppo i egula o indi ec
da a access exp essions. This inspec o /execu o echnique exchanges con ol da a be o e ac ual
communica ions o a oid a e sing he whole i e a ion space o he pa allelized loop on e e y
p ocess. Unlike ou app oach, hese solu ions canno be used, o example, in ou S encil-Op
benchma k, whe e he communica ion pa e n is ecalcula ed on each i e a ion.
PGAS (Pa i ioned Global Add ess Space) models p esen an abs ac ion o wo k wi h mixed
dis ibu ed- and sha ed-memo y en i onmen s simila o T asgo. The PGAS language ha is
mo e closely ela ed o ou wo k is Chapel [28]. I p oposes a sepa a ion o domain and mapping
modules o gene a e dis ibu ed a ays. Howe e , he bes communica ion agg ega ion me hods
p esen ed so a o Chapel abs ac ions a e es ic ed o speci ic ope a ions, o domain mapping
p ope ies. Fo example, he wo k in [29] is es ic ed o global a ay assignmen s wi h block o
cyclic dis ibu ions. The wo k in [30] p esen s a symbolic subs i u ion o mapping a ibu es in
a ine access exp essions wi h he same inspi a ion as ou app oach. Howe e , he Chapel un ime
canno agg ega e se e al exp essions ac oss di e en loops o gene a e he ull ask oo p in .
Also, i needs o ely on non-agg ega e communica ions when he whole se o da a accessed by
an exp ession is no ully alloca ed in he same emo e p ocesso . I only wo ks o cyclic o
block-cyclic dis ibu ions.
7. Conclusion
This pape p esen s an ex ension o he T asgo pa allel p og amming and compiling ame-
wo k. This ex ension includes echniques ha , o a ine exp essions o pa allelism on da a
accesses, au oma ically de e mine a un ime ad-hoc communica ion pa e ns o dis ibu ed-
memo y p ocesses ac oss wo consecu i e SPMD blocks. This new echnique uses he esul s
o a pa i ion policy o compu e a un ime exac coa se-g ained communica ion pa e ns o dis-
ibu ed message-passing p ocesses. I is based on in e sec ions o emo e and local oo p in s
in e ms o he esul s o he mapping unc ions chosen. Ou app oach allows he au oma ic
gene a ion o p e-compiled mul i-le el pa allel lib a ies o p og ams, ha can adap hei com-
munica ion and synch oniza ion s uc u es o he a ge sys em. Expe imen al esul s, o se e al
ep esen a i e cases o s udy, show ha ou echnique p oduces e icien codes, despi e he o e -
head o ou un ime communica ion calcula ion, compa ed wi h a compile- ime s a e-o - he-a
ool ha gene a es communica ion codes, and wi h manually implemen ed and op imized pu e-
MPI e e ences codes.
Fu u e wo k includes he applicabili y o he ans o ma ion model in he con ex o cu en
polyhed al model amewo ks, using mo e i egula domains, o ex ending i o non-comple ely
a ine exp essions. The new T asgo amewo k is a ailable a h p:// asgo.in o .u a.es.
25