scieee Open visual document viewer

A Technique to Automatically Determine Ad-hoc Communication Patterns at Runtime

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

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