On he un- ime cos o dis ibu ed-memo y
communica ions gene a ed using he polyhed al
model
Ana Mo e on-Fe nandez
Dep . In o m´
a ica,
Uni e sidad de Valladolid
Campus Miguel Delibes, s/n
Valladolid, Spain
Email: [email p o ec ed]a.es
A u o Gonzalez-Esc ibano
Dep . In o m´
a ica,
Uni e sidad de Valladolid
Campus Miguel Delibes, s/n
Valladolid, Spain
Email: a u o@in o .u a.es
Diego R. Llanos
Dep . In o m´
a ica,
Uni e sidad de Valladolid
Campus Miguel Delibes, s/n
Valladolid, Spain
Email: diego@in o .u a.es
Abs ac —The polyhed al model can be used o au oma ically
gene a e dis ibu ed-memo y communica ions o a ine nes ed
loops. Recen ly, new communica ion schemes ha educe he
communica ion olume ha e been p esen ed. In his pape we
s udy he ex a compu a ional e o in oduced a un- ime
by he code gene a ed o manage he communica ion de ails
ac oss dis ibu ed p ocesses. We ocus on he mos sophis ica ed
communica ion scheme so a in oduced ( he FOP scheme). We
p esen an asymp o ic cos s udy o he FOP scheme in e ms
o wo main un- ime pa ame e s: The p oblem size, and he
numbe o p ocesso s. Based on his s udy, we iden i y scalabili y
limi a ions in cu en implemen a ions o hese echniques, and
p opose a simple implemen a ion al e na i e o elimina e one o
hem. Expe imen al esul s a e p esen ed, showing he po en ial
impac on pe o mance o hese implemen a ion limi a ions when
using hese codes in la ge pa allel sys ems.
Keywo ds—Polyhed al model, dis ibu ed-memo y, un- ime,
complexi y
I. INTRODUCTION
The polyhed al model has been p o ed o be a use ul ool o
ans o m and gene a e pa allel p og ams o codes wi h a ine
nes ed loops [1]. I can also be used o au oma ically gene -
a e code o dis ibu ed-memo y pla o ms. The dependence
analysis suppo ed by he model allows o gene a e code ha
iden i ies which alues should be communica ed ac oss p o-
cesses, packing/unpacking he da a, and execu ing he p ope
communica ion ope a ions [2]. G iebl [5] p esen s a model
o use he polyhed al model o e dis ibu ed sys ems. How-
e e , his echnique p oduces many edundan communica ion.
Recen ly, se e al communica ion schemes has been p esen ed
in o de o educe he olume o da a communica ed [3], [6].
The au oma ically gene a ed codes a e capable o coo dina ing
he compu a ion and communica ion ac oss he e ogeneous
de ices. This allows he exploi a ion o pa allelism in he -
e ogeneous clus e s wi h GPUs o o he accele a o s, which
is he cu en end o build huge pa allel sys ems [8]. The
scale o he machines and p oblems ha can be cu en ly
aced, g ows by se e al o de s o magni ude compa ing wi h
hose ound in mos pe o mance e alua ions done p e iously
wi h dis ibu ed-memo y polyhed al gene a ed codes. And i
will con inue g owing up, wi h exascale compu ing being an
impo an esea ch ocus.
In his pape we s udy he codes gene a ed by he mos
sophis ica ed communica ion scheme in oduced so a ( he
FOP scheme [6]). We p esen a s udy o he ex a cos
in oduced a un- ime by he gene a ed codes o manage he
communica ions. We do an asymp o ic complexi y analysis in
e ms o wo main un- ime pa ame e s: The p oblem size N,
measu ed as he numbe o da a elemen s o be p ocessed,
and he numbe o p ocesso s P. Ou complexi y model
highligh s some po en ial scalabili y limi a ions in e ms o
he p oblem size N, and he numbe o p ocesso s P, in he
cu en implemen a ions o hese echniques. Speci ically, we
ocus on he implemen a ion included in he cu en Plu o
compile amewo k [4]. We iden i y and isola e one o his
limi a ions, ela ed o he applica ion o he dis ibu ion policy
used o schedule he i e a ions o a pa allelized loop. We
discuss ha o de e minis ic dis ibu ion policies, a simple
implemen a ion al e na i e, p e iously exploi ed in Hi map [7]
(a un- ime lib a y o managemen o dis ibu ed hie a chical
iling a ays), elimina es his speci ic scalabili y p oblem.
We p esen h ee cases o s udy: 1-D and 2-D Jacobi sol e s,
and Floyd-Wa shall algo i hm. The e e ence codes o hese
h ee examples a e included in Polybench [9]. Ou expe imen-
al esul s show he po en ial high impac o hese scalabili y
limi a ions on he pe o mance o he applica ion, o huge
da a sizes and clus e s. We also show how he p oposed
implemen a ion al e na i e highly alle ia es he pe o mance
p oblems, as p edic ed by ou complexi y model.
The es o he pape is o ganized as ollows: Sec ion 2
summa izes he main ea u es o FOP-scheme gene a ed codes.
Sec ion 3 p esen s he cos model. Sec ion 4 discusses ou
implemen a ion al e na i e. Sec ion 5 desc ibes he applica ion
o he model o one case o s udy. Sec ion 6 p esen s he
expe imen al s udy and esul s. Sec ion 7 concludes he pape .
II. THE COMMUNICATION SCHEME
The FOP communica ion scheme is a model o he au-
oma ic gene a ion o communica ion code o dis ibu ed-
memo y pa allel p og ams, in he con ex o polyhed al code
ans o ma ions. I is based in dependence analysis ac oss iles,
o pa allel p og ams ha dis ibu e he i e a ions o iled
loops. I has been designed o educe he olume o da a
communica ed ac oss p ocesso s, compa ed wi h o he s a e-
o - he-a sys ems o au oma ically gene a e communica ion
code o a ine-loop nes s. I has been also p o ed ha i
p o ides good pe o mance o small and medium sized da a
se s, and small numbe o dis ibu ed-memo y p ocesses [6].
In ha wo k, he au ho s desc ibe he concep ual app oach and
he solu ion design in de ail. An implemen a ion o his scheme
is included in he cu en e sion o he Plu o compile [4].
In his sec ion, we summa ize he main ea u es o he FOP
scheme, and some design and implemen a ion decisions o he
gene a ed codes.
The code dedica ed o calcula e and execu e communica-
ions is inse ed a he end o dis ibu ed loops which o igina e
a communica ion need. The analysis o RAW dependences
is done sepa a ely o each di e en da a s uc u e (a ay
a iable) in ol ed in a dis ibu ed loop nes . Fo a gi en
i e a ion o a dis ibu ed loop, he FO ( low-ou ) se is de ined
as he se o da a gene a ed/w i en du ing he i e a ion, ha
is equi ed/ ead du ing he execu ion o o he i e a ions. A
compile ime, he FOP scheme de e mines a dependences
pa i ion. A pa i ion o he low-ou se in e ms o subse s
o a ge i e a ions ha can be loca ed in di e en iles. Each
pa o he se is ea ed independen ly, leading o a di e en
piece o communica ion code. This pa i ion is applica ion
dependen . The gene a ed code o a gi en pa i ion can ha e
wo di e en la o s. Mul icas ope a ions ha send all he
da a in a pa i ion o e e y p ocesso ha equi es da a om i .
And Unicas ope a ions ha issue a di e en communica ion
o each i e a ion ha equi es da a om his pa i ion. To
a oid sending mo e han one ime he same da a o he
same p ocesso , unicas ope a ions a e only chosen when i
is possible o de e mine, a compile o un- ime, ha he
ecei ing i e a ions a e all scheduled a di e en p ocesso s.
The au ho s o FOP p opose some ules o de e mine when
i is sa e o in oduce unicas ope a ions, mul icas being he
de aul choice. In his wo k we will ocus on he de aul and
less complex mul icas ope a ions.
Fo mul icas ope a ions, FOP in oduces one piece o code
o each dis ibu ed loop, pa o he dependences pa i ion, and
a ay a iable. The code uses se e al auxilia y da a s uc u es.
One da a bu e pe p ocesso in ol ed in he compu a ion, o
s o e he da a o be sen . One single ecei e bu e o s o e
all he da a ecei ed om o he p ocesso s. Two coun e s pe
p ocesso , o s o e he amoun o da a o be sen o ecei ed
o/ om a emo e p ocesso . Each piece o code con ain h ee
s ages:
1) Pack: Pack da a while iden i ying a ge p ocesso s. The
i e a ions space o he dis ibu ed loop assigned o he
local p ocesso is a e sed again. Fo each i e a ion a
unc ion gene a ed wi h applica ion speci ic in o ma ion
(σ) is used o iden i y which o he i e a ions (and hus,
p ocesso s) equi e da a om his pa i ion and i e a ion.
The da a is packed (copied) in o he co esponding ou pu
bu e s and he coun e s ha measu e he da a o be sen
o each o he p ocesso a e upda ed.
2) Coo dina ion and communica ion: In e change o com-
munica ion sizes ac oss p ocesso s, and issue he equi ed
poin - o-poin communica ions. The coo dina ion s ep is
done wi h he s anda d all- o-all MPI collec i e ope a ion.
Each p ocesso sends he alue o each ou pu -bu e
coun e o he co esponding ecei ing p ocess. This
in e change a oids he need o a e se he i e a ion space
scheduled on any o he p ocesso , doing he same analysis
as o packing, only o ob ain he sizes o da a ha we
expec o ecei e om each o he p ocesso . A e he
coo dina ion s ep, asynch onous send and ecei e ope a-
ions a e issued o each p ocesso wi h a non-ze o alue
in he co esponding coun e . Wi h all he ecei e coun e s
a ailable, i is possible o compu e displacemen s o use
one single bu e o all he ecei e ope a ions.
3) Unpack: Unpack ecei ed da a. The whole i e a ion space
o he dis ibu ed loops is a e sed iden i ying which
i e a ions a e scheduled on emo e p ocesso s o which
we ha e ecei ed da a. Fo each one o hese i e a ions i
is es ed i he local p ocesso is one o he ecei e s o he
da a o his pa i ion and i e a ion (again wi h a unc ion
speci ically gene a ed o his applica ion). In ha case,
he da a is unpacked om he bu e o he ac ual a ay
a iable.
III. COST MODEL
In his sec ion we p esen a cos model o he un- ime
compu a ional e o o he communica ion managemen code
in oduced by he FOP scheme [6]. I is he scheme o poly-
hed al model compu a ions wi h less communica ion olume
in oduced so a . We use as e e ence o speci ic design
decisions he codes in oduced by he cu en implemen a ion
o Plu o. Ou model measu es he asymp o ic cos o he
calcula ions needed o issue he communica ions in e ms
o wo un- ime pa ame e s: Numbe o p ocesso s (P), and
P oblem size (N), measu ed as he numbe o da a elemen s
o be p ocessed. The model does no ake in o accoun he
ac ual cos o he communica ions, which is dependen on
ex e nal ac o s ela ed o he pla o m and he communica-
ion opology. We model only he ex a cos s in oduced by
he au oma ically gene a ed code o p epa e and launch he
communica ion ac i i ies (calcula ions o pack/unpack, and
o he local coo dina ion ac i i ies). The objec i e is o ind
scalabili y limi a ions in oduced when applying he scheme,
ha could be elimina ed by design o implemen a ion changes.
As commen ed in he p e ious sec ion, in his wo k we will
ocus on he cos o he lowe complexi y mul icas ope a ions.
Unicas ope a ions, ha may in oduce u he limi a ions, will
be co e ed in an ex ended u u e wo k.
An exce p o he code gene a ed o a 1D Jacobi pa allel
p og am included in he Polybench [9] benchma k is shown
in Fig. 1. I will be used as example while discussing he cos
model.
A. Gene al cos o a dis ibu ed loop
Fo simplici y o he discussion, le us conside a single
dis ibu ed loop L, wi h an i e a ion index . Le D( )⊂ P(Z)
be he Domain o , he subse o i e a ions ha a e a e sed
by he loop index. Le O(L◦)be he uppe -bound o he un-
ime cos o calcula ing he communica ions needed o all he
i e a ions o D( )scheduled o a gi en p ocesso . F om now
on, cons an ac o s ha a e applica ion dependen , and no
a ec ed by un- ime pa ame e s will be deno ed wi h cname,
whe e he name will be a single lowe -case le e .
Each combina ion o dis ibu ed-loop, a iable, and pa o
he dependences pa i ion, leads o one Ins ance o communi-
ca ion code. In Fig. 1, lines 9 o 35 a e he code gene a ed
o one communica ion ins ance associa ed o a ay a iable b.
Lines 36 o 61 con ain a simila code, gene a ed o a second
communica ion ins ance associa ed o a ay a iable a.
Le czbe he numbe o combina ions o di e en a ay
a iables, and pa s o he dependences pa i ion in oduced by
he FOP scheme o he loop L, and O(y◦) he uppe -bound
o he cos o one communica ion ins ance.
O(L◦) = cz×O(y◦)
The cos o one ins ance o communica ion y◦is he sum o
he cos s o i s h ee consecu i e s ages p e iously desc ibed:
packing, coo dina ion, and unpacking. In he ollowing sec-
ions we model he execu ion o one ins ance o code o a
gene ic i e a ion o he ou e loops. We will ocus on he ou e
loops i e a ions ha lead o maximum pa allelism po en ial o
he dis ibu ed loops.
B. P oblem size and numbe o i e a ions
The loops pa allelized by he polyhed al model ools ep e-
sen a ans o med space o he o iginal loops. The ca dinali y
o he i e a ions se o a dis ibu ed loop is a unc ion o
he p oblem size |D( )|= (N), ha can be de e mined in
e ms o he ans o ma ions applied. We a e mainly in e es ed
in he loops whe e he ca dinali y g ows asymp o ically wi h
N, allowing o exploi mo e pa allelism o bigge p oblem
sizes. Some cons an s a e in oduced by he ans o ma ions
ha educe he o e all cos . Fo example, when iling is
applied, he ile size c appea s as a di iso |D( )|= (N/c ),
wi h an asymp o ic uppe -bound s ill in he o de o N:
|D( )| ∈ O(N).
C. Dis ibu ion policy
ADis ibu ion Policy unc ion Π : D( ),N→ P(D( ))
is used o de e mine he subse o a domain D( ) ha is
scheduled on a p ocesso ank p∈[0, P −1]. The In e se Dis-
ibu ion Policy unc ion, π:Z→[0, P −1], maps each index
o he domain o he co esponding p ocesso ank. In gene al,
dis ibu ion policies y o ob ain a good load balance. Thus,
we assume ha he numbe o i e a ions scheduled on each
p ocesso is simila : ∀p∈[0, P −1],|Π(D( ), p)| ≃ (N)/P.
The un- ime cos o applying hese unc ions is deno ed wi h
Π◦, and π◦ espec i ely.
The unc ion Πis used o compu e he i e a ions o he
loop scheduled o he local p ocess. In he example code o
Fig. 1, lines 6 and 7 calcula e he lowe and uppe limi s
o he i e a ion space o be dis ibu ed lb dis and ub dis .
These a e he inpu s o he Π unc ion implemen ed in he
poly loop dis unc ion. The ou pu s, lbp and ubp, a e he
lowe and uppe limi s o he locally scheduled i e a ions.
The cos o his unc ion is associa ed o he pa alleliza ion
o he algo i hm. Thus, we do no conside i as a speci ic
cos in oduced by he communica ion calcula ions.
D. Packing s age
The packing s age a e ses he subse o locally scheduled
i e a ions. See he loop in lines 9 o 15 in Fig. 1, ha a e ses
i e a ions om lbp o ubp. I has wo main con ibu ions o he
o e all cos ha a e compu ed o each i e a ion conside ed.
Fi s , each i e a ion applies a unc ion σ(i), speci ically
gene a ed o each applica ion, o ob ain he lis o ecei ing
p ocesso s. The sigma unc ion con ains a cons an numbe
o condi ionals cc, dependen on he applica ion sou ce code.
Each condi ional po en ially applies π o ob ain he ank o he
p ocesso ha has a a ge i e a ion. Thus, ob aining he a ge
p ocesso s o all he i e a ions scheduled o a p ocesso , is
done in (N)/P ×cc×π◦.
Second, each i e a ion a e ses he lis o p ocesso s o
de ec he ones ha should ecei e da a om he local p ocess.
This is done in O(P), wi h a e y small cons an cs, as i
execu es a simple condi ional. See line 12 in Fig. 1. The ac ual
packing ope a ion is done only o he p ocesso s de ec ed as
ecei e s (condi ion e alua ed o ue). The code o packing
da a in he ou pu bu e s is applica ion dependen and only
a e ses he da a ha is going o be sen . Howe e , da a
a e packed (copied) in a di e en bu e o each ecei ing
p ocesso . Thus, he e could be mul iples copies o he same
da a. In he wo s case, all p ocesso s should ecei e he same
da a. This is dependen on he communica ion s uc u e o he
applica ion. Fo example, neighbo synch oniza ion communi-
ca ions ha e O(1) numbe o p ocesso s in ol ed o each da a
subse , while some communica ions in LU educ ions esul
in O(P)p ocesso s in ol ed. Le us model he ca dinali y o
he numbe o communica ions wi h an h- ela ion unc ion
h(P). Le c be he mean olume o da a o be sen by one
1i ((N >= 1) && (T >= 1) && (N >= 4)) {
2 o ( 2 = -1; 2 <= loo d (3 *T + N - 4, 32); 2++) {
3/*Sequen ial Code */
4.....
5/*End sequen ial code */
6_lb_dis = max (ceild (2 * 2, 3), ceild (32 * 2 - T + 1, 32));
7_ub_dis = min (min ( loo d (2 *T + N - 4, 32), loo d (64 * 2 + N + 60, 96)), 2);
8poly _loop_dis (_lb_dis , _ub_dis , np ocs, my_ ank, &lbp, &ubp);
9 o ( 4 = lbp; 4 <= ubp; 4++) {
10 clea _sende _ ecei e _lis s (np ocs);
11 sigma_b_1_0 ( 2, 4, T, N, my_ ank, np ocs);
12 o (__p = 0; __p < np ocs; __p++) {
13 i ( ecei e _lis [__p] != 0) {
14 send_coun s_b[__p] = pack_b_1_0 ( 2, 4, send_bu _b[__p], send_coun s_b[__p]);
15 }}}
16 i ( 2 <= loo d (3 *T + N - 5, 32)) {
17 MPI_All oall (send_coun s_b, ..., ec _coun s_b, ...);
18 eq_coun = 0;
19 o (__p = 0; __p < np ocs; __p++)
20 i (send_coun s_b[__p] >= 1)
21 MPI_Isend (send_bu _b[__p], send_coun s_b[__p],... );
22 o (__p = 0; __p < np ocs; __p++)
23 i ( ec _coun s_b[__p] >= 1)
24 MPI_I ec ( ec _bu _b + displs_b[__p], ...);
25 MPI_Wai all ( eq_coun , eqs, s a s);
26 o (__p = 0; __p < np ocs; __p++) {
27 send_coun s_b[__p] = 0;
28 cu _displs_b[__p] = displs_b[__p];
29 } }
30 o ( 4 = _lb_dis ; 4 <= _ub_dis ; 4++) {
31 p oc = pi_0 ( 2, 4, T, N, np ocs);
32 i ((my_ ank != p oc) && ( ec _coun s_b[p oc] > 0)) {
33 i (is_ ecei e _b_1_0 ( 2, 4, T, N, my_ ank, np ocs) !=0) {
34 cu _displs_b[p oc] = unpack_b_1_0 ( 2, 4, ec _bu _b, cu _displs_b[p oc]);
35 }}}
36 o ( 4 = lbp; 4 <= ubp; 4++) {
37 clea _sende _ ecei e _lis s (np ocs);
38 sigma_a_1_0 ( 2, 4, T, N, my_ ank, np ocs);
39 o (__p = 0; __p < np ocs; __p++) {
40 i ( ecei e _lis [__p] != 0) {
41 send_coun s_a[__p] = pack_a_1_0 ( 2, 4, send_bu _a[__p], send_coun s_a[__p]);
42 }}}
43 MPI_All oall (send_coun s_a, ..., ec _coun s_a, ...);
44 eq_coun = 0;
45 o (__p = 0; __p < np ocs; __p++)
46 i (send_coun s_a[__p] >= 1)
47 MPI_Isend (send_bu _a[__p], ...);
48 o (__p = 0; __p < np ocs; __p++)
49 i ( ec _coun s_a[__p] >= 1)
50 MPI_I ec ( ec _bu _a + displs_a[__p],...);
51 MPI_Wai all ( eq_coun , eqs, s a s);
52 o (__p = 0; __p < np ocs; __p++) {
53 send_coun s_a[__p] = 0;
54 cu _displs_a[__p] = displs_a[__p];
55 }
56 o ( 4 = _lb_dis ; 4 <= _ub_dis ; 4++) {
57 p oc = pi_0 ( 2, 4, T, N, np ocs);
58 i ((my_ ank != p oc) && ( ec _coun s_a[p oc] > 0)) {
59 i (is_ ecei e _a_1_0 ( 2, 4, T, N, my_ ank, np ocs) !=0) {
60 cu _displs_a[p oc] = unpack_a_1_0 ( 2, 4, ec _bu _a, cu _displs_a[p oc]);
61 }}}
62 } }
Figu e 1. Exce p o he gene a ed communica ion code o a 1D Jacobi sol e using he FOP scheme
i e a ion o he a ay a iable conside ed. This is ypically
a cons an de e mined by he applica ion and ans o ma ions
applied. Thus, he cos o his second pa o he packing s age
can be es ima ed wi h: (N)/P ×(cs×P+c ×h(P)). The
o e all cos o he whole packing s age is es ima ed as:
pack◦= (N)/P ×(cc×π◦+cs×P+c ×h(P))
E. Coo dina ion and communica ion s age
The coo dina ion s age includes se e al ac ions, see lines
16 o 19 in Fig. 1. I s a s wi h an MPI all- o-all collec i e
communica ion ope a ion o in e change coun e s. In gene al,
his ype o all- o-all communica ions a e assumed o be
done in O(P). Then, he ac ual poin - o-poin communica ions
needed a e launched a e sing he p ocesso anks in O(P).
The ac ual cos o he communica ions is no modelled o his
wo k, only he p epa a ion and launching ac i i ies. Finally, a
las loop is execu ed ha also a e ses he p ocesso anks in
O(P) o simple bookeeping ope a ions. We model he o e all
cos o his s age (wi hou ac ual communica ion cos s) by:
coo d◦=P
F. Unpacking s age
The da a ecei ed om a p ocesso has been packed in
i e a ion o de . Thus, hey should be unpacked in he same
o de . See lines 30 o 35 in Fig. 1. This s age a e ses he
whole i e a ion space o he dis ibu ed loop ( om lb dis
o ub dis in he example code), using he π unc ion o
de e mine which i e a ions a e scheduled in emo e p ocesso s.
The cos o his ope a ion is modelled wi h (N)×π◦.
A second pa o he cos appea s only o i e a ions on
emo e p ocesso s om which da a has been ecei ed a he
local p ocess du ing he communica ion s age. In he wo s
case his condi ion check, o a gi en i e a ion, can be sa is ied
o all he es o Pp ocesso s. Bu we can model again he
numbe o i e a ions ha a e going o be de ec ed as alid
ac oss he whole space wi h he h- ela ion unc ion h(P)o
he applica ion. Each locally scheduled i e a ion p oduces a
mean o h(P)communica ions ecei ed om o he i e a ions.
Fo hese se o alid i e a ions, a second check is done wi h a
applica ion ailo ed unc ion ha con ains one o mo e pieces
o code (a cons an numbe cdo hem, dependen on he
sou ce code) and in e nally applying he π unc ion. Finally,
he ac ual unpack ope a ion is done only once o each da a
elemen , and he cos di ec ly depends on he olume o da a
communica ed . The o e all cos o he whole unpacking
s age is modelled by:
unpack◦= (N)×π◦+ (N)/P ×h(P)×(π◦+cd+ )
G. To al cos
Ou inal cos model is dependen on wo unc ions, and
some cons an s, ha should be de e mined o each appli-
ca ion: (N),h(P), ,cc,cd. As we a e mainly in e es ed
in he asymp o ic beha iou , i should no be di icul o
de e mine he o de o he unc ions in e ms o Nand P.
The cons an s only gi e us a ough idea o he weigh o
each pa o he o mula, bu hey canno be conside ed alone
o a eally p ecise model, as he amoun o a i hme ical
ope a ions gene a ed by he loop ans o ma ions o access he
da a elemen s, pack/unpack hem, and simila ope a ions has
no been conside ed.
The o e all cos o calcula ing a gene ic communica ion
ins ance y◦, can be es ima ed as he accumula ion o he h ee
s ages: y◦=pack◦+coo d◦+unpack◦.
y◦= (N)/P ×(cc×π◦+cs×P+c ×h(P))
+P
+ (N)×π◦+ (N)/P ×h(P)×(π◦+cd+ )
A e mul iplica i e cons an ac o s elimina ion, and some
simpli ica ion he asymp o ic uppe -bound can be modelled as:
O(y◦) = O( (N)×π◦+ (N)/P ×π◦×h(P) + P)
IV. PROPOSAL: IMPLEMENTATION ALTERNATIVE
As i can be obse ed in he cos model o mula, a key
ope a ion is he iden i ica ion o he p ocesso ha owns an
i e a ion o he dis ibu ed loop, using he in e se dis ibu ion
policy unc ion π. I appea s se e al imes in he cos model,
as a mul iplie ac o .
Gi en an unknown dis ibu ion policy unc ion Π, a simple
way o build πis o execu e a loop ha applies Π o
each p ocesso ank, checking i he i e a ion pa ame e is
in he esul ing se . See pseudocode in Fig. 2 (le ). The
implemen a ion solu ion p e iously p esen ed, makes he π
unc ion independen on he Πpolicy implemen ed, as a
as he policy e u ns a block o con iguous i e a ions. The
un- ime Poly lib a y e sion included in he cu en Plu o
dis ibu ion, con ains only one Π unc ion: A classical block
dis ibu ion policy. See pseudocode in Fig. 2 (middle).
Wi h he cu en Plu o’s implemen a ion, he cos o he
unc ions is: O(Π◦) = O(1), and O(π◦) = O(P). Fo
mo e gene ic pa i ion policies he cos may inc ease, because
checking i an index is inside a block ange can be done in
O(1), bu o a gene ic se o nindexes he sea ch cos is a
leas O(log n)i i is so ed, o O(n)i i is is no . In his las
case he cos o πcould go up o O(π◦) = O(P×N).
We p opose o use a di e en app oach p e iously used
in Hi map [7], a un- ime lib a y o dis ibu ed hie a chical
iling a ays managemen . In Hi map, he p og amme o he
dis ibu ion policies is o ced o de elop plug-ins ha include
bo h he di ec and he in e se dis ibu ion policies unc ions.
In Hi map, he classical pa i ion policies (block, cyclic, e c.)
ha e implemen a ions whe e he cos o Π◦and π◦is qui e
simila , and i is always O(1). This solu ion can be exploi ed
o any de e minis ic dis ibu ion policy based on an in e ible
unc ion. Fo non-in e ible unc ions he p og amme may
unc ion pi(Dom d,in i,in P)
do p = 0, P-1
d’ = PI(d,p)
i i in d’ hen e u n p
enddo
unc ion PI(Dom d,in p,in P)
i (p < |d|%P)
.lb = d.lb + (|d|/P)*p+p
.ub = .lb + (|d|/P)
else
.lb = d.lb + (|d|/P)*p + |d|%P
.ub = .lb + (|d|/P) - 1
endi
e u n
unc ion pi_Al (Dom d,in i,in P)
o = i - d.lb;
lim = (|d|/P + 1)*(|d|%P)
i ( o < lim )
e u n o /(|d|/P + 1)
else
e u n (o -lim)/(|d|/P) + |d|%P
endi
Figu e 2. Pseudo-codes o he o iginal π(le ) and Π(middle) unc ions, and ou al e na i e implemen a ion p oposed o π( igh ). Dom < lb, ub >
ep esen s a uple wi h he lowe and uppe bound o a con iguous 1-dimensional i e a ion space. |d|=d.ub −d.lb + 1 ep esen s he domain ca dinali y.
chose o pay he ex a un- ime cos ac o , o pay an ex a
memo y cos . I is always possible o s o e in an a ay he index
o he assigned p ocesso o all he elemen s in he i e a ion
space, keeping he O(1) un- ime cos o he π unc ion.
We ha e in oduced in Poly (Plu o’s un ime helpe unc-
ions) a di ec implemen a ion o he in e se dis ibu ion
policy o block pa i ions, elimina ing a mul iplie ac o o
Pin se e al s ages o he communica ion calcula ion. See
pseudocode in Fig. 2 ( igh ).
The asymp o ic impac o his change can be seen in he cos
model. A e subs i u ing he cos s o he π unc ion de i ed
om he cu en implemen a ion, he esul is:
O(y◦)=( (N)×P+ (N)×h(P) + P)
Wi h he al e na i e implemen a ion, mul iplie P ac o s com-
ing om he π unc ion disappea :
O(y◦) = ( (N) + (N)/P ×h(P) + P)
I is specially ema kable ha in he o iginal implemen a ion,
he size p oblem is mul iplied by he numbe o p ocesso s du -
ing he unpacking s age. In he ollowing sec ions we p esen
empi ical e idence o he impac o c ea ing a speci ic π
unc ion o each dis ibu ion policy Π, di ec ly implemen ing
he in e se unc ion wi h a cos bounded by O(1).
V. CASE STUDY: 1-D JACOBI
To show how o apply he cos model, we ha e chosen as
case s udy he 1-dimensional Jacobi p og am. This applica ion
is a good example o s udy because he code p oduced by Plu o
includes only one dis ibu ed loop wi h mul icas ope a ions, i
has a simple neighbo synch oniza ion communica ion s uc-
u e, and i is e y easy o ind p ope app oxima ions o he
applica ion dependen unc ions.
A. Cos model pa ame iza ion
The code has been gene a ed using he de aul ile sizes
in he Plu o example (c = 32 i e a ions o any iled loop).
The unc ion ha compu es he numbe o i e a ions in he
ans o med pa allel loop, has wo inpu pa ame e s: (N, T).
Whe e Nis he a ay size, and Tis he numbe o i e a ions
o he o iginal sequen ial code be o e ans o ma ions. The
o mula used o compu e he limi s o he dis ibu ed loop
index 4depend on he alue o he ou e loop index 2
(see lines 6 o 8 in Fig. 1). These loops c ea e a pipelined
execu ion. Du ing he applica ion p og ess, he amoun o
dis ibu ed i e a ions o he 4loop g ows, i keeps s able o a
while, and hen dec eases. The maximum deg ee o pa allelism
ob ained in he s able phase is ela ed o he p oblem size
pa ame e s, and can be app oxima ed wi h: i (3T≥N), hen
(N, T )≃0.01N; i (3T < N), hen lim (N, T)T→∞ =
3.125T. Thus, (N, T )g ows linea ly wi h he p oblem size
pa ame e s O( (N, T )) = O(min(N, 3T)). Fo simplici y,
le us assume ha Tis always big enough o ob ain he
maximum deg ee o pa allelism o a gi en inpu a ay size.
Thus, O( (N)) = O(N).
The e a e wo communica ions ins ances, one o a ay a,
and one o a ay b. Thus, cz= 2. The h- ela ion unc ion
h(P)is ypically O(1) in neighbo synch oniza ion applica-
ions. Indeed, expe imen al measu es wi h he gene a ed code
o he 1-D Jacobi p og am show ha he mean alues o he h-
ela ion ac oss i e a ions and p ocesso s a e: ¯
h(P)≃1 o he
code ins ance gene a ed o he a ay a, and ¯
h(P)≃0.25 o
he code ins ance gene a ed o he a ay b. The da a olume
c communica ed by each dis ibu ed i e a ion has been also
measu ed: c = 188 da a elemen s o aa ay, and c = 63
da a elemen s o ba ay. Inspec ing he gene a ed code, we
obse e ha he o he cons an alues a e he ollowing. Fo
he aa ay cc= 7, cd= 7, and o he ba ay cc= 1, cd= 1.
Fo an asymp o ic beha iou s udy, we can ne e heless
igno e he applica ion cons an s, and simpli y he esul ing
model o he o e all cos o he communica ions needed o
one i e a ion o he ou e loop as:
O(L◦) = O(N×π◦+N/P ×π◦+P)
A e subs i u ing he cos s o he π unc ion de i ed om
he cu en implemen a ion, he esul is:
O(L◦) = O(N×P+N+P)
Wi h he al e na i e implemen a ion, mul iplie P ac o s com-
ing om he π unc ion disappea :
O(L◦) = O(N+N/P +P)
B. Simula ion s udy
Doing eal expe imen s o big da a sizes, and la ge numbe
o p ocesso s, may equi e a huge amoun o compu a ion
ime in c i ical supe compu e in as uc u es. Fo una ely, we
1e-06
1e-05
0.0001
0.001
0.01
0.1
1
10
100
1000
1000 5000 10000 50000 100000 500000 1000000
Time (sec.)
N (a ay size)
O iginal, 1d-Jacobi ( P = 5,000 )
Packing
Unpacking
Compu a ion
1e-06
1e-05
0.0001
0.001
0.01
0.1
1
10
100
1000
1000 5000 10000 50000 100000 500000 1000000 5000000
Time (sec.)
N (a ay size)
Al e na i e, 1d-Jacobi ( P = 5,000 )
Packing
Unpacking
Compu a ion
Figu e 3. Execu ion imes wi h he o iginal and al e na i e π unc ion wi h di e en p oblem sizes N
1e-06
1e-05
0.0001
0.001
0.01
0.1
1
10
100
1000
100 500 1000 5000 10000
Time (sec.)
P (# p ocesso s)
O iginal, 1d-Jacobi ( N = 1,000,000 )
Packing
Unpacking
Compu a ion
1e-06
1e-05
0.0001
0.001
0.01
0.1
1
10
100
1000
100 500 1000 5000 10000
Time (sec.)
P (# p ocesso s)
Al e na i e, 1d-Jacobi ( N = 1,000,000 )
Packing
Unpacking
Compu a ion
Figu e 4. Execu ion imes wi h he o iginal and al e na i e π unc ion wi h di e en numbe o p ocesses P
can modi y he codes gene a ed by Plu o o simula e a gi en
amoun o he ou e loop i e a ions in a chosen p ocesso , wi h
he desi ed Nand Ppa ame e s, using a educed amoun
o memo y. This allow us o pe o m an empi ical s udy o
in es iga e he e ec s o scaling he Nand Ppa ame e s
o sizes ha esemble high-end supe compu e s. Expe imen al
esul s in a smalle eal case and machine a e p esen ed in
Sec . VI.
The modi ica ions need in he gene a ed code o he 1D
Jacobi example include: (1) Adding some code o ead pa am-
e e s o he chosen limi s o he ou e loop ( 2index); (2)
change he decla a ions o he aand ba ays o ha e a small
ixed size (4096 elemen s); (3) modi y all a ay accesses o
use he esul ing index modulo 4096 o s ay in o he ixed
a ays bounda ies; (4) elimina e he MPI calls; (5) a he
s a o each 2i e a ion, locally compu e he send coun e s
o all he emo e p ocesso s, in o de o simula e he all-
o-all MPI communica ion elimina ed, ha coo dina es he
communica ion sizes ac oss p ocesso s. This is done ou o
he code sec ions ha a e measu ed wi h ime coun e s.
We p ese e he same ime measu ing mechanisms included
in he o iginal code, o he compu a ion sec ion, and o
each one o he h ee communica ion calcula ion s ages. The
da a esul s p oduced by his simula ion a e no co ec . The
communica ion codes pack and unpack dummy alues in
he bu e s, and in he cons ic ed a ays. Bu all he com-
munica ion p epa a ion calcula ions, and packing/unpacking
ope a ions, a e done exac ly as in he o iginal code. Thus,
he ime measu es a e consis en wi h he eal case, excep o
he ac ual communica ion cos s which a e in en ionally no
included o conside ed in he s udy.
We discuss esul s ob ained using he simula ion p og am
in an PC machine wi h an In el-i3 M370 (2.4 GHz) CPU,
unning a Linux 3.2.29 ke nel. The na i e compile used
is GCC 4.7.1, wi h he op imiza ion lag −O3. We ha e
compiled wo e sions o he simula ion code: One using he
o iginal implemen a ion o he π unc ion (O g); and one using
he al e na i e implemen a ion o he in e se unc ion (Al ).
The p og ams a e execu ed wi h a la ge ange o p oblem
sizes (N∈[103,106], T =N/3), and numbe o p ocesso s
(P∈[102,104]). The simula ion s a s a he i s i e a ion o
he ou e loop 2, whe e he maximum ange o he dis ibu ed
loop 4is achie ed. Then, he code uns 100 consecu i e
i e a ions o 2. Measu es ha e been eplica ed wi h a bi a y
p ocesso numbe s p= 4,17,29, ..., ob aining he same
esul s.
Figu e 3 and 4 shows he measu ed cos o 100 i e a ions o
he communica ion code s ages o he o iginal (O g) p og am
and he al e na i e code (Al ), when ixing one o he un- ime
pa ame e s (No P). The execu ion imes o he sequen ial
pa o he code, ha do he ac ual compu a ion, a e also shown
wi h a line. No ice he loga i hmic scale on y-axis.
We can obse e wi h he o iginal π unc ion implemen a ion
how as he calcula ions associa ed wi h communica ion code
exceed by o de s o magni ude he compu a ion ime, when he
pa ame e s g ow. The p oduc o Nand Pin he unpacking
code due o he cos o he π unc ion domina es he cos ,
g owing o mo e han one minu e o clock ime o big p oblem
sizes, o a high numbe o p ocesso s.
Wi h ou p oposed al e na i e implemen a ion, he Pmul-
iplie in oduced by he π unc ion disappea s. I can be seen
in he igh o igu e 4 how he unpacking pa o he code
is no mo e a ec ed by i . When he numbe o p ocesso s P
g ows, he amoun o wo k o be done by each local p ocess
is p opo ionally educed. Ne e heless, he communica ion
code cos is s ill dependen on he o e all p oblem size. In
ou expe imen s, i exceeds he cos o he compu a ion in
one o de o magni ude o enough numbe o p ocesso s. We
can see in he igu e 3 how he cos o he unpacking unc ion
g ows as e han he compu a ion e o o big p oblem sizes.
VI. EXPERIMENTAL STUDY
In his sec ion we discuss a eal expe imen al s udy pe -
o med o e i y ha he asymp o ic beha iou o eal codes
execu ed in eal machines ollows he same beha iou as he
simula ion esul s, and can be p edic ed using he p oposed
cos model.
A. Expe imen al en i onmen
We ha e chosen h ee s udy cases included wi h Plu o
compile as examples, and also included in he Polybench
benchma ks. The i s one is he al eady discussed 1D Jacobi
p og am. The second one is a 2D Jacobi p og am, and he
hi d one a Floyd-Wa shall’s algo i hm implemen a ion. This
p og ams ep esen s examples o he classes o p og ams in
Polybench ha gene a es communica ion code. Linea algeb a
examples do no de i e in ac ual communica ions because
Plu o ans o ma ions assume ha he whole da a s uc u es
a e no dis ibu ed, bu eplica ed on each p ocesso , de i ing
in emp y se s o low-ou dependences ac oss p ocesso s.
We ha e compiled wo e sions o each gene a ed p og am.
One using he o iginal implemen a ion o he π unc ion (O g);
and one using he al e na i e implemen a ion o he in e se
π unc ion p oposed (Al ). The expe imen s we e execu ed
in a sha ed-memo y machine (He acles), 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, adding up o 64 co es in o al. Fo
his expe imen a ion using a eal pla o m, we ha e limi ed
he p oblem size N, and he numbe o p ocesso s P o he
maximum suppo ed by he a ge machine.
B. Resul s
Figu e 5 show he expe imen al pe o mance measu es
ob ained o he h ee s udy cases. We can obse e he same
p edic ed esul s han in he simula ion s udy in Sec . V-B,
bu in a smalle scale due o he smalle Nand Ppa ame e
alues.
The impac on he pe o mance o changing he o iginal
π unc ion by ou p oposed al e na i e is mo e no iceable in
some p oblems han in o he s. I depends on he a io be ween
sequen ial compu a ion and communica ions imes. Fo he
h ee cases o s udy, he mos no iceable e ec appea s o
he Floyd-wa shall case, whe e he packing/unpacking cos is
almos 30% o he o al execu ion ime, as epo ed in [3].
This is due o: (1) A highe h(P) ac o o his algo i hm
compa ing wi h he neighbo synch oniza ion s uc u e o he
Jacobi p og ams; and (2) a highe numbe o communica ion
ins ances in he loop. This e ec is also p edic ed by he
p oposed cos model.
We can also obse e ha , as p edic ed by he model, e en
a e applying ou p oposed al e na i e implemen a ion o he
π unc ion, he e is s ill a p opo ional inc emen o he com-
munica ion calcula ion cos a un- ime, wi h he p oblem size
N. The FOP scheme elays on checking i communica ions
a e needed o he whole space o dis ibu ed iles, on each
p ocesso .
Ou esul s show ha he cos model is an use ul ool o
p edic he asymp o ic beha iou o he code in oduced o
manage he communica ion. I can be used o loca e scalabili y
limi a ions, in o de o ake design o implemen a ion decisions
o a oid hem. We also show how ou p oposed al e na i e o
he implemen a ion o he π unc ion leads o he elimina ion
o one o hese scalabili y p oblems.
VII. CONCLUSION
This pape p esen s a model o he un- ime cos o
he codes gene a ed by a s a e-o - he-a polyhed al-model
echnique (FOP scheme), o communica ion managemen in
a dis ibu ed-memo y en i onmen . The model allows o s udy
he asymp o ic beha iou o he pe o mance o hese pa s o
he code, in e ms o he p oblem size N, and he numbe
o p ocesso s P. I highligh s po en ial scalabili y limi a ions,
helping esea che s o iden i y hem and possibly elimina e
hem in u u e designs and implemen a ions.
This s udy shows how he model is used o de ec scalabili y
limi a ions. We also p opose and al e na i e way o imple-
0.0001
0.001
0.01
0.1
1
10000 50000 100000 500000 1000000
Time (sec.)
N (a ay size)
O iginal, Jacobi 1D ( P = 64 )
Packing
Coo dina ion
Unpacking
Compu a ion
0.0001
0.001
0.01
0.1
1
10000 50000 100000 500000 1000000
Time (sec.)
N (a ay size)
Al e na i e, Jacobi 1D ( P = 64 )
Packing
Coo dina ion
Unpacking
Compu a ion
0.1
1
10
100
8000 10000 12000 14000 16000
Time (sec.)
√N (a ay size)
O iginal, Jacobi 2D ( P = 64 )
Packing
Coo dina ion
Unpacking
Compu a ion
0.1
1
10
100
8000 10000 12000 14000 16000
Time (sec.)
√N (a ay size)
Al e na i e, Jacobi 2D ( P = 64 )
Packing
Coo dina ion
Unpacking
Compu a ion
1
10
100
1000
10000
5000 8000 10000 12000 14000
Time (sec.)
√N (a ay size)
O iginal, Floyd Wa shall ( P = 64 )
Packing
Coo dina ion
Unpacking
Compu a ion
1
10
100
1000
10000
5000 8000 10000 12000 14000
Time (sec.)
√N (a ay size)
Al e na i e, Floyd Wa shall ( P = 64 )
Packing
Coo dina ion
Unpacking
Compu a ion
Figu e 5. Execu ion imes o he codes gene a ed using he FOP scheme, wi h he o iginal and he al e na i e π unc ion implemen a ion, o he h ee s udy
cases: 1D Jacobi, 2D Jacobi, and Floyd-Wa shall’s algo i hm.