scieee Open visual document viewer

HitFlow: A Dataflow Programming Model for Hybrid Distributed- and Shared-Memory Systems

Fresno Bausela, Javier,Barba Gutiérrez, Daniel,González Escribano, Arturo,Llanos Ferraris, Diego Rafael

Abstract

Producción Científica

Full text

In e na ional Jou nal o Pa allel P og amming manusc ip No. (will be inse ed by he edi o ) Hi Flow: A Da a low P og amming Model o Hyb id Dis ibu ed- and Sha ed-Memo y Sys ems Ja ie F esno ·Daniel Ba ba · A u o Gonzalez-Esc ibano · Diego R. Llanos Recei ed: da e / Accep ed: da e Abs ac Da a low p og amming consis s in de eloping a p og am by de- sc ibing i s sequen ial s ages and he in e ac ions be ween hem. The un- imes suppo ing his kind o p og amming a e esponsible o exploi ing he pa allelism by concu en ly execu ing he di e en s ages as soon as hei de- pendencies a e me . In his pape we in oduce a new pa allel p og amming model and ame- wo k based on he da a low pa adigm. I p esen s a new combina ion o ea- u es ha allows o easily map p og ams o sha ed o dis ibu ed memo y, exploi ing da a locali y and a ini y o ob ain he same pe o mance han op- imized coa se-g ain MPI p og ams. These ea u es include: I is a unique one- ie model ha suppo s hyb id sha ed- and dis ibu ed-memo y sys ems wi h he same abs ac ions; i can exp ess ac i i ies a bi a ily linked, including non-nes ed cycles; i uses in e nally a dis ibu ed wo k-s ealing mechanism o allow Mul iple-P oduce /Mul iple-Consume con igu a ions; and i has a un- ime mechanism o he econ igu a ion o he dependences and communica- ion channels which also allows he c ea ion o ask- o- ask da a a ini ies. We p esen an e alua ion using examples o di e en classes o applica ions. Ex- pe imen al esul s show ha p og ams gene a ed using his amewo k deli e good pe o mance in hyb id dis ibu ed- and sha ed-memo y en i onmen s, wi h a simila de elopmen e o as o he da a low p og amming models o i- en ed o sha ed-memo y. Keywo ds Dis ibu ed sys ems ·Dynamic compu a ion ·Pa allel p og am- ming models ·S eaming compu a ion J. F esno ·D. Ba ba ·A. Gonzalez-Esc ibano ·D. R. Llanos Depa amen o de In o m´a ica, Uni e sidad de Valladolid, Valladolid, Spain J. F esno E-mail: [email p o ec ed]a.es D. Ba ba E-mail: [email p o ec ed]a.es A. Gonzalez-Esc ibano E-mail: [email p o ec ed]a.es D. R. Llanos E-mail: [email p o ec ed]a.es 2 Ja ie F esno e al. 1 In oduc ion The mos common p og amming ools o pa allel machines a e based on mes- sage passing lib a ies, such as MPI [1], o sha ed memo y APIs like OpenMP [2]. These ools allow he p og amme s o exploi he capabili ies o he machines by explici ly de ine he pa allel sec ions inse ed in he sequen ial code and p og am in e -p ocess synch oniza ions and communica ions. On he o he hand, s eam and da a low lib a ies and languages (such as Fas Flow [3], CnC [4], OpenS eam [5], o S-Ne [6]) educe he complexi y o c ea ing a pa allel p og am because he p og amme only has o de ine he sequen ial s ages and i s dependencies. I is he esponsibili y o he un ime o con ol he sequen ial s ages execu ion and pe o m he da a synch oniza- ions. Howe e , hese models do no p esen speci ic ea u es o exp ess some compu a ional pa e ns, o o ob ain communica ion-e icien implemen a ions on dis ibu ed p ocesses. In his wo k we p opose a no el combina ion o ea u es o da a low p og amming models: (a) A single one- ie ep esen a ion o sha ed- and dis ibu ed-memo y a chi ec u es; (b) Desc ip ion o a p og am as a econ- igu able ne wo k o ac i i ies and yped da a con aine s a bi a ily in e con- nec ed, wi h a gene ic sys em o ep esen dis ibu ed Mul iple-P oduce /Mul- iple-Consume (MPMC) con igu a ions; (c) Suppo o dependence s uc- u es ha in ol e non-nes ed eedback loops; (d) A mechanisms o econ igu e dependences a un ime wi hou c ea ing new asks; and (e) A mechanism o in ui i ely exp ess ask- o- ask a ini ies which would allow a be e exploi a- ion o da a locali y ac oss s a e-d i en ac i i ies. As a p oo o concep we ha e de ised Hi Flow, a new da a low pa allel p og amming model and ame- wo k ha ex ends a p e ious p oposal [7] o include all hese ea u es. Table 1 shows a compa ison o di e en da a low solu ions in e ms o hese ea u es. This combina ion o ea u es allows he c ea ion o ne wo ks o asks ha can be mapped o message-passing p ocesses wi h a ixed scheduling. The capaci y o econ igu ing he dependences and ac i i ies o a ask allows he un ime modi ica ion o he communica ion pa e n used a each compu a ion s age, wi hou he need o c ea ing o scheduling new asks. Tasks can alloca e on hei local con ex s bu e s, o da a pa s assigned wi h a classical da a pa i ion policy, ha pe sis ac oss di e en s ages. In his way, da a can main ain he a ini y wi h he message-passing p ocesses and ac oss ela ed asks, a oiding cos ly mig a ions and op imizing he communica ions. This scheme leads o implemen a ions wi h simila pe o mance and scalabili y han p og ams manually de eloped and op imized using message-passing models, such as MPI. We p esen an e alua ion o ou p oposal using examples o ou di e - en applica ion classes. We desc ibe how hey a e ep esen ed in ou model, showing how o exp ess di e en ypes o pa allel pa adigms including s a ic and dynamic synch oniza ion s uc u es. Expe imen al wo k has been ca ied ou o p o e ha he p og ams gene a ed using ou amewo k achie e good pe o mance in compa ison wi h manually de eloped implemen a ions using Hi Flow: A Da a low P og amming Model 3 Fas Flow CnC OpenS eam S-Ne Hi Flow Single ie model X X X Recon igu e dependencies X X Allows asks a ini ies X MPMC con igu a ions X X X X Feedback loops X X X X Table 1: Compa ison o da a low lib a ies. bo h message-passing lib a ies such as MPI, and s a e-o - he-a ools o pa - allel da a low p og amming, like Fas Flow [3] o CnC [4]. These expe imen s show ha he o e heads in oduced by he new abs ac ions do no ha e a signi ican impac on pe o mance. Finally, an analysis o di e en de elop- men e o me ics shows ha he cos o p og amming using ou p oposal, a ge ing hyb id dis ibu ed- and sha ed-memo y sys ems, is simila o o he sha ed-memo y da a low app oaches, highly educing he p og amming cos compa ing wi h using message passing di ec ly. The es o he pape is o ganized as ollows. Sec ion 2 desc ibes ou p o- posal. A discussion abou i s usage is gi en in Sec . 3 while Sec . 4 shows he implemen a ion de ails. Sec ion 5 p esen s he expe imen al wo k ca ied ou o es he implemen a ion. Sec ion 6 desc ibes some ela ed wo k in he ield. Finally, he conclusions o he pape a e in Sec . 7. 2 The Hi Flow model In his sec ion we p esen Hi Flow, a new pa allel p og amming amewo k implemen ed in C++ ha exploi s da a low pa allelism o bo h sha ed- and dis ibu ed-memo y sys ems. The Hi Flow p og amming model akes i s no a- ion om Colo ed Pe i ne s [8]. A Hi Flow p og am is a ne wo k composed o wo kinds o nodes, called places and ansi ions. The places a e sha ed-da a con aine s ha keep okens, while he ansi ions a e he sequen ial p ocessing componen s o he sys em. T ansi ions a e connec ed by di ec ed channels o places, wi h he di ec ion de e mining he inpu o ou pu ole o places o each ansi ion (see Fig. 1). A ansi ion akes one oken om each o i s inpu places and pe o ms some ac i i y wi h hem. I may hen add okens o any/all o i s ou pu places. This ac i i y is epea ed while he e a e okens a i ing o he inpu places. We p opose he compu a ion inside he ansi ions o be mode-d i en. Us- ing a ma hema ical no a ion, a p og am o compu a ion is ep esen ed by P={p1, p2, . . . , pn}a ini e se o places, and T={ 1, 2, . . . , m}a ini e se o ansi ions. The ansi ions a e composed o modes: i={m1 i, m2 i, . . . , mo i}. Each mode mj iis a uple h , I, O, nex i, whe e is a sequen ial unc ion, I⊆Pa e he inpu channels, O⊆Pa e he ou pu channels, and nex ∈ {m1 i, m2 i, . . . , mo i}∪ END is he mode ha will be ac i a ed a e he cu en mode miends. 4 Ja ie F esno e al. T ansi ion A T ansi ion B modes: Mode mAMode mA Ne wo k c ea ion Ne wo k execu ion T ansi ion C B C A B C A p1 p2 12 Fig. 1: Ne wo k example wi h modes. T ansi ion A has wo modes (A1and A2), each mode enables a di e en ou pu channel connec ing A wi h B o A wi h C. Modes a e used o de ine mu ually-exclusi e ac i i ies inside he ansi- ions, and dynamically econ igu e he ne wo k. A mode enables a subse o connec ions o inpu places o ou pu places. The sequen ial unc ion is exe- cu ed when okens a i e in he inpu places o he ac i e mode. A ansi ion changes i s mode when all he okens om he ac i e mode ha e been p o- cessed. To de ec ha he e a e no mo e okens emaining o pending o a i e o he inpu places, special signal okens a e used o in o m o a mode change (mode-change signal). The change o mode in a ansi ion au oma ically sends mode-change signals o all i s ou pu places. Thus, signals a e p opaga ed au oma ically ac oss he ne wo k, lushing okens p oduced on he p e ious mode, be o e changing u he ansi ion o he new mode. When a ansi ion change i s mode, inpu and ou pu places a e econ igu ed acco ding o he new mode speci ica ion. An example o a ne wo k wi h modes can be seen in Fig. 1. The ne wo k has a ansi ion (A) wi h wo modes. On each mode, he ansi ion will send okens o a di e en des ina ion B o C. Finally, he modes can be used o enable da a locali y, de ining ask- o- ask a ini ies. Tasks implemen ed as unc ions o di e en modes in he same ansi ion a e mu ually exclusi e and a e execu ed by he same h ead so hey can sha e da a s uc u es. Fo example, da a a ini y is used in he Smi h- Wa e man algo i hm, which is one o he benchma ks discussed in he expe - imen al sec ion. This benchma k pe o ms a wo-phase wa e on algo i hm (see Fig. 2). In he i s phase, asks calcula e he elemen s o a ma ix s a - ing om he op le elemen . The second phase is a back acking sea ch ha s a s om he bo om- igh elemen , and each ask wo ks on a pa o he he ma ix ob ained in he i s phase. As i is shown in Fig. 2, i is possible o c ea e a ne wo k o model his kind o p oblems wi hou using he modes. Howe e , using he modes, we can old ha ne wo k by adding wo di e en ac i i ies in he ansi ions, one o each phase o he algo i hm. Thus, each ansi ion can pe o m he wo equi ed s ages by sha ing i s assigned po ion o he ma ix, a oiding communica ions o he ma ix po ions ha would imply sending big okens h ough places. Hi Flow: A Da a low P og amming Model 5 Ne wo k wi hou modes Ne wo k wi h modes Ma ix M1,0 M1,1 M1,2 M0,0 M0,1 M0,2 M2,0 M2,1 M2,2 M1,0 M1,1 M1,2 M0,0 M0,1 M0,2 M2,0 M2,1 M2,2 M1,0 M1,1 M1,2 M0,0 M0,1 M0,2 M2,0 M2,1 M1,0 M1,1 M1,2 M0,0 M0,1 M0,2 M2,0 M2,1 M2,2 Fig. 2: Smi h-Wa e man ne wo k s uc u e wi h and wi hou modes. 3 P og amming wi h Hi Flow We ha e de eloped a p o o ype o a amewo k o implemen pa allel p og ams in acco dance wi h he p oposed model. The cu en p o o ype elies on POSIX Th eads P og amming (P h eads) and he s anda d Message Passing In e ace (MPI) o suppo bo h sha ed- and dis ibu ed-memo y a chi ec u es. We de- cided o use P h eads in he p o o ype because he C++ S anda d Lib a y h eads whe e no ully suppo ed a he ime he de elopmen began. Po ing he cu en code o use na i e C++ h eads would be s aigh o wa d. This sec ion explains he key ea u es o he p og amming amewo k. I con ains a summa y o he Hi Flow API, a desc ip ion o how o build a p og am ne wo k, and de ails abou he mode seman ics. The main Hi Flow classes a e shown in he UML diag am in Fig. 3. A able wi h he API me hods can be ound in [9]. 3.1 Building ansi ions To use his amewo k, he use has o c ea e a class which ex ends he p o ided T ansi ion class wi h he sequen ial ac i i ies o he p og am (See example in Fig. 4). The ini and end me hods can be ex ended o execu e s a ing and ending ac ions be o e and a e he execu ion o he p og am. The use classes should in oduce one o mo e new me hods wi h a bi a y names o encapsula e he code o pa icula mode ac i i ies. The associa ion be ween modes and ac i i y me hods is es ablished when building he ne wo k (see sec ion 3.2). The ac i i y me hod is au oma ically called when he e a e okens o be p ocessed in he inpu places decla ed o i s mode. I he e a e no inpu places o a pa icula mode, i will be called jus once. The use -de ined ac i i y me hods can use he T ansi ion::ge o T ansi ion::pu me hods o e ie e okens om, o append okens o he cu en mode places. The ge me hod e ie es one oken o each o he ac i e inpu places. On each ac i i y me hod in oca ion, Hi Flow ensu es ha he ge me hod can be called once. Addi ional calls o ge will block un il he e is a leas one oken in each inpu 6 Ja ie F esno e al. <<bind>> Place # Place(s ing name, ReduceOp op) + se MaxSize(in s) T Ne add(T ansi ion* ) un() T ansi ion + addMe hod(Me hod * me hod, s ing mode = "de aul ") + addInpu (Place * p, s ing mode = "de aul ") + addOu pu (Place * p, s ing mode = "de aul ", bool eedback = alse) + ini () + end() # ge (T1* 1, [T2* 2=NULL,...]) # pu (T& , in placeid=0) # mode(s ing name) Place<in > Mode s ing name Use T ansi ion Fig. 3: UML diag am o he amewo k. place. The pu me hod adds a oken o a speci ic ou pu place. The ou pu place can be selec ed by i s iden i ie using he second a gumen o he pu me hod. I can be omi ed i he e is only one ac i e ou pu place in he mode. A mode au oma ically inishes when: (a) The p oduce ansi ions ha e sen a mode-end signal indica ing ha hey ha e inished he ac i i y in ha mode; and (b) All he okens ha we e gene a ed in he p e ious mode ha e been consumed om he inpu places. A his momen , he ansi ion sends end-mode signal okens o he ac i e ou pu places and au oma ically e ol es o he nex -p og ammed mode. The nex -p og ammed mode can be changed by calling he me hod T ansi ion::mode a any ime. I i is no changed by he use , he de aul nex mode is END, ha is used o inish he compu a ion. The example in Fig. 4 ex ends he T ansi ion class by decla ing a use ac i i y me hod. The me hod e ie es a oken om one place, p ocesses i , and sends he esul o an ou pu place. The okens a e C++ a iables o any ype, handled using empla e me h- ods. The ma shaling and unma shaling is done in e nally wi h MPI unc ions. The basic ypes (cha ,in , loa , ...) a e enabled by de aul . Use -de ined ypes equi e he p og amme o decla e a da a ype in oking he Hi Flow unc ion (hi TypeC ea e) ha in e nally gene a es and egis e s he p ope MPI de i ed ype. 3.2 Building he ne wo k Once he ansi ion classes a e de ined, he p og amme builds he ne wo k in he main unc ion o he C++ p og am. This implies c ea ing ansi ion and place objec s, associa ing he ac i i y me hods, inpu , and ou pu places o modes on he ansi ions, and inally adding he ansi ions o a Ne objec . Hi Flow: A Da a low P og amming Model 7 1class MyT ansi ion: public T ansi ion { 2public: 3 oid execu e(){ // Use ac i i y me hod 4double in ask; 5ge (&in ask); // Re ie e a oken om he place 6double ou ask = p ocess(in ask); 7pu (&ou ask); // Pu he oken in o he ou pu 8} 9}; Fig. 4: Hi Flow example o he c ea ion o a T ansi ion ex ending he basic T ansi ion class. 1Place<double> placeA, placeB; // Decla e he places 2placeA.se MaxSize(10); // Se he place size 3 4MyT ansi ion ansi ion; 5 6// Add he me hod and places o modeA 7 ansi ion.addMe hod(&MyT ansi ion::execu e,"modeA"); 8 ansi ion.addInpu (&placeA,"modeA"); 9 ansi ion.addOu pu (&placeB,"modeA"); 10 ... 11 12 Ne ne ; // Decla e he ne 13 ne .add(& ansi ion); // Add he ansi ion 14 ne . un(); // Run he ne Fig. 5: Hi Flow example o he ne wo k c ea ion. Fig. 5 shows a simple code o build a ne wo k using he p e iously shown MyT ansi ion ansi ion. The i s s ep is o c ea e he places ha will be used in he applica ion (line 1). The Place class is a empla e class used o build he in e nal commu- nica ion channels. The size o he place de ines he g anula i y o he in e nal communica ions: I is an op imiza ion pa ame e ha ep esen s he num- be o packed okens ha will be ans e ed oge he . The use can se i in acco dance wi h he oken gene a ion a io o he ansi ion. The nex s ep is o se he ac i i y me hod and he inpu s and ou pu s o each mode. The addInpu ,addOu pu , and addMe hod me hods, ha e an op ional pa ame e o speci y he mode. When his pa ame e is no speci ied, a de aul END mode is implici ly selec ed. Lines s a ing a 7 se he ac i i y me hod, an inpu place, and an ou pu place o he de aul mode. Mul iple calls o he addInpu o addOu pu o he same ansi ion mode, allow MPMC cons uc ions o be buil . Finally, all he ansi ions a e added o a Ne class ha con ols he map- ping and he execu ion (lines 12 and 13). Line 14 in okes he Ne :: un me hod ha s a s he compu a ion. 8 Ja ie F esno e al. 3.3 Mapping Using Hi Flow, he p og amme can p o ide a mapping policy o assign an- si ions o he a ailable MPI p ocesses. I i is no p o ided, he e is a de aul allback policy implemen ing a simple ound- obin algo i hm. MPI p ocesses wi h mo e han one mapped ansi ion au oma ically spawn addi ional h eads o concu en ly execu e all he ansi ions. Hi Flow implemen a ion sol es he po en ial concu ency p oblems in oduced by synch oniza ion and communi- ca ion when mapping ansi ions o he same p ocess (see Sec . 4.3). In he cu en p o o ype, he mapping policies should p o ide an a ay associa ing indexes o ansi ions o MPI p ocess iden i ie s. 4 Implemen a ion de ails This sec ion discusses some o he implemen a ion challenges associa ed wi h he model, and how hey ha e been sol ed in he cu en amewo k imple- men a ion. 4.1 Ta ge ing bo h sha ed and dis ibu ed sys ems One o he main goals o he amewo k is o suppo bo h sha ed- and dis ibu ed-memo y sys ems wi h a single p og amming le el o abs ac ion. The use -de ined ansi ion objec s ha con ain he logic o he p oblem a e mapped in o he a ailable MPI p ocesses. Since he e may no be enough p o- cesses o all o he ansi ions, h eads a e spawned inside he p ocesses i needed. Only one h ead is spawned o each ansi ion, o execu e he use unc ion and i s po ac i i ies asynch onously o o he ansi ions. The main h ead on each MPI p ocess ini ializes he un ime da a s uc u es, launch he h eads o he ansi ions mapped o i , and wai o hem o inish. Coo dina- ion be ween he spawned eads, o use he sha ed s uc u es o he un ime sys em, is done using mu exes and condi ion a iables. 4.2 Dis ibu ed places The Hi Flow places a e no physically loca ed in a single p ocess. Ins ead, hey a e dis ibu ed oken con aine s. A place is implemen ed as mul iple queues o okens loca ed in he ansi ions ha use ha place as inpu . When load balance equi es i , he okens a e ansmi ed and ea anged be ween he queues on he ansi ions1. This solu ion builds a dis ibu ed MPMC queue mechanism ha exploi s da a locali y, and is mo e scalable han a cen alized 1In he cu en implemen a ion, MPI communica ions a e always used o mo e okens om queue o queue, e en when he queue objec s a e mapped in o he same MPI p ocess. Al hough his simpli ies he implemen a ion, and MPI communica ions a e highly op imized in sha ed-memo y, his decision clea ly opens possibili ies o u he op imiza ion. Hi Flow: A Da a low P og amming Model 9 a) b) c) d) WS e) Fig. 6: T ansla ion om he model design o i s implemen a ion. WS: wo k- s ealing. scheme whe e a single p ocess manages all he okens o a place. Howe e , his is a solu ion ha in oduces coo dina ion challenges ha will be discussed below. In e nally, he dis ibu ed places a e implemen ed using po s ha manage he mo emen o he okens om he sou ce o one o he des ina ion an- si ions. Inpu and ou pu po s a e linked using channels. Fig. 6 shows how he a cs o he model a e implemen ed using po s. The e a e i e possible si ua ions: (a) When a place connec s wo ansi ions, a channel will be cons uc ed o send he okens om he sou ce o he des ina ion. (b) When he e a e wo o mo e inpu places in a ansi ion, he ansi ion will ha e se e al inpu po s, each o hem connec ed o he co esponding sou ce. (c) When wo o mo e ansi ions send okens o a common place, he des- ina ion will ha e a single po ha will ecei e okens, ega dless o he ac ual sou ce. (d) I a place has se e al ou pu ansi ions, any o hem can consume he okens. To allow his beha io , when a place is sha ed by se e al des ina- ions, he sou ce will send okens in a ound- obin ashion o each ou pu po . This can lead o load unbalance i he ime o consume okens in he des ina ions is no compensa ed. To sol e his, a wo k-s ealing mechanism is used o edis ibu e okens be ween he des ina ion ansi ions. (e) When a ansi ion uses he same place as inpu and ou pu , he oken will low di ec ly o he inpu po o e iciency easons. 4.3 Po s, bu e s, and communica ions This sec ion desc ibes he in e nals o he po objec s and explains he de ails abou he communica ions and bu e ing. Fig. 7 shows an example o a wo- ansi ion ne wo k. The e is a p oduce ha gene a es okens which a e sen o a consume using he place A. The consume p esumably pe o ms a il e 16 Ja ie F esno e al. 0 50 100 150 200 250 300 0 10 20 30 40 50 60 70 Time (seconds) Wo ke s/P ocesses Clo eSW (Sha ed memo y) Hi Flow Manual Fas Flow CnC 40 45 50 55 60 65 70 75 80 85 0 10 20 30 40 50 60 70 Time (seconds) Wo ke s/P ocesses Clo eSW (Dis ibu ed memo y) Hi Flow Manual Fig. 11: Clo e’s Smi h-Wa e man benchma k esul s. 5.2.4 2D Jacobi sol e The las benchma k es ed is a Jacobi sol e ha pe o ms 1,000 i e a ions o a 4-poin s encil compu a ion in a 10,000 x 10,000 bidimensional g id. We compa e Hi Flow agains a manually de eloped e sion using C and MPI (Manual) in he dis ibu ed-memo y sys em, and we also compa e agains OpenMP, Fas Flow, and he e sion p o ided by CnC in he sha ed-memo y case. Fas Flow benchma k is de eloped using he s encil da a pa allel skele on which, among o he pa ame e s, accep s he app op ia e unc ion o upda e each cell. When ying o de elop his benchma k using dis ibu ed Fas Flow, we encoun e ed he same p oblems as in Clo e’s e sion o Smi h-Wa e man and we we e unable o implemen i . The Manual e sion is a classical s en- cil implemen a ion ha di ides he g id in o po ions and uses a neighbo synch oniza ion communica ion s uc u e o exchange bo de da a on each compu a ion i e a ion. The Hi Flow e sion uses he same pa i ion policy. Each pa i ion is assigned o a ansi ion ha communica es wi h i s neigh- bo s, sending and ecei ing he da a o he bo de s using places in wo di e en modes. The esul s in Fig. 12 show ha Hi Flow ob ains a simila pe o mance o manual C+MPI in dis ibu ed memo y, and also simila o Fas Flow, and OpenMP e sions o sha ed-memo y. The implemen a ion p o ided by CnC does no show a good pe o mance, and i is no p epa ed o unning on dis- ibu ed memo y. As expec ed, he dis ibu ed expe imen s show a deg ada ion o pe o mance when changing om one single node o se e al he e ogeneous nodes. Howe e hey ob ain a good scalabili y when mo e nodes a e added. These esul s show ha he Hi Flow model can also be applied o p oblems ha a e usually sol ed using s a ic da a pa allel models. Hi Flow: A Da a low P og amming Model 17 0 10 20 30 40 50 60 70 80 0 10 20 30 40 50 60 70 Time (seconds) Wo ke s/P ocesses Jacobi 2D (Sha ed memo y) Hi Flow Manual CNC OpenMP Fas Flow 90 100 110 120 130 140 150 160 0 10 20 30 40 50 60 70 Time (seconds) Wo ke s/P ocesses Jacobi 2D (Dis ibu ed memo y) Hi Flow Manual Fig. 12: Jacobi 2D esul s. 5.3 Code complexi y In his sec ion, we use se e al code complexi y and de elopmen -e o me - ics o compa e Hi Flow codes wi h o he p oposals. Fo his compa ison, we use h ee classical de elopmen e o me ics: The numbe o code okens, McCabe’s cycloma ic complexi y [18], and Hals ead’s de elopmen e o [19]. The numbe o okens de ec ed by he p og amming-language pa se , mea- su es he code olume o C/C++ p og ams be e han he numbe o code lines. McCabe’s cycloma ic complexi y is a quan i a i e measu e o he num- be o linea ly independen pa hs h ough a p og am’s sou ce code. Finally, he Hals ead’s de elopmen -e o me ic is also a quan i a i e measu e based on he numbe o ope a o s and ope ands in he sou ce code. They a e ela ed o he men al ac i i y needed by a p og amme o de elop he code, and o he amoun o es cases needed o check he p og am co ec ness. Low cycloma ic complexi y and Hals ead’s de elopmen e o indica e codes which a e simple o de elop and debug. These me ics a e ypically used in he assessmen o so wa e design complexi y. We ha e selec ed he Mandelb o benchma k because i is a simple bench- ma k, and we ha e mo e implemen a ions using di e en p og amming ools. Table 2 shows he measu es ob ained o each me ic, o he di e en e sions o he benchma k. The me ics clea ly show ha da a low abs ac ions allow he ep esen a ion o he a ge p og am wi h less de elopmen e o han us- ing di ec ly MPI. The sha ed-memo y Fas Flow and CnC e sions, ollowed by he Hi Flow e sion a e he simples implemen a ions. Howe e , he egu- la Fas Flow e sion canno be used in dis ibu ed-memo y sys ems, and CnC needs some uning o un i in a dis ibu ed en i onmen . The e sion ha uses he dis ibu ed-memo y suppo o Fas Flow leads o he bigge me ics alues. This is due o he use o he wo- ie model, ha o ces o implemen sepa a ely he coo dina ion logic o he dis ibu ed p ocesses, and he logic used o sha ed memo y inside he nodes. A ull example o a simple pipeline applica ion implemen ed in Hi Flow, and in Fas Flow suppo ing only sha ed-memo y, can be seen in Fig. 13. I 18 Ja ie F esno e al. Me ic Manual MPI Fas Flow Dis . FF CNC Hi Flow Tokens 552 471 935 565 518 McCabe 34 33 57 24 32 Hals ead 8.65E+5 4.46E+5 13.1E+5 4.16E+5 4.93E+5 Table 2: Complexi y compa ison can be obse ed ha he codes o he nodes ac i i ies and coo dina ion a e e y simila , while Fas Flow p esen nea highe -le el abs ac ions. I educes he code complexi y hanks o he use o skele ons o build he ne wo k. This app oach could also be used on op o Hi Flow. This is u he ly discussed a he end o he Rela ed Wo k sec ion. The esul s indica e ha , using he echniques p esen ed in his wo k, da a low abs ac ions in gene al can e icien ly exploi hyb id sha ed- and dis ibu ed-memo y using a one- ie p og amming model, and educing he de elopmen e o compa ing wi h di ec ly using message-passing in e aces. 6 Rela ed wo k In his sec ion we i s commen he di e ences be ween ou cu en p oposal and he p e ious wo k o ou g oup in he same esea ch line. Then, we discuss concep ual simila i ies and di e ences wi h o he da a low o ask-ne wo k o ien ed p og amming models. We ocus he discussion on ea u es ha ha e implica ions in he p og amming s a egies, he implemen a ion echniques used, and he mapping o he asks in he con ex o dis ibu ed p ocesses. Hi Flow is a complemen o Hi map, a lib a y o au oma ic bu s a ic hie a chical mapping, wi h suppo o dense and spa se da a s uc u es [20, 21,22]. The Hi map lib a y ocuses on da a-pa allel echniques and does no ha e a na i e suppo o da a low applica ions. In a p e ious wo k [7], we in oduced a i s app oach o a da a low model ha could be used as a Hi map ex ension. The model in oduced in his pape gene alizes se e al es ic ions o he p e ious one, in oducing a comple e gene ic model o ep esen any kind o combina ions o pa allel s uc u es and pa adigms. The di e ences wi h he p e ious Hi map ex ension can be summa ized as: (1) We p esen a gene al MPMC sys em whe e consume s can consume di e en ask ypes om di e en p oduce s. (2) I suppo s cycles in he ne wo k cons uc ion. (3) The new model in oduces a concep o mode inside he p ocessing uni s o econ igu e he ne wo k, allowing mu ually exclusi e unc ions in a ansi ion, and o in ui i ely de ine ask- o- ask a ini y wi h an easie mapping o ixed- scheduled MPI p ocesses. S-Ne [6] is a decla a i e coo dina ion language. I de ines he s uc u e o a p og am as a se o connec ed asynch onous componen s called boxes. S-Ne only akes ca e o he coo dina ion: The ope a ions done inside boxes a e de- ined using con en ional languages. Boxes a e s a eless componen s wi h only a single inpu and a single ou pu s eam. F om he p og amme s’ pe spec i e, Hi Flow: A Da a low P og amming Model 19 he implemen a ion o s eams on he language le el by ei he sha ed memo y bu e s o dis ibu ed memo y message passing is en i ely anspa en . Hi Flow has se e al simila i ies wi h Fas Flow [3], a s uc u ed pa allel p og amming amewo k a ge ing sha ed memo y mul i-co e a chi ec u es. Fas Flow is s uc u ed as a s ack o laye s ha p o ide di e en le els o abs ac ion, p o iding he pa allel p og amme wi h a se o eady- o-use, pa ame ic algo i hmic skele ons, modeling he mos common pa allelism ex- ploi a ion pa e ns. Hi Flow ansi ion API is simila o Fas Flow. Fig. 13 shows a ull example o a simple pipeline applica ion o compa e bo h o hem. The main di e ences a e ha he Hi Flow amewo k is designed o suppo bo h sha ed- and dis ibu ed-memo y wi h a single ie model. I in- cludes a anspa en mechanism o he co ec e mina ion o ne wo ks e en in he p esence o eedback-edges, and mode-d i en con ol o c ea e a ini y be ween ansi ions in dis ibu ed memo y en i onmen s. The Fas Flow g oup has de eloped an ex ension o Fas Flow o a ge dis ibu ed nodes, using a wo ie model [14]. Howe e , his solu ion o ces he p og amme o imple- men sepa a ely he coo dina ion logic o he dis ibu ed p ocesses, and he logic used o sha ed memo y inside he nodes. I uses a di e en mechanism o ex e nal channels o communica e he asks. In his sense, Hi Flow makes he p og am design independen om he mapping be ween sha ed-memo y and dis ibu ed-memo y le els. Hi Flow ne wo ks a e simila o CnC (Concu en Collec ions [4]) g aphs. CnC is a pa allel p og amming model whe e he compu a ion is de ined by se- ial unc ions called compu a ion s eps and hei seman ic o de ing cons ain s. Like Hi Flow ansi ions, CnC s eps communica e h ough message-passing as well as sha ed memo y using sha ed en i ies called i em collec ions. One o he di e ences be ween Hi Flow and CnC is ha CnC allows he p og amme o gi e he schedule hin s abou he h ead a ini y. Howe e , CnC s eps only execu e one ac i i y each one wi h i s own memo y space. Thus i is no possi- ble o de ine ask o ask a ini ies in he way Hi Flow ansi ions do, o be e map he ask ne wo ks o MPI p ocesses wi hou incu ing in communica ion cos penal ies. The e a e some p oposals ha suppo ask pa allelism in oducing an- no a ions in he sequen ial sou ce code. Fo example, he OpenMP 3.0 ask p imi i es and he dependency ex ensions in oduced in e sion 4.0 o he s anda d [23]. The p og amme exposes da a low in o ma ion using p agmas o de ine he s eam inpu and ou pu ask. The un ime ensu es he coo di- na ion o he di e en elemen s. O he da a low p oposals based on anno a- ions a e: OpenS eam [5], OmpSs [24], and S a PU [25]. All hese p oposals simpli y he de elopmen o ask pa allel p og ams in sha ed-memo y, and OmpSs and S a PU also suppo en i onmen s wi h accele a o de ices, and e en dis ibu ed memo y. These models ely in load balancing mechanisms ha dynamically map asks wi h an a bi a y g anula i y le el de ined by he p og amme . On he o he hand, ou app oach is designed o simpli y he ex- p ession o ask ne wo ks wi h a lexible g anula i y, and o allow he c ea ion o a ini ies be ween asks and dis ibu ed p ocesses. The asks can econ ig- 20 Ja ie F esno e al. u e hei ac i i y and/o communica ion channels o change he compu a ion and communica ion s uc u e ac oss di e en s ages. The pu pose is o be e exploi locali y and o educe he communica ion cos s. Skele on lib a ies p esen an app oach ha ha e a highe le el o abs ac- ion han ou da a low model, wi h lexible implemen a ions o di e en a ge a chi ec u es and hyb id pla o ms (see e.g. SkePU [26] o Muesli [27]). They ypically use a wo- ie model, no ela ed o he a ge pla o m, bu o he p og amming pa adigm. They dis inguishing be ween wo di e en ypes o skele ons ( ask- o da a-pa allel o ien ed). These ypes canno be composed in any o m. Da a-pa allel skele ons can only be he lea es o he composi ion ee. Da a a ini ies ac oss di e en s ages, which means di e en hie a chies o skele ons, a e no p ope ly de ined. Finally, he amoun o included skele ons do no suppo all he applica ions classes suppo ed by a gene ic da a low p og amming model ha can exp ess ask o da a-pa allel compu a ions wi h he same abs ac ion, suppo ing a bi a ily connec ed ansi ions and places, wi h dependences loops. In he con ex o Hi Flow, skele ons could be used as highe -le el abs ac ions o anspa en ly gene a e common asks ne wo ks, by combining a limi ed se o s uc u es. Fas Flow al eady exploi s his app oach as we discussed a he end o he p e ious sec ion. 7 Conclusions This pape p esen s a pa allel p og amming model and amewo k wi h a no el combina ion o ea u es designed o easily map da a low p og ams o dis ibu ed-memo y p ocesses. I allows p og ams o be desc ibed as a ne - wo k o communica ing ac i i ies in an abs ac o m. The sys em allows he implemen a ion o applica ions om simple s a ic pa allel s uc u es, o com- plex combina ions o da a low and dynamic pa allel p og ams. The desc ip ion is decoupled om he mapping echniques o policies, which can be e icien ly applied a un ime, au oma ically adap ing s a ic o dynamic s uc u es o di e en esou ce combina ions. Ou cu en amewo k anspa en ly a ge s hyb id sha ed- and dis ibu ed-memo y pla o ms. We p esen an e alua ion wi h examples o di e en classes o dynamic and s a ic applica ions. Expe imen al pe o mance esul s show ha he o e head in oduced by ou abs ac ions has minimal impac compa ed wi h manu- ally de eloped implemen a ions using MPI. Compa isons wi h o he da a low p og amming ools show ha Hi Flow can be e exp ess some classes o p o- g ams designed o dis ibu ed en i onmen s, while i s implemen a ion can be imp o ed o sha ed-memo y. Compa isons o de elopmen e o me ics indi- ca e ha Hi Flow codes ha e a simila de elopmen cos han o he da a low abs ac ions. Hi Flow codes p esen a much lowe complexi y han manually de eloped MPI codes, and ob ain he same pe o mance and scalabili y. This gene ic amewo k will allow us o ocus ou esea ch on he bes mapping policies ha can anspa en ly a ge he e ogeneous pla o ms o Hi Flow: A Da a low P og amming Model 21 1#include <hi low.h> 2using namespace hi low; 3 4 5class S ageA: public T ansi ion { 6in num asks; 7public: 8S ageA(in ): num asks( ){}; 9 10 oid c ea e(){ 11 long ask; 12 o (in i=0; i<num asks; i++){ 13 ask = i; 14 pu ( ask); 15 } 16 17 } 18 }; 19 20 class S ageB: public T ansi ion { 21 long sum; 22 public: 23 oid ini (){ 24 sum = 0; 25 26 } 27 oid p ocess(){ 28 long ask; 29 ge (& ask); 30 sum += ask; 31 } 32 oid end(){ 33 cou << "Sum " << sum << endl; 34 } 35 }; 36 37 in main(in na gs, cha * a gs[]){ 38 39 Hi Flow::ini (&na gs, & a gs); 40 41 S ageA s _a(10); 42 S ageB s _b; 43 44 Place<long> place("long con aine "); 45 46 s _a.addOu pu (&place,"c ea eTasks"); 47 s _a.addMe hod(&S ageA::c ea e,"c ea eTasks"); 48 s _b.addMe hod(&S ageB::p ocess,"p ocessTasks"); 49 s _b.addInpu (&place,"p ocessTasks"); 50 51 Ne ne ; 52 ne .add(&s _a); ne .add(&s _b); 53 ne . un(); 54 55 e u n 0; 56 } 1#include < /pipeline.hpp> 2using namespace ; 3 4class S ageA: public _node { 5in num asks; 6public: 7S ageA(in ): num asks( ){}; 8 9long * s c( oid * in ask){ 10 11 o (in i=0; i<num asks; i++){ 12 long * ask = new long(i); 13 _send_ou ( ask); 14 } 15 e u n NULL; 16 } 17 }; 18 19 class S ageB: public _node { 20 long sum; 21 public: 22 in s c_ini (){ 23 sum = 0; 24 e u n 0; 25 } 26 long * s c(long * in ask){ 27 sum += *in ask; 28 dele e in ask; 29 e u n GO_ON; 30 } 31 oid s c_end(){ 32 cou << "Sum " << sum << endl; 33 } 34 }; 35 36 in main() { 37 38 _pipeline pipe; 39 pipe.add_s age(new S ageA(10)); 40 pipe.add_s age(new S ageB()); 41 i (pipe. un_and_wai _end()<0) 42 e u n -1; 43 e u n 0; 44 } Fig. 13: A ull pipeline example in bo h Hi Flow (le ) and sha ed-memo y only Fas Flow ( igh ) amewo ks. 22 Ja ie F esno e al. speci ic o gene ic combina ions o pa allel pa adigms, allowing us o build powe ul pa allel pa e ns using a common and gene ic amewo k. Acknowledgemen s This esea ch has been pa ially suppo ed by MICINN (Spain) and ERDF p og am o he Eu opean Union: HomP og-He Sys p ojec (TIN2014-58876- P), PCAS p ojec (TIN2017-88614-R), CAPAP-H6 (TIN2016-81840-REDT), and COST P og am Ac ion IC1305: Ne wo k o Sus ainable Ul ascale Com- pu ing (NESUS). By Jun a de Cas illa y Le´on, p ojec PROPHET (VA082P17). And by he compu ing acili ies o Ex emadu a Resea ch Cen e o Ad anced Technologies (CETA- CIEMAT), unded by he Eu opean Regional De elop- men Fund (ERDF). CETA-CIEMAT belongs o CIEMAT and he Go e n- men o Spain. Re e ences 1. W. G opp, E. Lusk, A. Skjellum, Using MPI : Po able Pa allel P og amming Wi h he Message-passing In e ace, 2nd Edi ion, MIT P ess, 1999. 2. R. Chand a, L. Dagum, D. Koh , D. Maydan, J. McDonald, R. Menon, Pa allel p o- g amming in OpenMP, 1s Edi ion, Mo gan Kau mann, 2001. 3. M. Aldinucci, M. Danelu o, P. Kilpa ick, M. To qua i, Fas Flow: high-le el and e i- cien s eaming on mul i-co e, in: S. Pllana, F. Xha a (Eds.), P og amming Mul i-co e and Many-co e Compu ing Sys ems, 1s Edi ion, Pa allel and Dis ibu ed Compu ing, Wiley, 2017. 4. Z. Budimli´c, M. Bu ke, V. Ca ´e, K. Knobe, G. Lowney, R. New on, J. Palsbe g, D. Peixo o, V. Sa ka , F. Schlimbach, S. Tasi la , Concu en collec ions, Scien i ic P og amming 18 (3-4) (2010) 203–217. doi:10.1155/2010/521797. URL h p://dx.doi.o g/10.1155/2010/521797 5. A. Pop, A. Cohen, A OpenS eam: Exp essi eness and Da a-Flow Compila ion o OpenMP S eaming P og ams, ACM T ansac ions on A chi ec u e and Code Op i- miza ion (TACO) 9 (4) (2013) 53:1–53:25. 6. C. G elck, S.-B. Scholz, A. Sha a enko, A gen le in oduc ion o S-Ne : Typed s eam p ocessing and decla a i e coo dina ion o asynch onous componen s, Pa allel P ocess- ing Le e s 18 (2) (2008) 221–237. 7. J. F esno, A. Gonzalez-Esc ibano, D. R. Llanos, Run ime Suppo o Dynamic Skele ons Implemen a ion, in: P oceedings o he 19 h In e na ional Con e ence on Pa allel and Dis ibu ed P ocessing Techniques and Applica ions (PDPTA), July 22-25, Las Vegas, NV, USA., 2013, pp. 320–326. 8. K. Jensen, L. M. K is ensen, L. Wells, Colou ed Pe i Ne s and CPN Tools o mod- elling and alida ion o concu en sys ems, In e na ional Jou nal on So wa e Tools o Technology T ans e 9 (3-4) (2007) 213–254. 9. J. F esno, A. Gonzalez-Esc ibano, D. R. Llanos, Addi ional ma e ial o he pape : Da a low p og amming model o hyb id dis ibu ed and sha ed memo y sys ems, Tech. Rep. IT-DI-2015-0003, Depa men o Compu e Science, Uni e si y o Valladolid, Spain (2015). URL h p://www.in o .u a.es/~j esno/ epo s/IT-DI-2015-0003.pd 10. D. Holmes, Ideas o pe sis an poin o poin communica ion, Tech. ep., MPI Fo um Mee ings (2014). URL h p://mee ings.mpi- o um.o g/2014-11-scbo -p2p.pd Hi Flow: A Da a low P og amming Model 23 11. J. Dinan, P. Balaji, D. Goodell, D. Mille , M. Sni , R. Thaku , Enabling MPI in e - ope abili y h ough lexible communica ion endpoin s, in: 20 h Eu opean MPI Use s’s G oup Mee ing, Eu oMPI ’13, Mad id, Spain - Sep embe 15 - 18, 2013, 2013, pp. 13– 18. doi:10.1145/2488551.2488553. URL h p://doi.acm.o g/10.1145/2488551.2488553 12. A. Szalkowski, C. Lede ge be , P. K ¨ahenb¨uhl, C. Dessimoz, SWPS3 - as mul i- h eaded ec o ized Smi h-Wa e man o IBM Cell/B.E. and x86/SSE2., BMC esea ch no es 1 (107). 13. P. Clo e, Biologically signi ican sequence alignmen s using Bol zmann p obabili ies, Tech. ep. (2003). URL h p://bioin o ma ics.bc.edu/~clo e/pub/bol zmannPa is03.pd 14. M. Aldinucci, S. Campa, M. Danelu o, Ta ge ing dis ibu ed sys ems in Fas Flow, in: P oceedings o he Eu o-Pa 2012 Pa allel P ocessing Wo kshops, Vol. 7640 o Lec u e No es in Compu e Science, Sp inge Be lin Heidelbe g, 2013, pp. 47–56. 15. J. F esno, A. Gonzalez-Esc ibano, D. R. Llanos, One Tie Da a low P og amming Model o Hyb id Dis ibu ed- and Sha ed-Memo y Sys ems, in: P oc. o In l. Wo k- shop on High-le el P og amming o He e ogeneous and Hie a chical Pa allel Sys ems (HLPGPU), Janua y 19, P ague, Czech Republic., 2016. 16. M. Aldinucci, M. Meneghin, M. To qua i, E icien Smi h-Wa e man on Mul i-co e wi h Fas Flow, in: M. Danelu o, J. Bou geois, T. G oss (Eds.), P oceedings o he 18 h Eu omic o Con e ence on Pa allel, Dis ibu ed and Ne wo k-based P ocessing (PDP), IEEE Compu e Socie y, Pisa, I aly, 2010, pp. 195–199. 17. UniP o Knowledgebase (UniP o KB), www.unip o .o g/. 18. T. J. McCabe, A complexi y measu e, IEEE T ansac ions on so wa e Enginee ing SE- 2 (4) (1976) 308–320. doi:10.1109/TSE.1976.233837. 19. M. H. Hals ead, Elemen s o So wa e Science (Ope a ing and P og amming Sys ems Se ies), Else ie Science Inc., New Yo k, NY, USA, 1977. 20. J. F esno, A. Gonzalez-Esc ibano, D. R. Llanos, Ex ending a hie a chical iling a ays lib a y o suppo spa se da a pa i ioning, The Jou nal o Supe compu ing 64 (1) (2013) 59–68. 21. A. Gonzalez-Esc ibano, Y. To es, J. F esno, D. R. Llanos, An Ex ensible Sys em o Mul ile el Au oma ic Da a Pa i ion and Mapping, IEEE T ansac ions on Pa allel and Dis ibu ed Sys ems 25 (5) (2014) 1145–1154. 22. A. Mo e on-Fe nandez, A. Gonzalez-Esc ibano, D. Llanos, Exploi ing dis ibu ed and sha ed memo y hie a chies wi h hi map, in: High Pe o mance Compu ing Simula ion (HPCS), 2014 In e na ional Con e ence on, 2014, pp. 278–286. doi:10.1109/HPCSim. 2014.6903696. 23. OpenMP A chi ec u e Re iew Boa d, OpenMP applica ion p og am in e ace e sion 4.5 (No embe 2015). URL h p://www.openmp.o g/mp-documen s/openmp-4.5.pd 24. Ba celona Supe compu ing Cen e , OmpSs speci ica ion 4.5 (Decembe 2015). URL h p://pm.bsc.es/ompss-docs/specs 25. C. Augonne , S. Thibaul , R. Namys , P.-A. Wac enie , S a pu: A uni ied pla o m o ask scheduling on he e ogeneous mul ico e a chi ec u es, Concu . Compu . : P ac . Expe . 23 (2) (2011) 187–198. doi:10.1002/cpe.1631. URL h p://dx.doi.o g/10.1002/cpe.1631 26. C. W. K. Johan Enmy en, Skepu: a mul i-backend skele on p og amming lib a y o mul i-gpu sys ems, in: P oceedings o he ou h in e na ional wo kshop on High-le el pa allel p og amming and applica ions (HLPP 2010), ACM, 2010, pp. 5–14. 27. H. K. Philipp Ciechanowicz, Enhancing muesli’s da a pa allel skele ons o mul i-co e compu e a chi ec u es, in: P oceedings o 12 h IEEE In e na ional Con e ence on High Pe o mance Compu ing and Communica ions (HPCC 2010), IEEE, 2010.