scieee Open visual document viewer

On the run-time cost of distributed memory communications generated using the polyhedral model

Moreton Fernández, Ana,González Escribano, Arturo,Llanos Ferraris, Diego Rafael

Abstract

Producción Científica

Full text

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.