Model of Distributed Computations
Model of Distributed Computations
Computations
Ajay Kshemkalyani and Mukesh Singhal
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 1/ 1
Distributed Computing: Principles, Algorithms, and Systems
A Distributed Program
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 2/ 1
Distributed Computing: Principles, Algorithms, and Systems
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 3/ 1
Distributed Computing: Principles, Algorithms, and Systems
H i = (hi , → i )
hi is the set of events produced by pi and
binary relation →i defines a linear order on these events.
Relation → i expresses causal dependencies among the events
of pi .
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 4/ 1
Distributed Computing: Principles, Algorithms, and Systems
The send and the receive events signify the flow of information between
processes and establish causal dependency from the sender process to the
receiver process.
A relation →msg that captures the causal dependency due to message
exchange, is defined as follows. For every message m that is exchanged
between two processes, we have
send (m) → msg rec (m).
Relation →msg defines causal dependencies between the pairs of
corresponding send and receive events.
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 5/ 1
Distributed Computing: Principles, Algorithms, and Systems
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 6/ 1
Distributed Computing: Principles, Algorithms, and Systems
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 7/ 1
Distributed Computing: Principles, Algorithms, and Systems
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 8/ 1
Distributed Computing: Principles, Algorithms, and Systems
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 9/ 1
Distributed Computing: Principles, Algorithms, and Systems
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 10 / 1
Distributed Computing: Principles, Algorithms, and Systems
Concurrent events
For any two events ei and ej , if ei /→ ej and ej /→ ei ,
then events ei and ej are said to be concurrent (denoted as ei ǁ ej ).
In the execution of Figure 2.1, e3 ǁ e3 and e4 ǁ e1.
1
3
2
3
The relation ǁ is not transitive; that is, (ei ǁ ej ) ∧ (ej ǁ ek ) /⇒ ei ǁ
ek .
For example, in Figure 2.1, e3 ǁ e4 and e4 ǁ e5, however, e3 /ǁ e5.
3 2
2
1
3
1
For any two events ei and ej in a distributed execution,
ei → ej or ej → ei , or ei ǁ ej .
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 11 / 1
Distributed Computing: Principles, Algorithms, and Systems
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 12 / 1
Distributed Computing: Principles, Algorithms, and Systems
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 13 / 1
Distributed Computing: Principles, Algorithms, and Systems
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 14 / 1
Distributed Computing: Principles, Algorithms, and Systems
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 15 / 1
Distributed Computing: Principles, Algorithms, and Systems
Not at ions
LSix denotes the state of process pi after the occurrence of event ex and
before
i the event eix+1.
LSi 0 denotes the initial state of process
pi .x is a result of the execution of all the events executed by process pi till ex .
LS
i
i
Let send (m)≤LSx denote the fact that ∃y :1≤y ≤x :: ey =send (m).
i i
Let rec (m)/≤LSx denote the fact that ∀y :1≤y ≤x :: ey /=rec (m).
i i
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 16 / 1
Distributed Computing: Principles, Algorithms, and Systems
A Channel State
The state of a channel depends upon the states of the processes it
connects.
Let SCx , y denote the state of a channel
ij
Cij . of a channel is defined as
The state
follows: V
SCijx, y ={ mij | send (mij ) ≤ i e
x
rec (m ) /≤ ye
} j ij
Thus, channel state SCijx,y denotes all messages that pi sent upto event ex and
which process pj had not i received until event ey .
j
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 17 / 1
Distributed Computing: Principles, Algorithms, and Systems
Global State
The global state of a distributed system is a collection of the local states
of the processes and the channels.
Notationally, global state GS is defined as,
S
i i j, k SC yj , zk }
S jk
For a global state to be GSmeaningful,
= { LS the
xi states of all the components of
,
the distributed system must be recorded at the same instant.
This will be possible if the local clocks at processes were perfectly
synchronized or if there were a global system clock that can be
instantaneously read by the processes. (However, both are impossible.)
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 18 / 1
Distributed Computing: Principles, Algorithms, and Systems
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 19 / 1
Distributed Computing: Principles, Algorithms, and Systems
p e11 e21 e 13 e 14
1
m
e12 12
e22 e23 e42 m 21
p2
1
p3 e3 e 32 e 33 e 34 e 53
e41 e42
p4
time
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 20 / 1
Distributed Computing: Principles, Algorithms, and Systems
In Figure
2.2:
A global state GS1 = {LS 11, LS23, LS3 3, LS4 2 } is
because the state of p2 has recorded the receipt of message m12,
inconsistent
however,
the state state
A global of p1 GS
has not recorded
consisting of its send.
local states {LS 2, LS 4, LS 4, LS
2 1 2 3 4
is} consistent; all the channels are empty except C21 that
2
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 21 / 1
Distributed Computing: Principles, Algorithms, and Systems
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 22 / 1
Distributed Computing: Principles, Algorithms, and Systems
. . . Cuts of a Distributed
Computation
Figure 2.3: Illustration of cuts in a distributed execution.
C1 C2
e11 e21 e13 e 14
p
1
1
p3 e3 e 32 e 33 e34 e53
e41 e42
p4
time
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 23 / 1
Distributed Computing: Principles, Algorithms, and Systems
. . . Cuts of a Distributed
Computation
In a consistent cut, every message received in the PAST of the cut was sent
in the PAST of that cut. (In Figure 2.3, cut C2 is a consistent cut.)
All messages that cross the cut from the PAST to the FUTURE are in transit
in the corresponding consistent global state.
A cut is inconsistent if a message crosses the cut from the FUTURE to
the PAST. (In Figure 2.3, cut C1 is an inconsistent cut.)
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 24 / 1
Distributed Computing: Principles, Algorithms, and Systems
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 25 / 1
Distributed Computing: Principles, Algorithms, and Systems
PAST( ej ) FUTURE( e j )
e
j
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 26 / 1
Distributed Computing: Principles, Algorithms, and Systems
Let Pasti (ej ) be the set of all those events of Past(ej ) that are on process
pi .
Pasti (ej ) is a totally ordered set, ordered by the relation →i , whose
maximal element is denoted by max (Pasti (ej )).
max (Pasti (ej )) is the latest event at process pi that affected event ej
(Figure 2.4).
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 27 / 1
Distributed Computing: Principles, Algorithms, and Systems
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 28 / 1
Distributed Computing: Principles, Algorithms, and Systems
Define Futurei (ej ) as the set of those events of Future(ej ) that are on
process
pi.
define min(Futurei (ej )) asSthe first event on process pi that is affected by
Define
ej. Min Future(ej ) as (∀i ) { min(Futurei (ej ))} ,
which consists of the first
event at every process that is causally affected by
event ej .
Min Future(ej ) is referred to as the surface of the
future cone of ej .
All events at a process pi that occurred after max
(Past i (ej )) but before
min(Futurei (ej )) are concurrent with ej .
Therefore, all and only those events of computation H that belong to the set
“H − Past(ej ) − Future(ej )” are concurrent with event ej .
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 29 / 1
Distributed Computing: Principles, Algorithms, and Systems
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 30 / 1
Distributed Computing: Principles, Algorithms, and Systems
. . . M odels of Process
Communications
Neither of the communication models is superior to the other.
Asynchronous communication provides higher parallelism because the sender
process can execute while the message is in transit to the receiver.
However, A buffer overflow may occur if a process sends a large number of
messages in a burst to another process.
Thus, an implementation of asynchronous communication requires more
complex buffer management.
In addition, due to higher degree of parallelism and non-determinism, it is
much more difficult to design, verify, and implement distributed algorithms
for asynchronous communications.
Synchronous communication is simpler to handle and implement.
However, due to frequent blocking, it is likely to have poor performance and
is likely to be more prone to deadlocks.
A. Kshemkalyani and M. Singhal (Distributed Comput A Model of Distributed Computations CUP 2008 31 / 1