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.